ouroboros/tests/test_headless_task_events.py
Ouroboros 57230ee2ef v7next F1: domain D17 - headless split and three test-giant splits, proof-green
headless.py (1573 at tip) gives up its two ledger-assigned leaves again:
ouroboros/headless_status.py (50, 11 symbols) and
ouroboros/workspace_patch_capture.py (668, 19 symbols). All 30 spans are
byte-identical between the reference leaves and git show HEAD bytes — the
hardened transplant --check (mandatory byte gate, undeclared-top-level check,
def681bd) is green on every symbol of both leaves. The facade (947) replays
the oracle's exact edit script over tip bytes: its only divergence from the
reference facade is genuine upstream residue drift (child_ref promotion
machinery, import changes), verified hunk by hunk.

Test splits per the ledger, upstream bytes as truth:
- test_headless_cli.py 2824 -> 462 + five themed siblings + shared fixtures;
  93 test functions preserved exactly (lossless set equality), one adapted
  span kept (the _PATCH_MAX_UNTRACKED_FILE_BYTES monkeypatch retargeted to
  the new leaf, ledger row 739's own adaptation); nine oracle spans carrying
  OTHER domains' v7 spellings reverse-mapped to upstream signatures keyed to
  git show HEAD (registry._run_shell_safety_check string form, tools.core
  _repo_read, queue.init 3-arg, queue.QUEUE_SNAPSHOT_PATH).
- test_workspace_executor.py 1995 -> 541 + three siblings + shared builder;
  two reference spans byte-falsified by upstream drift (06339bb7 readiness
  truth, a849c9a6 probe uncertainty) re-emitted from tip bytes and recorded
  in docs/v7next/LEDGER_CORRECTIONS.md.
- test_agent_task_pipeline.py 1658 -> 1515: only the five ledger-assigned
  _store_task_result rows carved into tests/test_store_task_result.py; the
  other siblings belong to their own lanes.
Thirteen post-cutoff upstream tests have no ledger rows; placed by the
split's theme rule (task_api x4, task_artifacts x1, docker x6, services x2),
disclosed in LEDGER_CORRECTIONS for F5 row-minting.

Pins and mirrors: test_headless_extraction.py transplanted with the
tool_module_inventory clause reduced under an oracle-SHA note (that leaf
belongs to the tools lane); conftest serial table gains the executor family;
the process-custody Popen allowlist row moves headless.py ->
workspace_patch_capture.py with the oracle's justification. All 14 non-split
D17 runtime modules re-proven zero-v7-delta (task_results.py included:
upstream-hot drift stands, no ledger split assigned).

size-ratchet manifest regenerated with the official tool (three test giants
leave GIANT_PATHS, no new band entries); --check green. ruff F clean.
135/105/244/113 tests green in 4-var isolation; HEAD held after every run.

(cherry picked from commit 8dac8303006085bfdb636bacbfc402d6905fae7c)
2026-08-30 17:27:56 +00:00

326 lines
11 KiB
Python

"""Task event replay, log tails and effective child status projection.
Split verbatim out of ``tests/test_headless_cli.py`` by theme. This module
owns what readers see after a task runs: event replay, lineage-filtered log
tails, SSE finalization order, task listing, and the effective-status/result
projection that waits for workspace artifacts.
"""
from __future__ import annotations
import json
import pytest
from starlette.applications import Starlette
from starlette.routing import Route
from starlette.testclient import TestClient
from ouroboros.gateway.tasks import (
api_task_events,
api_task_get,
api_tasks_list,
iter_task_events,
)
from ouroboros.task_results import write_task_result
from tests._headless_cli_shared import ( # noqa: F401 (autouse fixture applies on import)
_managed_worker_pool_available,
)
def test_task_event_replay_uses_existing_logs_and_result(tmp_path):
data = tmp_path / "data"
logs = data / "logs"
logs.mkdir(parents=True)
task_id = "abc123"
(logs / "progress.jsonl").write_text(
json.dumps({"ts": "2026-01-01T00:00:00Z", "task_id": task_id, "content": "working"}) + "\n",
encoding="utf-8",
)
result_dir = data / "task_results"
result_dir.mkdir()
(result_dir / f"{task_id}.json").write_text(
json.dumps({"task_id": task_id, "status": "completed", "result": "done", "ts": "2026-01-01T00:00:01Z"}),
encoding="utf-8",
)
events = iter_task_events(data, task_id)
assert [event["type"] for event in events] == ["progress", "task_result"]
assert events[0]["seq"] == 1
assert events[1]["data"]["result"] == "done"
def test_task_event_replay_parent_includes_child_lineage_events(tmp_path):
data = tmp_path / "data"
logs = data / "logs"
logs.mkdir(parents=True)
parent_id = "parent1"
child_id = "child1"
(logs / "progress.jsonl").write_text(
"\n".join([
json.dumps({"ts": "2026-01-01T00:00:00Z", "task_id": parent_id, "content": "parent"}),
json.dumps({
"ts": "2026-01-01T00:00:01Z",
"task_id": child_id,
"parent_task_id": parent_id,
"root_task_id": parent_id,
"delegation_role": "subagent",
"subagent_task_id": child_id,
"content": "child progress",
}),
]) + "\n",
encoding="utf-8",
)
write_task_result(
data,
parent_id,
"running",
result="parent pending",
ts="2026-01-01T00:00:00Z",
)
write_task_result(
data,
child_id,
"running",
result="child pending",
parent_task_id=parent_id,
root_task_id=parent_id,
delegation_role="subagent",
ts="2026-01-01T00:00:01Z",
)
events = iter_task_events(data, parent_id)
progress_events = [event for event in events if event["type"] == "progress"]
assert [event["task_id"] for event in progress_events] == [parent_id, child_id]
assert progress_events[1]["data"]["content"] == "child progress"
def test_logs_tail_parent_filter_includes_child_lineage_events(tmp_path):
from ouroboros.gateway.logs import api_logs_tail
data = tmp_path / "data"
logs = data / "logs"
logs.mkdir(parents=True)
(logs / "progress.jsonl").write_text(
"\n".join([
json.dumps({"ts": "2026-01-01T00:00:00Z", "task_id": "parent1", "content": "parent"}),
json.dumps({
"ts": "2026-01-01T00:00:01Z",
"task_id": "child1",
"subagent_task_id": "child1",
"parent_task_id": "parent1",
"root_task_id": "parent1",
"delegation_role": "subagent",
"content": "child",
}),
json.dumps({"ts": "2026-01-01T00:00:02Z", "task_id": "other", "content": "other"}),
]) + "\n",
encoding="utf-8",
)
app = Starlette(routes=[Route("/api/logs/{name}", endpoint=api_logs_tail, methods=["GET"])])
app.state.drive_root = data
response = TestClient(app).get("/api/logs/progress?task_id=parent1&limit=10")
payload = response.json()
assert response.status_code == 200
assert [row["content"] for row in payload["entries"]] == ["parent", "child"]
def test_workspace_event_replay_suppresses_task_done_until_artifacts_terminal(tmp_path):
data = tmp_path / "data"
logs = data / "logs"
logs.mkdir(parents=True)
task_id = "abc123"
(logs / "events.jsonl").write_text(
json.dumps({"ts": "2026-01-01T00:00:01Z", "type": "task_done", "task_id": task_id}) + "\n",
encoding="utf-8",
)
write_task_result(
data,
task_id,
"completed",
workspace_root=str(tmp_path / "workspace"),
artifact_status="finalizing",
child_status="completed",
)
events = iter_task_events(data, task_id)
assert "task_done" not in [event["type"] for event in events]
assert events[-1]["type"] == "task_result"
def test_effective_child_completion_waits_for_artifacts(tmp_path):
data = tmp_path / "data"
child = tmp_path / "child"
for root in (data, child):
(root / "task_results").mkdir(parents=True)
write_task_result(
data,
"task-artifacts",
"scheduled",
child_drive_root=str(child),
workspace_root=str(tmp_path / "workspace"),
artifact_status="pending",
result="queued",
)
write_task_result(
child,
"task-artifacts",
"completed",
result="done",
ts="2026-01-01T00:00:02Z",
outcome_axes={
"lifecycle": {"status": "completed"},
"artifacts": {"status": "not_applicable"},
},
)
app = Starlette(routes=[Route("/api/tasks/{task_id}", endpoint=api_task_get, methods=["GET"])])
app.state.drive_root = data
payload = TestClient(app).get("/api/tasks/task-artifacts").json()
assert payload["status"] == "running"
assert payload["artifact_status"] == "finalizing"
assert payload["child_status"] == "completed"
assert payload["outcome_axes"]["lifecycle"]["status"] == "running"
assert payload["outcome_axes"]["artifacts"]["status"] == "finalizing"
write_task_result(data, "task-artifacts", "completed", artifact_status="ready", child_drive_root=str(child), workspace_root=str(tmp_path / "workspace"))
payload = TestClient(app).get("/api/tasks/task-artifacts").json()
assert payload["status"] == "completed"
assert payload["artifact_status"] == "ready"
def test_public_task_result_strips_nested_legacy_result_status(tmp_path):
data = tmp_path / "data"
(data / "task_results").mkdir(parents=True)
write_task_result(
data,
"legacy-loop",
"completed",
result="done",
loop_outcome={"result_status": "failed", "compat_result_status": "failed", "reason_code": "legacy"},
verification_ledger={
"entries": [
{"kind": "legacy", "result_status": "partial"},
{"kind": "nested", "payload": {"compat_result_status": "infra_failed"}},
{"kind": "list", "items": [{"result_status": "failed"}]},
],
},
)
app = Starlette(routes=[Route("/api/tasks/{task_id}", endpoint=api_task_get, methods=["GET"])])
app.state.drive_root = data
payload = TestClient(app).get("/api/tasks/legacy-loop").json()
assert "result_status" not in payload
assert "result_status" not in payload["loop_outcome"]
assert "compat_result_status" not in payload["loop_outcome"]
rendered = json.dumps(payload)
assert "result_status" not in rendered
assert "compat_result_status" not in rendered
def test_effective_child_failure_waits_for_artifacts(tmp_path):
data = tmp_path / "data"
child = tmp_path / "child"
for root in (data, child):
(root / "task_results").mkdir(parents=True)
write_task_result(
data,
"task-failed",
"failed",
child_drive_root=str(child),
workspace_root=str(tmp_path / "workspace"),
artifact_status="finalizing",
child_status="failed",
result="boom",
)
write_task_result(child, "task-failed", "failed", result="boom", ts="2026-01-01T00:00:02Z")
app = Starlette(routes=[Route("/api/tasks/{task_id}", endpoint=api_task_get, methods=["GET"])])
app.state.drive_root = data
payload = TestClient(app).get("/api/tasks/task-failed").json()
assert payload["status"] == "running"
assert payload["artifact_status"] == "finalizing"
assert payload["child_status"] == "failed"
def test_task_sse_emits_final_result_after_cursor_saw_scheduled_result(tmp_path):
data = tmp_path / "data"
(data / "task_results").mkdir(parents=True)
task_id = "abc123"
(data / "task_results" / f"{task_id}.json").write_text(
json.dumps({"task_id": task_id, "status": "completed", "result": "done", "ts": "2026-01-01T00:00:01Z"}),
encoding="utf-8",
)
app = Starlette(routes=[Route("/api/tasks/{task_id}/events", endpoint=api_task_events, methods=["GET"])])
app.state.drive_root = data
response = TestClient(app).get(f"/api/tasks/{task_id}/events?cursor=1&wait=0")
assert response.status_code == 200
assert '"type": "task_result"' in response.text
assert '"status": "completed"' in response.text
def test_task_list_filters_on_effective_child_status(tmp_path):
data = tmp_path / "data"
child_running = tmp_path / "child-running"
child_done = tmp_path / "child-done"
for root in (data, child_running, child_done):
(root / "task_results").mkdir(parents=True)
write_task_result(data, "task-running", "scheduled", child_drive_root=str(child_running), result="queued")
write_task_result(child_running, "task-running", "running", result="working", ts="2026-01-01T00:00:01Z")
write_task_result(data, "task-done", "scheduled", child_drive_root=str(child_done), result="queued")
write_task_result(child_done, "task-done", "completed", result="done", ts="2026-01-01T00:00:02Z")
app = Starlette(routes=[Route("/api/tasks", endpoint=api_tasks_list, methods=["GET"])])
app.state.drive_root = data
client = TestClient(app)
running = client.get("/api/tasks?status=running").json()["tasks"]
completed = client.get("/api/tasks?status=completed").json()["tasks"]
assert [task["task_id"] for task in running] == ["task-running"]
assert running[0]["result"] == "working"
assert [task["task_id"] for task in completed] == ["task-done"]
assert completed[0]["result"] == "done"
@pytest.mark.parametrize("status", ["cancelled", "failed"])
def test_effective_task_result_preserves_parent_terminal_status(tmp_path, status):
data = tmp_path / "data"
child = tmp_path / "child"
for root in (data, child):
(root / "task_results").mkdir(parents=True)
write_task_result(
data,
"task-terminal",
status,
child_drive_root=str(child),
result="parent terminal",
ts="2026-01-01T00:00:02Z",
)
write_task_result(
child,
"task-terminal",
"running",
result="child stale",
ts="2026-01-01T00:00:03Z",
)
app = Starlette(routes=[Route("/api/tasks/{task_id}", endpoint=api_task_get, methods=["GET"])])
app.state.drive_root = data
payload = TestClient(app).get("/api/tasks/task-terminal").json()
assert payload["status"] == status
assert payload["result"] == "parent terminal"
assert payload["ts"] == "2026-01-01T00:00:02Z"