ouroboros/tests/test_native_post_task_wait_browser.py
Ouroboros b9bb1fe311 Restore ordinary conversation capabilities and complete chosen inputs
Keep independent native Main and Project turns addressable through post-task
work and model-access waits, with Project operations bound to their folder.
Separate web-tool restrictions from network access and preserve owner controls.

Deliver complete chosen delegation inputs and keep full operative plan sources
outside the bounded review-state index, retaining paid and legacy custody.

Co-authored-by: Ouroboros <311266734+ouroboros-agent@users.noreply.github.com>
2026-09-08 19:26:53 +03:00

104 lines
5.5 KiB
Python

"""Native post-task wait state and real Stop ingress through the browser card."""
import json
import threading
from types import SimpleNamespace
import pytest
from starlette.applications import Starlette
from starlette.routing import Route
from starlette.testclient import TestClient
from ouroboros import agent_task_pipeline as pipeline, model_wait
from ouroboros.gateway.history import make_chat_history_endpoint
from ouroboros.gateway.state import _chat_activities_snapshot_safe
from ouroboros.gateway.tasks import api_task_cancel
from ouroboros.post_task_checkpoint import post_task_model_wait
from ouroboros.task_results import load_task_result
from supervisor import active_activity, workers
from tests.test_post_task_model_wait import phase as post_phase_fixture, until
from tests.test_project_chat_continuity import _write_history_rows
from tests.test_subscription_setup_browser import capture, subscription_ui as ui_fixture
phase = post_phase_fixture
subscription_ui = ui_fixture
pytestmark = [pytest.mark.ui_browser, pytest.mark.serial]
def test_browser_stops_native_post_task_wait_without_losing_answer(subscription_ui, phase, monkeypatch):
from supervisor import queue as task_queue
f, ui = phase, subscription_ui
task_id, page = f.task["id"], ui["page"]
f.task["_is_direct_chat"] = True
monkeypatch.setattr(task_queue, "DRIVE_ROOT", f.root)
monkeypatch.setattr(task_queue, "RUNNING", {})
monkeypatch.setattr(task_queue, "PENDING", [])
monkeypatch.setattr(workers, "WORKERS", {})
pending = [{"type": "send_message", "task_id": task_id, "chat_id": 1, "text": "Already answered"}]
actor = SimpleNamespace(_busy=True, _accepting_owner_messages=False)
registry = active_activity.get_direct_activity_registry()
def run():
registry.register(task_id, 1, actor=actor)
try:
with model_wait.task_model_wait_scope(task=f.task, drive_root=f.root, event_queue=f.events,
worker_slot_held=False):
pipeline._dispatch_root_post_task(f.env, f.task, "Already answered", f.events, pending,
{"rounds": 3}, {}, {}, f.root / "logs", budget_drive_root="", split_drive=False,
project_scoped=False, project_task=False, parent_env=None, parent_task=None)
finally:
registry.unregister(task_id)
_write_history_rows(f.root, task_id)
# The earlier running progress frame already carried the ordinary host's
# cancel authority. Retain it through history, then use the real endpoint.
path = f.root / "logs/progress.jsonl"
progress = json.loads(path.read_text())
path.write_text(json.dumps({**progress, "cancelable": True}) + "\n")
thread = threading.Thread(target=run)
thread.start()
responses = []
app = Starlette(routes=[Route("/api/chat/history", make_chat_history_endpoint(f.root)),
Route("/api/tasks/{task_id}/cancel", api_task_cancel, methods=["POST"])])
app.state.drive_root = f.root
try:
until(lambda: post_task_model_wait(f.root, task_id) and any(
row.get("credential_harness") for row in post_task_model_wait(f.root, task_id).waits.values()))
with TestClient(app) as client:
def history(route):
response = client.get("/api/chat/history?chat_id=1")
route.fulfill(status=response.status_code, content_type="application/json", body=response.content)
def cancel(route):
response = client.post(f"/api/tasks/{task_id}/cancel", json=route.request.post_data_json)
responses.append(response)
route.fulfill(status=response.status_code, content_type="application/json", body=response.content)
page.route("**/api/chat/history*", history)
page.route("**/api/state", lambda route: route.fulfill(content_type="application/json", body=json.dumps({
"supervisor_ready": True, "projects": [], "active_chat_activities": _chat_activities_snapshot_safe(f.root)})))
page.route(f"**/api/tasks/{task_id}", lambda route: route.fulfill(content_type="application/json",
body=json.dumps(load_task_result(f.root, task_id))))
page.route(f"**/api/tasks/{task_id}/cancel", cancel)
page.goto(ui["url"] + "/")
card = page.locator(f'.chat-live-card[data-task-id="{task_id}"]')
card.locator('[data-wait-id]').first.wait_for(timeout=5000)
assert "worker slot" not in card.locator('.model-wait-footnote').inner_text()
capture(page, "native-post-task-wait-before-stop")
card.locator('[data-cancel-run]').click()
page.locator('[data-task-control="stop_now"]').click()
page.wait_for_function("() => !document.querySelector('.task-control-menu')")
assert responses and responses[-1].status_code in (200, 503)
thread.join(5)
assert not thread.is_alive() and len(f.engine.creates) == 1
assert load_task_result(f.root, task_id)["root_phase_checkpoint"]["post_task_synthesis"] == "degraded"
page.reload()
page.get_by_text("final answer", exact=True).wait_for()
page.wait_for_function("() => !document.querySelector('[data-wait-id]')")
capture(page, "native-post-task-wait-stopped-answer-preserved")
finally:
f.task["_skip_post_task_synthesis"] = True
f.ready.set()
thread.join(5)
assert not thread.is_alive()