fix: resume task events by physical cursor and reuse current result facts

This commit is contained in:
Ouroboros 2026-09-06 02:19:21 +00:00
parent b6c847e31f
commit eea01dd929
17 changed files with 935 additions and 182 deletions

View file

@ -327,12 +327,12 @@ server.py (Starlette+uvicorn) ← HTTP + WebSocket on configurable host:port (de
├── gateway/ ← Gateway Boundary v1: browser-facing route ownership + frontend contract SSOT (see below)
│ ├── contracts.py ← Active WS/HTTP envelope contract owner
│ ├── endpoint_index.py ← `HTTP_ENDPOINTS` index (re-exported by contracts.py); routers own the Route objects
│ ├── schema.py, task_list_scan.py ← The executable gateway contract — JSON Schema derived from the TypedDicts, validating ingress — and the raw result-name scan behind the unfiltered `GET /api/tasks` slice path
│ ├── schema.py, task_list_scan.py ← The executable gateway contract — JSON Schema derived from the TypedDicts, validating ingress — and the stat-invalidated compact result facts shared by list ordering, SSE discovery and Main routing
│ ├── router.py ← Starlette route collector for /api/* and /ws
│ ├── ws.py ← WS manager, extension WS dispatch, broadcast
│ ├── state.py ← /api/health + /api/state
│ ├── tasks.py ← Headless task create/list/get/cancel/events; cancel accepts `stop_policy` (empty = immediate; `finalize_then_cancel` → 202 + open intent → supervisor/owner_stop.py; unknown → 400)
│ ├── task_events.py ← Task-event SSE endpoint
│ ├── task_events.py ← Task-event SSE endpoint: legacy GET ranks plus read-only POST v2 physical-chain cursors
│ ├── task_hurry.py ← POST hurry ingress: exact one-field `{request_id}` body — extra fields refused, because hurry carries no text by design and a smuggled field must not become a side channel; queue-owned admission; idempotent projection (semantics: owner_hurry.py)
│ ├── task_decision.py ← ONE `POST /api/decisions` ingress with family-parsed ids (`quiz:` here, `routing:` → routing_decision.py, `interaction:` reserved); writes `KIND_QUIZ_ANSWER`, broadcasts `quiz_state` (lifecycle: owner_quiz.py; ABI: §11.1)
│ ├── routing_decision.py ← Validates a click against the durable `needs_manual_target` row, recovers the original text, dispatches the existing `steer_task`/`promote_chat_to_task`, confirms through routing_wait receipts; replay-stable derived identities
@ -705,7 +705,9 @@ Project unread state is the durable comparison `visible_revision > project_seen_
A chat instance has an explicit resource lifecycle: `destroy()` marks it dead, disposes subscriptions, listeners, observers, and timers, and removes the DOM last; late async work checks the destroyed flag. `app.js` keeps at most one live Project chat instance; closing or switching stashes scroll intent and destroys it. The narrow exception is client state the server cannot reconstruct (staged `File` objects or an upload in flight): such an instance is hidden and marked pending, reused on reopen, and returns to the destroy policy after that work settles; typed but unsent text survives separately in per-thread session storage. This keeps hidden Project rooms from accumulating listeners or acknowledging unseen revisions without discarding data.
Chat and progress logs rotate as one timeline; archived segments stay durable. Interactive history uses bounded, archive-aware readers expanding only far enough to satisfy the requested thread's filtered quota; Project history cannot be satisfied by unrelated Main rows, and display reads avoid materializing or rebasing artifacts. The task-event stream replays archive-aware and then follows appended bytes, discovering children by a scandir name-diff over the main root's task-results directory; its terminal event performs the one materializing task read that delivers artifact-bearing final truth, and a follow tick decodes only result files not yet successfully read (disclosed residual: the per-tick scan stays proportional to directory size). Subagent admission/custody dedupe is the supervisor queue's transition (§5). The point is bounded UI latency without treating rotation as conversation loss.
All five task-event sources (progress, chat, events, tools and supervisor) read their retained archive chains plus live files. Interactive history uses bounded, archive-aware readers until the requested thread's filtered quota is met; Project history cannot be satisfied by unrelated Main rows, and display reads avoid materializing or rebasing artifacts. Task-event lineage discovery, task-list ordering and Main's newest-result selection share `gateway/task_list_scan.py`'s compact process-local memo: dev/inode/size/mtime/ctime invalidate changed files, failed or torn reads are never cached, and full selected results still pass the schema reader. Directory enumeration remains proportional to the number of results. Main's dialogue projection reads twenty nonempty text rows from a live tail and at most two archives, retaining message ids and the 500-character text preview; its unvisited-row count is explicitly unknown.
`POST /api/tasks/{id}/events` is the read-only SSE v2 transport (`TaskEventsRequest`): JSON `{v:2, wait, cursor}` carries `{v,seq,view,positions}`, with one byte offset per root/source across the complete archive+live chain. Events use physical root/source order, never timestamp ranks; each log row carries a stable root/source/byte-position `event_id` and the cursor after THAT row. Filter changes (task ids, roots, task_done suppression or a proven creation floor) emit `cursor_replay` and replay the new view. The CLI advances only after consuming a frame and deduplicates its most recent 4096 log identities; older replay duplicates remain possible. A `cursor_checkpoint` at wait expiry also advances skipped bytes without incrementing delivery sequence or claiming task progress. Live partial lines wait; immutable malformed/partial lines emit `history_gap` and advance. Unreadable or shorter chains emit `cursor_unavailable` without resetting. Archives must remain immutable and retained: manual prefix removal masked by subsequent growth is outside the offset guarantee. Only explicit creation facts may skip older archives, and skipped bytes still count in positions. Each connection emits a fresh task-result projection and materializes the terminal result once; synthetic results have no log identity and are never deduplicated. Legacy GET and `iter_task_events` keep timestamp-sorted integer ranks, whose retroactive insertions may repeat or omit rows across reconnects; a from-zero replay recovers retained history. The CLI falls back to GET only on a first-connection HTTP 405. Subagent admission/custody dedupe remains the supervisor queue's transition (§5).
The agent-facing `chat_history` reader uses the same live-plus-rotated timeline and may narrow by exact provider, account, conversation, thread, actor, and inclusive date bounds before the count/offset/text-search window — presence provenance is searchable as structured transport fact, not only flattened prose.
@ -853,7 +855,8 @@ Every `/api/files/*` operation resolves its requested path and refuses the opera
| POST | `/api/tasks` | `gateway.tasks.api_tasks_create` |
| GET | `/api/tasks` | `gateway.tasks.api_tasks_list` |
| GET | `/api/tasks/{task_id}` | `gateway.tasks.api_task_get` |
| GET | `/api/tasks/{task_id}/events` | `gateway.tasks.api_task_events` |
| GET | `/api/tasks/{task_id}/events` | `gateway.tasks.api_task_events` (legacy integer rank) |
| POST | `/api/tasks/{task_id}/events` | `gateway.tasks.api_task_events` (read-only v2 cursor) |
| GET | `/api/tasks/{task_id}/artifacts/{name}` | `gateway.tasks.api_task_artifact` |
| POST | `/api/tasks/{task_id}/cancel` | `gateway.tasks.api_task_cancel` |
| POST | `/api/tasks/{task_id}/hurry` | `gateway.tasks.api_task_hurry` |
@ -931,7 +934,7 @@ The browser reconnects with bounded exponential delay, shows the reconnect overl
Each Chat instance handles `open` by resynchronizing archive-aware durable history and `close` by withdrawing online/accounting presentation; reconnect deduplication covers overlap between live frames and REST replay, Logs merges the same way, and large history parsing runs off the server event loop. Delivery is live plus replay, not a promise that every transient frame is persisted: durable chat rows, task results, queue snapshots, Project revisions, review ledgers, cost ledgers, and lifecycle state remain the recovery authorities.
## 5. Supervisor Loop
`server.py::_run_supervisor()` is the single scheduler for pooled tasks. A healthy tick publishes liveness, rotates the paired chat and progress logs, checks worker health, drains worker, direct-chat, and consciousness events, accepts owner bridge input, enforces deadlines and schedules, runs throttled reconciliation and evolution admission, assigns eligible work, and persists `state/queue_snapshot.json`. Bridge intake precedes timeout, maintenance, evolution, and assignment work so a slow control-plane step cannot make a new owner message invisible. Three consecutive loop failures clear supervisor readiness, stop its watchdog generation, and notify the owner instead of leaving a healthy-looking server that no longer assigns work; a failure raised while a shutdown or restart is already in progress (the lifespan teardown sets a process-local stop event first and joins the loop for a bounded window before workers, bridge and event bus go down) is not a crash — the loop exits quietly, without the counter, the error, or the alarm — and the crash backoff waits on that stop event so a shutdown is never held by it.
`server.py::_run_supervisor()` is the single scheduler for pooled tasks. A healthy tick publishes liveness, rotates the runtime logs, checks worker health, drains worker, direct-chat, and consciousness events, accepts owner bridge input, enforces deadlines and schedules, runs throttled reconciliation and evolution admission, assigns eligible work, and persists `state/queue_snapshot.json`. Bridge intake precedes timeout, maintenance, evolution, and assignment work so a slow control-plane step cannot make a new owner message invisible. Three consecutive loop failures clear supervisor readiness, stop its watchdog generation, and notify the owner instead of leaving a healthy-looking server that no longer assigns work; a failure raised while a shutdown or restart is already in progress (the lifespan teardown sets a process-local stop event first and joins the loop for a bounded window before workers, bridge and event bus go down) is not a crash — the loop exits quietly, without the counter, the error, or the alarm — and the crash backoff waits on that stop event so a shutdown is never held by it.
`PENDING` and `RUNNING`, guarded by `supervisor.queue._queue_lock`, are the live task-lifecycle authority. Admission reserves identity before project, workspace, attachment, or routing side effects can create a duplicate; refuses a disabled pool, duplicate task, project deletion, accepted or sealed root, or exhausted root budget; attaches the task contract; and preserves stable priority order. Assignment runs against the same locked state and skips reaping slots, budget-paused work, closed project roots, conflicting project writers, and tasks exceeding the root's subagent capacity or depth reservation; evolution tasks are dropped there when `evolution_block_reason()` is set (Light runtime mode, `supervisor/workers.py`). That is the last of three evolution-only runtime-mode fences: owner and post-task entry points refuse a campaign start, `enqueue_evolution_task_if_needed()` independently pauses and disables a carried campaign before queueing it, and assignment drops what still slipped through (`supervisor/evolution_lifecycle.py`). Generic `supervisor.queue.enqueue_task()` has no runtime-mode predicate at all. Configured worker count is therefore not available capacity: the truthful value is the currently assignable idle count after custody, reaping, and admission fences.
@ -968,7 +971,7 @@ Startup and throttled maintenance reconcile three residue classes. Process custo
Cooperative project checkpointing has two equivalent quiescence triggers: a host-minted genesis or cooperative tree is checked when its root settles with no live descendants, and again when the last child settles beneath an already-terminal root. The second trigger exists because a root-scope budget stop terminalizes the root before its children reach their own dispatch boundaries; a root-only trigger would see a live tree once and never return. Event dispatch detects the condition only after removing the finishing task from RUNNING. The bounded git chain runs on a daemon thread, revalidates quiescence under the queue lock immediately before mutation, and uses a per-root latch that replays a trigger arriving during an in-flight check. Only host-minted project roots are eligible; owner-attached folders are never auto-committed, credential-shaped files stay excluded and disclosed, and every material success, skip, or error receives a durable receipt.
The bridge recognizes `/panic`, `/restart`, `/review`, `/evolve [on|off]`, `/bg [start|stop|status]`, and `/status`; all other text enters ordinary agent routing. External transports may invoke these commands only with positive owner identity and a transport-specific owner-chat binding. The commands reuse runtime-mode, queue, cancellation, and typed-result authority rather than implementing parallel control paths. Chat and progress logs rotate on the same supervisor tick and archive readers preserve their joint timeline. Only explicitly isolated devtool roots may use the narrow rotation sentinel from §1; normal runtime roots never inherit it.
The bridge recognizes `/panic`, `/restart`, `/review`, `/evolve [on|off]`, `/bg [start|stop|status]`, and `/status`; all other text enters ordinary agent routing. External transports may invoke these commands only with positive owner identity and a transport-specific owner-chat binding. The commands reuse runtime-mode, queue, cancellation, and typed-result authority rather than implementing parallel control paths. Runtime logs rotate on the same supervisor tick and archive readers preserve their retained timelines. Only explicitly isolated devtool roots may use the narrow rotation sentinel from §1; normal runtime roots never inherit it.
## 6. Agent Core

View file

@ -475,6 +475,12 @@ filtered down to the answer.
unbounded event log (`ouroboros/delegate_custody.py`); the fingerprint-keyed
render cache in `ouroboros/_usage_rows_memo.py` — a projection cached while
its input is unchanged, invalidated only by advance/refold, never by TTL.
Interactive result discovery reuses `gateway/task_list_scan.py`'s compact
stat-invalidated facts; decode changed/new files, never cache failed or torn
reads, and fetch selected full rows through the existing schema owner.
Task-event v2 cursors advance after each consumed row; preserve physical byte
positions through rotation and disclose replay/gaps rather than resetting an
unavailable cursor. Keep the legacy GET rank contract separate.
Enforcement: Repo Commit Checklist item 24 (advisory) triggers on diffs that
add or change an endpoint/poller/subscription/timer or read a growing store;

View file

@ -13,6 +13,7 @@ import time
import urllib.error
import urllib.parse
import urllib.request
from collections import OrderedDict
from typing import Any, Dict, Iterable, Iterator, List, Optional
@ -98,8 +99,13 @@ class OuroborosHTTPClient:
except urllib.error.URLError as exc:
raise ConnectionCLIError(f"cannot download from Ouroboros server at {self.base_url}: {exc}") from exc
def stream_sse(self, path: str, timeout: float = 120.0) -> Iterator[Dict[str, Any]]:
req = urllib.request.Request(self.base_url + path, headers={"Accept": "text/event-stream"})
def stream_sse(self, path: str, timeout: float = 120.0, *, body: Optional[dict] = None) -> Iterator[Dict[str, Any]]:
headers = {"Accept": "text/event-stream"}
if body is not None:
headers["Content-Type"] = "application/json"
req = urllib.request.Request(self.base_url + path, headers=headers,
data=json.dumps(body).encode("utf-8") if body is not None else None,
method="POST" if body is not None else "GET")
try:
with urllib.request.urlopen(req, timeout=timeout) as resp:
yield from _parse_sse_lines(resp)
@ -765,7 +771,13 @@ def _watch_task(
quiet: bool,
timeout_sec: float,
) -> None:
cursor = 0
cursor: Optional[dict] = None
legacy_seq = 0
legacy = False
first_connection = True
# Bounded overlap suppression across explicit view replays. Older repeats
# outside this recent window may reappear; no infinite exactly-once claim.
seen: OrderedDict[str, None] = OrderedDict()
final = False
deadline = time.time() + timeout_sec if timeout_sec and timeout_sec > 0 else None
while not final:
@ -777,21 +789,49 @@ def _watch_task(
remaining = max(0.0, deadline - time.time())
wait_param = max(0, min(30, int(remaining)))
request_timeout = max(1.0, min(40.0, remaining + 1.0))
path = f"/api/tasks/{urllib.parse.quote(task_id)}/events?cursor={cursor}&wait={wait_param}"
path = f"/api/tasks/{urllib.parse.quote(task_id)}/events"
saw_event = False
for event in client.stream_sse(path, timeout=request_timeout):
if deadline is not None and time.time() >= deadline:
raise TaskTimeoutCLIError(f"task {task_id} did not finish within {timeout_sec:g}s")
saw_event = True
cursor = max(cursor, int(event.get("seq") or cursor))
if jsonl:
print(json.dumps(event, ensure_ascii=False), flush=True)
elif not quiet:
rendered = _render_event_for_stderr(event)
if rendered:
print(rendered, file=sys.stderr, flush=True)
if event.get("type") == "task_result":
final = _is_terminal_result((event.get("data") or {}))
try:
events = (client.stream_sse(f"{path}?cursor={legacy_seq}&wait={wait_param}", timeout=request_timeout)
if legacy else client.stream_sse(path, timeout=request_timeout,
body={"v": 2, "wait": wait_param, "cursor": cursor}))
for event in events:
if deadline is not None and time.time() >= deadline:
raise TaskTimeoutCLIError(f"task {task_id} did not finish within {timeout_sec:g}s")
saw_event = True
if event.get("type") == "error":
raise CLIError(str(event.get("error") or "task stream failed"))
identity = str(event.get("event_id") or "")
duplicate = bool(identity) and identity in seen
if not duplicate:
if jsonl:
print(json.dumps(event, ensure_ascii=False), flush=True)
elif not quiet:
rendered = _render_event_for_stderr(event)
if rendered:
print(rendered, file=sys.stderr, flush=True)
if event.get("type") == "task_result":
final = _is_terminal_result((event.get("data") or {}))
if identity:
seen[identity] = None
seen.move_to_end(identity)
if len(seen) > 4096:
seen.popitem(last=False)
# Advance only after the complete envelope was consumed. A
# result snapshot has no log identity and is never deduplicated.
if not legacy:
next_cursor = event.get("cursor")
if not isinstance(next_cursor, dict):
raise CLIError("v2 task event is missing its cursor")
cursor = next_cursor
legacy_seq = max(legacy_seq, int(event.get("seq") or legacy_seq))
except CLIError as exc:
if (first_connection and not saw_event and not legacy
and isinstance(exc.__cause__, urllib.error.HTTPError) and exc.__cause__.code == 405):
legacy, first_connection = True, False
continue
raise
first_connection = False
if not saw_event:
time.sleep(0.5)
@ -848,6 +888,8 @@ def _render_event_for_stderr(event: Dict[str, Any]) -> str:
return f"tool: {data.get('tool', '?')}"
if etype in {"task_done", "task_metrics"}:
return f"{etype}: {data.get('task_id', '')}"
if etype in {"cursor_replay", "history_gap"}:
return f"{etype}: {event.get('reason', '')}"
return ""

View file

@ -1237,6 +1237,19 @@ class ClaudexorLoginJobProblem(TypedDict, total=False):
required_actions: List[str]
class TaskEventCursor(TypedDict):
v: Literal[2]
seq: int
view: str
positions: Dict[str, Dict[str, int]]
class TaskEventsRequest(TypedDict):
v: Literal[2]
wait: NotRequired[int]
cursor: NotRequired[Optional[TaskEventCursor]]
class TaskEvent(TypedDict, total=False):
seq: int
source: str
@ -1246,6 +1259,10 @@ class TaskEvent(TypedDict, total=False):
task_id: str
root: str
data: Dict[str, Any]
event_id: str
cursor: TaskEventCursor
reason: str
error: str
class TaskCancelResponse(TypedDict, total=False):
@ -1571,6 +1588,8 @@ __all__ = [
"ClaudexorVendorCredentialDisposition",
"ClaudexorCredentialProfileDeleteResponse",
"TaskEvent",
"TaskEventCursor",
"TaskEventsRequest",
"TaskCancelResponse",
"TaskHurryRequest",
"TaskHurryResponse",

View file

@ -29,6 +29,7 @@ HTTP_ENDPOINTS: tuple[str, ...] = (
"GET /api/tasks/{task_id}",
"GET /api/tasks/{task_id}/artifacts/{name}",
"GET /api/tasks/{task_id}/events",
"POST /api/tasks/{task_id}/events",
"POST /api/tasks/{task_id}/cancel",
"POST /api/tasks/{task_id}/hurry",
"POST /api/tasks/{task_id}/resume",

View file

@ -73,14 +73,9 @@ async def api_logs_tail(request: Request) -> JSONResponse:
rows: List[Dict[str, Any]] = []
for root in roots:
path = pathlib.Path(root) / "logs" / filename
# Bounded tail (v6.90.x P2, events/tools added v6.109.29): a window-doubled
# byte tail of the live file plus a newest-first archive/<name>_*.jsonl
# backfill until `limit` MATCHING rows are collected. chat/progress/events/
# tools rotate on the supervisor tick (see _ROTATED_LOG_PREFIXES for the
# SSE heal); supervisor.jsonl does NOT rotate, so its archive glob is empty.
# `_line` is the entry's position
# within the read window (archive chain -> live tail) — a chronological
# tie-break within one root, no longer an absolute file line number.
# A bounded live byte tail plus newest-first archive backfill until
# the filtered quota is met. All runtime sources below rotate.
# ``_line`` is a position within this read window, not the whole file.
entries = read_rotated_jsonl_entries(
path, pathlib.Path(root) / "archive", name, limit, _matches_task_filter
)

View file

@ -238,7 +238,7 @@ def collect_routes(
Route("/api/tasks", endpoint=api_tasks_list, methods=["GET"]),
Route("/api/tasks/{task_id}/artifacts/{name}", endpoint=api_task_artifact, methods=["GET"]),
Route("/api/tasks/{task_id}", endpoint=api_task_get, methods=["GET"]),
Route("/api/tasks/{task_id}/events", endpoint=api_task_events, methods=["GET"]),
Route("/api/tasks/{task_id}/events", endpoint=api_task_events, methods=["GET", "POST"]),
Route("/api/tasks/{task_id}/cancel", endpoint=api_task_cancel, methods=["POST"]),
Route("/api/tasks/{task_id}/hurry", endpoint=api_task_hurry, methods=["POST"]),
Route("/api/tasks/{task_id}/resume", endpoint=api_task_resume, methods=["POST"]),

View file

@ -1,15 +1,14 @@
"""Task-event SSE endpoint and follow state (split out of gateway/tasks.py).
"""Task-event SSE: legacy sorted replay and additive physical-cursor streaming.
Extracted verbatim from ``ouroboros/gateway/tasks.py`` when the v6.9x P2
slice/discovery work pushed that module past the 1600-line size gate
(tests/test_smoke.py::test_no_oversized_modules). ``gateway/tasks.py``
re-exports this module's surface, so route wiring, CLI imports, and
monkeypatch pins keep addressing ``ouroboros.gateway.tasks`` unchanged.
The existing ``gateway.tasks`` exports remain the route and injection surface.
Both transports share current result/lineage projections and all five sources;
v2 owns only read positions, never event persistence or task state.
"""
from __future__ import annotations
import asyncio
import hashlib
import json
import os
import pathlib
@ -17,13 +16,15 @@ import time
from typing import Any, Dict, List, Optional
from starlette.requests import Request
from starlette.responses import StreamingResponse
from starlette.responses import JSONResponse, StreamingResponse
from ouroboros.gateway._helpers import coerce_int, request_drive_root
from ouroboros.headless import ARTIFACT_STATUS_FINALIZING, ARTIFACT_STATUS_PENDING
from ouroboros.outcomes import public_task_result
from ouroboros.task_results import load_task_result, task_results_dir, validate_task_id
from ouroboros.task_status import FINAL_STATUSES
from ouroboros.gateway.task_list_scan import raw_result_facts
from ouroboros.utils import jsonl_chain_handles
def _tasks_namespace():
@ -61,6 +62,8 @@ async def api_task_events(request: Request) -> StreamingResponse:
async def _bad_id():
yield _sse({"type": "error", "error": message, "seq": 1}, event_id=1)
return StreamingResponse(_bad_id(), media_type="text/event-stream", status_code=400)
if getattr(request, "method", "GET") == "POST":
return await _api_task_events_v2(request, task_id)
cursor = max(0, coerce_int(request.query_params.get("cursor"), 0))
wait_sec = max(0, min(coerce_int(request.query_params.get("wait"), 30), 120))
drive_root = request_drive_root(request)
@ -171,9 +174,8 @@ async def api_task_events(request: Request) -> StreamingResponse:
# Live logs that the supervisor rotates into archive/<prefix>_<ts>.jsonl
# (supervisor/state.rotate_jsonl_log_if_needed). supervisor / task_reflections /
# chat_annotations do NOT rotate and never need archive-chain healing here.
_ROTATED_LOG_PREFIXES = {"progress": "progress", "chat": "chat", "events": "events", "tools": "tools"}
# (supervisor/state.rotate_jsonl_log_if_needed). Every source served here rotates.
_ROTATED_LOG_PREFIXES = {source: source for source, _parts in _LOG_SOURCES}
def _event_sort_key(item: Dict[str, Any]) -> tuple:
@ -252,16 +254,12 @@ class _TaskEventFollower:
self.suppress_task_done = False
self.filter_grew = False
self._queue_snapshot_mtime: Any = None
# Child discovery state (v6.9x P2): result filenames already read (and
# lineage-classified) by the scandir name-diff in _discover_roots.
self._results_dir = task_results_dir(self.drive_root, create=False)
self._seen_result_names: set = set()
# Archive floor (P2 review, fix 4): the RAW result's ts is the first
# write's timestamp (creation/admission — no production writer passes an
# explicit ts), so an archive whose rotation stamp predates it cannot
# contain this task's rows. Empty floor = no bound (fail open).
# Only an explicit creation fact may bound the scan. Legacy ``ts`` can
# describe finalization and must never hide earlier task history.
raw = load_task_result(self.drive_root, task_id) or {}
self._created_floor = _compact_ts_stamp(str(raw.get("created_at") or raw.get("ts") or ""))
self._created_floor = _compact_ts_stamp(str(raw.get("created_at") or ""))
def refresh_result(self) -> None:
self.result = _tasks_namespace().load_effective_task_result(
@ -284,78 +282,35 @@ class _TaskEventFollower:
return changed
def _discover_roots(self) -> bool:
"""Refresh roots + lineage filter ids; True when something new appeared.
"""Rebuild the filter from current compact facts, reusing unchanged rows.
Child discovery is a scandir NAME-DIFF over the main root's
task_results/ (v6.9x P2). Invariant this relies on: schedule_subagent
durably writes the child's ``task_results/<tid>.json`` — already
carrying lineage and child_drive_root — into the MAIN data root BEFORE
emitting any event and BEFORE the child is enqueued
(ouroboros/tools/control.py, the STATUS_REQUESTED write; a failed write
means the child was never scheduled). A new child is therefore always
visible as a new filename no later than its first log row. Each tick
reads ONLY names outside the seen-set; a name is committed to the
seen-set ONLY after read_json_dict succeeds (a torn/mid-write file is
retried next tick), and successfully-read NON-lineage names are
committed too so a busy shared store is not re-read every tick. The
lineage match reproduces find_child_tasks' subtree semantics exactly:
direct parent OR root equals the watched id, delegation_role ==
"subagent", child_drive_root collection, NO recursion (mid-stream
grandchildren of a non-root watched task did not match before either).
A child missed through a transient failure is recovered by the next
tick or, at worst, a client reconnect's full re-merge (at-least-once —
pre-existing property)."""
changed = False
Results are written with lineage before enqueue. Stat invalidation also
sees later changes to an existing row; reconnects share that same memo.
Discovery is a read-only projection, never result/schema authority.
"""
candidates = [self.drive_root]
child = str(
self.result.get("child_drive_root")
or self.result.get("headless_child_drive_root")
or ""
).strip()
child = str(self.result.get("child_drive_root") or self.result.get("headless_child_drive_root") or "").strip()
if child:
candidates.append(pathlib.Path(child))
try:
with os.scandir(self._results_dir) as entries:
names = [entry.name for entry in entries if entry.name.endswith(".json")]
except OSError:
names = []
for name in names:
if name in self._seen_result_names:
facts, _malformed = raw_result_facts(
self._results_dir, reader=_tasks_namespace().read_json_dict,
)
self._seen_result_names = {name for name, row in facts.items() if not row["schema_refusal"]}
ids = {self.task_id}
for row in facts.values():
if row["schema_refusal"] or row["delegation_role"] != "subagent":
continue
row = _tasks_namespace().read_json_dict(self._results_dir / name)
if row is None:
continue # torn write: not committed, re-read next tick
self._seen_result_names.add(name)
if str(row.get("delegation_role") or "") != "subagent":
child_id = row["task_id"] or row["id"]
if not child_id or not (row["parent_task_id"] == self.task_id or row["root_task_id"] == self.task_id):
continue
child_id = str(row.get("task_id") or row.get("id") or "").strip()
if not child_id:
continue
if not (
str(row.get("parent_task_id") or "") == self.task_id
or str(row.get("root_task_id") or "") == self.task_id
):
continue
if child_id not in self.task_filter_ids:
self.task_filter_ids.add(child_id)
changed = True
# A new FILTER ID over already-consumed bytes is lossy: rows
# matching only via subagent_task_id were filtered out when
# those bytes were read, so only a full re-merge recovers them
# (new ROOTS are fine — their logs join at offset 0). The
# stream checks this flag after every poll; full_merge resets it.
self.filter_grew = True
child_root = str(
row.get("child_drive_root")
or row.get("headless_child_drive_root")
or ""
).strip()
ids.add(child_id)
child_root = row["child_drive_root"] or row["headless_child_drive_root"]
if child_root:
candidates.append(pathlib.Path(child_root))
for path in candidates:
if path not in self.roots:
self.roots.append(path)
changed = True
roots = sorted(set(candidates), key=str)
changed = ids != self.task_filter_ids or roots != self.roots
self.filter_grew = self.filter_grew or changed
self.task_filter_ids, self.roots = ids, roots
return changed
def _log_state(self, root: pathlib.Path, source: str) -> Dict[str, Any]:
@ -463,8 +418,7 @@ class _TaskEventFollower:
self.logs = {}
self.roots = []
self.task_filter_ids = {self.task_id}
# Reset the discovery baseline too: the merge below re-reads every
# consumed byte, so every result name must be re-read and re-classified.
# The process-local stat memo survives a connection/replay rebuild.
self._seen_result_names = set()
self.refresh_result()
self.queue_snapshot_changed()
@ -503,6 +457,200 @@ class _TaskEventFollower:
return rows, advanced
def _cursor_input(value: Any) -> Dict[str, Any]:
"""Validate only the versioned transport; paths never select read authority."""
if value is None:
return {"v": 2, "seq": 0, "view": "", "positions": {}}
if not isinstance(value, dict) or set(value) != {"v", "seq", "view", "positions"}:
raise ValueError("cursor must contain v, seq, view and positions")
if value["v"] != 2 or type(value["seq"]) is not int or value["seq"] < 0:
raise ValueError("cursor version or sequence is invalid")
if not isinstance(value["view"], str) or not value["view"] or not isinstance(value["positions"], dict):
raise ValueError("cursor view or positions is invalid")
for root, sources in value["positions"].items():
if not isinstance(root, str) or not isinstance(sources, dict):
raise ValueError("cursor root positions must be objects")
if any(source not in _ROTATED_LOG_PREFIXES or type(offset) is not int or offset < 0
for source, offset in sources.items()):
raise ValueError("cursor source position is invalid")
return value
class _TaskEventCursorFollower(_TaskEventFollower):
"""Physical append positions over immutable archives plus the live file.
Positions count the complete chain, independently of creation-floor skips.
No timestamp sorting or history-sized event batches occur on this path.
Archives must remain immutable and retained: a shorter/unreadable chain
refuses continuation. Manual prefix removal masked by new appended bytes
cannot be detected by this offset protocol and is outside its guarantee.
"""
def __init__(self, drive_root: pathlib.Path, task_id: str, cursor: dict) -> None:
super().__init__(drive_root, task_id)
self.seq = cursor["seq"]
self.view = cursor["view"]
self.positions = {root: dict(sources) for root, sources in cursor["positions"].items()}
def checkpoint(self) -> dict:
return {"v": 2, "seq": self.seq, "view": self.view,
"positions": {root: dict(sources) for root, sources in self.positions.items()}}
def envelope(self, event: dict, *, identity: str = "", delivery: bool = True) -> dict:
if delivery:
self.seq += 1
return {**event, "seq": self.seq, "event_id": identity, "cursor": self.checkpoint()}
def refresh_view(self) -> Optional[dict]:
self.refresh_result()
self._discover_roots()
facts = {"task_id": self.task_id, "task_ids": sorted(self.task_filter_ids),
"roots": [str(root) for root in self.roots],
"suppress_task_done": self.suppress_task_done, "creation_floor": self._created_floor}
view = hashlib.sha256(json.dumps(facts, sort_keys=True).encode()).hexdigest()
expected = {str(root): {source: 0 for source, _parts in _LOG_SOURCES} for root in self.roots}
changed = bool(self.view) and view != self.view
if not self.view or changed:
self.positions = expected
elif (set(self.positions) != set(expected)
or any(set(self.positions[root]) != set(sources) for root, sources in expected.items())):
raise ValueError("cursor positions do not cover this view")
self.view = view
if changed:
return self.envelope({"type": "cursor_replay", "task_id": self.task_id,
"reason": "view_changed"}, delivery=False)
return None
def read_events(self):
"""Yield one row and its own post-row checkpoint, in root/source order."""
for root in self.roots:
for source, parts in _LOG_SOURCES:
live = root.joinpath(*parts)
position = self.positions[str(root)][source]
with jsonl_chain_handles(live, strict=True) as handles:
segments = []
total = 0
for path, handle in handles:
size = os.fstat(handle.fileno()).st_size
segments.append((path, handle, total, size))
total += size
if position > total:
raise ValueError(f"cursor_unavailable: shortened or missing {source} chain at {root}")
for path, handle, start, size in segments:
if position >= start + size:
continue
if (path != live and self._created_floor
and _archive_stamp_predates(path.name, source, self._created_floor)):
position = start + size
self.positions[str(root)][source] = position
continue
handle.seek(max(0, position - start))
while handle.tell() < size:
row_start = start + handle.tell()
raw = handle.readline(size - handle.tell())
if not raw:
break
if not raw.endswith(b"\n") and path == live:
break # an append still owns this unfinished row
position = start + handle.tell()
self.positions[str(root)][source] = position
identity = json.dumps([str(root), source, row_start], separators=(",", ":"))
try:
if not raw.endswith(b"\n"):
raise ValueError("incomplete_archive_line")
if not raw.strip():
continue
entry = json.loads(raw.decode("utf-8"))
if not isinstance(entry, dict):
raise ValueError("non_object_line")
except (ValueError, UnicodeDecodeError):
yield self.envelope({"type": "history_gap", "source": source,
"root": str(root), "task_id": self.task_id,
"reason": "invalid_archive_line" if path != live else "invalid_jsonl_line"},
identity=identity, delivery=False)
continue
if (str(entry.get("task_id") or "") not in self.task_filter_ids
and str(entry.get("subagent_task_id") or "") not in self.task_filter_ids
and str(entry.get("parent_task_id") or "") != self.task_id
and str(entry.get("root_task_id") or "") != self.task_id):
continue
event = _event_from_log_entry(source, 0, entry, root)
# Legacy ``line`` counts parsed rows. A resumed byte
# reader cannot reconstruct that count without replay.
event.pop("line")
if self.suppress_task_done and event["type"] == "task_done":
continue
yield self.envelope(event, identity=identity)
async def _api_task_events_v2(request: Request, task_id: str):
try:
body = await request.json()
if not isinstance(body, dict) or body.get("v") != 2 or set(body) - {"v", "wait", "cursor"}:
raise ValueError("expected {v: 2, wait?, cursor?}")
cursor = _cursor_input(body.get("cursor"))
wait = body.get("wait", 30)
if type(wait) is not int or not 0 <= wait <= 120:
raise ValueError("wait must be an integer from 0 to 120 seconds")
except (ValueError, TypeError) as exc:
return JSONResponse({"error": str(exc)}, status_code=400)
drive_root = request_drive_root(request)
if not load_task_result(drive_root, task_id):
return JSONResponse({"error": "task not found", "task_id": task_id}, status_code=404)
follower = _TaskEventCursorFollower(drive_root, task_id, cursor)
async def stream():
deadline = time.monotonic() + wait
first = True
try:
while True:
replay = await asyncio.to_thread(follower.refresh_view)
if replay:
yield _sse(replay, event_id=follower.seq)
rows = follower.read_events()
try:
while True:
read = asyncio.create_task(asyncio.to_thread(next, rows, None))
try:
event = await asyncio.shield(read)
except asyncio.CancelledError:
# A cancelled HTTP waiter does not stop its file
# read thread. Settle it before closing the handles.
await read
raise
if event is None:
break
yield _sse(event, event_id=follower.seq)
finally:
rows.close()
terminal = follower.result_is_final()
if first or terminal:
result = follower.result
if terminal:
result = await asyncio.to_thread(
_tasks_namespace().load_effective_task_result, drive_root, task_id,
)
event = follower.envelope({"source": "task_result", "type": "task_result",
"task_id": task_id, "data": public_task_result(result)})
# Synthetic result snapshots have no log identity. They are
# always consumed, including the fresh terminal materialization.
yield _sse(event, event_id=follower.seq)
first = False
if terminal:
break
if time.monotonic() >= deadline:
event = follower.envelope({"type": "cursor_checkpoint", "task_id": task_id}, delivery=False)
yield _sse(event, event_id=follower.seq)
break
await asyncio.sleep(0.5)
except (OSError, ValueError) as exc:
event = follower.envelope({"type": "error", "task_id": task_id,
"error": str(exc), "reason": "cursor_unavailable"}, delivery=False)
yield _sse(event, event_id=follower.seq)
return StreamingResponse(stream(), media_type="text/event-stream")
def iter_task_events(drive_root: pathlib.Path, task_id: str) -> List[Dict[str, Any]]:
"""Return synthesized replayable events for a task from existing logs.

View file

@ -1,10 +1,8 @@
"""Raw result-name scan for the unfiltered GET /api/tasks slice path.
"""Compact, stat-invalidated result facts for interactive read projections.
Split out of ``ouroboros/gateway/tasks.py`` at its module-size ceiling: one
coherent concern — the creation-ts sort scan behind slice-before-projection
(v6.9x P2) plus the ABI-2 admission routing for candidates whose bytes fail
to parse. ``tasks.py`` re-imports these names (same objects), so the endpoint
wiring and tests keep their historical surface.
The files and their schema readers remain authoritative. This process-local
memo serves name ordering, SSE lineage discovery and Main's newest-result
selection; selected full results still pass the existing admission reader.
"""
from __future__ import annotations
@ -19,21 +17,71 @@ from ouroboros.task_result_schema import (
)
from ouroboros.utils import read_json_dict
# Process-wide {(results_dir, filename) -> raw ts} memo for the unfiltered list
# path. The raw `ts` is CREATION-STABLE (write_task_result sets it on the first
# write; later updates touch only updated_at), so entries never need
# invalidation — only deletions are dropped and new names decoded. Keyed by the
# directory too, so multiple drive roots (tests, child drives) never collide.
# Concurrency note: worst case a race re-reads a file and stores the identical
# creation-stable value; no lock needed.
_RAW_TS_MEMO: Dict[tuple, str] = {}
# Never retain bodies: a cached row only chooses which authoritative files to
# read. Immutable tuple values publish atomically; concurrent scans may repeat
# a read, while the next stat invalidates a superseded observation.
_RAW_TS_MEMO: Dict[tuple, tuple] = {}
_RESULT_FACT_KEYS = (
"task_id", "id", "ts", "updated_at", "delegation_role", "parent_task_id",
"root_task_id", "child_drive_root", "headless_child_drive_root",
)
def _result_stat(path: pathlib.Path) -> tuple:
stat = path.stat()
return (stat.st_dev, stat.st_ino, stat.st_size, stat.st_mtime_ns, stat.st_ctime_ns)
def raw_result_facts(results_dir: pathlib.Path, *, reader=None) -> tuple[Dict[str, dict], List[str]]:
"""Read changed/new files only, never caching failed or concurrent reads.
Parseable inadmissible rows retain their refusal for callers to apply their
own schema-reader contract; the memo never admits or quarantines a row.
``reader`` keeps the legacy gateway.tasks read seam injectable.
"""
reader = reader or read_json_dict
try:
with os.scandir(results_dir) as entries:
names = sorted(entry.name for entry in entries if entry.name.endswith(".json"))
except FileNotFoundError:
names = []
dir_key = str(results_dir)
present = set(names)
for key in [k for k in list(_RAW_TS_MEMO) if k[0] == dir_key and k[1] not in present]:
_RAW_TS_MEMO.pop(key, None)
rows: Dict[str, dict] = {}
malformed: List[str] = []
for name in names:
key = (dir_key, name)
path = results_dir / name
try:
signature = _result_stat(path)
cached = _RAW_TS_MEMO.get(key)
if cached is not None and cached[0] == signature:
rows[name] = dict(cached[1])
continue
_RAW_TS_MEMO.pop(key, None)
data = reader(path)
if data is None or _result_stat(path) != signature:
malformed.append(name)
continue
except OSError:
_RAW_TS_MEMO.pop(key, None)
malformed.append(name)
continue
facts = {field: str(data.get(field) or "") for field in _RESULT_FACT_KEYS}
facts["schema_refusal"] = task_result_schema_refusal(data)
rows[name] = facts
if not facts["schema_refusal"]:
_RAW_TS_MEMO[key] = (signature, tuple(facts.items()))
return rows, malformed
def _raw_sorted_result_names(results_dir: pathlib.Path) -> tuple[List[str], List[str]]:
"""``(sorted_names, malformed_names)`` for the unfiltered list scan.
``sorted_names`` is every parseable result filename, newest-first by RAW
creation ts (memoized); a row whose file lacks `ts` sorts as
raw ts (stat-invalidated); a row whose file lacks `ts` sorts as
minus-infinity (oldest), tie-broken by filename for determinism.
``malformed_names`` is every candidate whose bytes failed to parse: ABI-2
forbids silently dropping it — the caller MUST route it through the same
@ -43,28 +91,8 @@ def _raw_sorted_result_names(results_dir: pathlib.Path) -> tuple[List[str], List
is re-read on the next request (and the quarantine primitive itself
re-checks under the row's write lock — a row a concurrent writer just
made admissible is KEPT, never moved)."""
try:
with os.scandir(results_dir) as entries:
names = [entry.name for entry in entries if entry.name.endswith(".json")]
except OSError:
return [], []
dir_key = str(results_dir)
present = set(names)
for key in [k for k in list(_RAW_TS_MEMO) if k[0] == dir_key and k[1] not in present]:
_RAW_TS_MEMO.pop(key, None)
decorated: List[tuple] = []
malformed: List[str] = []
for name in names:
key = (dir_key, name)
raw_ts = _RAW_TS_MEMO.get(key)
if raw_ts is None:
data = read_json_dict(results_dir / name)
if data is None:
malformed.append(name)
continue
raw_ts = str(data.get("ts") or "")
_RAW_TS_MEMO[key] = raw_ts
decorated.append((raw_ts, name))
rows, malformed = raw_result_facts(results_dir)
decorated = [(row["ts"], name) for name, row in rows.items()]
decorated.sort(reverse=True) # "" (no ts) sorts after every real timestamp
return [name for _ts, name in decorated], malformed

View file

@ -269,9 +269,10 @@ def _latest_project_task_result(ctx: Any, project_id: str) -> Optional[Dict[str,
def _main_routing_manifest(ctx: Any) -> Dict[str, Any]:
"""Bounded canonical facts for one Main-chat LLM routing decision."""
from ouroboros.gateway._helpers import read_rotated_jsonl_entries
from ouroboros.gateway.task_list_scan import raw_result_facts
from ouroboros.projects_registry import list_projects
from ouroboros.task_results import list_task_results
from ouroboros.utils import iter_jsonl_objects
from ouroboros.task_results import load_task_result, task_results_dir
projects = [{
"project_id": str(row.get("id") or ""),
@ -284,27 +285,36 @@ def _main_routing_manifest(ctx: Any) -> Dict[str, Any]:
} for row in list_projects(ctx.DRIVE_ROOT)]
roots = _addressable_root_tasks(ctx, None)
all_results = list_task_results(ctx.DRIVE_ROOT)
all_results.sort(key=lambda row: str(row.get("ts") or row.get("updated_at") or ""), reverse=True)
finals = [_task_result_ground_truth(row) for row in all_results[:16]]
facts, unreadable = raw_result_facts(task_results_dir(ctx.DRIVE_ROOT, create=False))
ordered = sorted(facts, key=lambda name: facts[name]["ts"] or facts[name]["updated_at"], reverse=True)
finals = []
for name in ordered:
if facts[name]["schema_refusal"]:
continue
row = load_task_result(ctx.DRIVE_ROOT, pathlib.Path(name).stem)
if row is not None:
finals.append(_task_result_ground_truth(row))
if len(finals) == 16:
break
dialogue_rows: list = []
chat_paths = sorted(
(pathlib.Path(ctx.DRIVE_ROOT) / "archive").glob("chat_*.jsonl"),
key=lambda path: path.name,
)[-2:] + [pathlib.Path(ctx.DRIVE_ROOT) / "logs" / "chat.jsonl"]
for path in chat_paths:
for row in iter_jsonl_objects(path):
text = str(row.get("text") or "").strip()
if text:
dialogue_rows.append({
"ts": str(row.get("ts") or ""),
"direction": str(row.get("direction") or ""),
"chat_id": int(row.get("chat_id") or 1),
"text": _clip_marked(text, 500),
"task_id": str(row.get("task_id") or ""),
"client_message_id": str(row.get("client_message_id") or ""),
})
root = pathlib.Path(ctx.DRIVE_ROOT)
rows, gaps = read_rotated_jsonl_entries(
root / "logs" / "chat.jsonl", root / "archive", "chat", 20,
lambda row: bool(str(row.get("text") or "").strip()),
max_archives=2, include_gaps=True,
)
for row in rows:
text = str(row.get("text") or "").strip()
if text:
dialogue_rows.append({
"ts": str(row.get("ts") or ""),
"direction": str(row.get("direction") or ""),
"chat_id": row.get("chat_id", 1),
"text": _clip_marked(text, 500),
"task_id": str(row.get("task_id") or ""),
"client_message_id": str(row.get("client_message_id") or ""),
})
dialogue = dialogue_rows[-20:]
return {
"projects": projects[:40],
@ -314,8 +324,12 @@ def _main_routing_manifest(ctx: Any) -> Dict[str, Any]:
"omissions": {
"projects": max(0, len(projects) - 40),
"root_tasks": max(0, len(roots) - 40),
"final_results": max(0, len(all_results) - 16),
"dialogue_rows": max(0, len(dialogue_rows) - 20),
"final_results": None if unreadable else max(0, len(facts) - len(finals)),
# A bounded read cannot count bytes/rows it deliberately did not
# visit. The exact historical messages remain available by id.
"dialogue_rows": None,
"dialogue_source": "chat_history (canonical chat and archives)",
"dialogue_gaps": sorted(gaps),
},
}

View file

@ -114,6 +114,7 @@ BAND_PATHS = {
"ouroboros/artifacts.py": "Exact task-source persistence and repository-diff materialization now share the actor-readable artifact seam for acceptance evidence.",
"ouroboros/cancel_intents.py": "Entered the band from 929 lines: reciprocal timeout-retry lineage validation and physical-leaf/logical-root aliasing stay with the durable cancel-intent mutation authority so Stop-now hardens the same request across retry races.",
"ouroboros/capability_evidence.py": "Grew INTO the band by the #284 fix: a fresh exact-model density witness may honestly undercut the cold floor \u2014 evidence logic belongs beside the witness store it reads.",
"ouroboros/cli.py": "The existing command-line transport keeps task-event negotiation, bounded replay deduplication and result rendering together; the additive cursor does not introduce a second CLI or task engine.",
"ouroboros/consciousness.py": "Durable Background Consciousness observation inbox and bounded truthful replay",
"ouroboros/context.py": "Entered the band from the 1501-1600 zone (1590 lines) by the v7 D03 extraction of the runtime-section fact builders into ouroboros/context_runtime_facts.py; shrink-only residue of the split, not new growth.",
"ouroboros/deep_self_review.py": "Deep self-review moved from one packed call to the three reviewer-row deliveries inside ONE surface module: the packed Atlas assembler and the retrieving runner (route-aware availability, provenance header, mandatory-read coverage, typed failures) share the memory whitelist, the prompt constants and the failure typing, so a split would separate the surface from its own pack contract; shrink next touch.",

View file

@ -46,6 +46,9 @@ from ouroboros.gateway.contracts import (
StateResponse,
TaskCostBreakdown,
TaskDetailResponse,
TaskEvent,
TaskEventCursor,
TaskEventsRequest,
TaskHurryRequest,
TaskHurryResponse,
TypingOutbound,
@ -252,6 +255,7 @@ def test_gateway_contract_endpoint_index_matches_router_and_types(tmp_path):
UpdateApplySuccessResponse, UpdateApplyErrorResponse,
UpdateStatusReadyOutbound, TaskCostBreakdown, TaskDetailResponse,
TaskHurryRequest, TaskHurryResponse, OwnerHurryProjection,
TaskEvent, TaskEventCursor, TaskEventsRequest,
OwnerSkillPresenceRuntimeRequest, OwnerSkillPresenceRuntimeResponse,
OnboardingCompleteRequest, OnboardingPresetProjection,
OnboardingSubagentsPreviewResponse, OnboardingCompleteResponse,

View file

@ -327,8 +327,8 @@ def test_cli_watch_caps_sse_wait_by_timeout(monkeypatch):
times = iter([100.0, 100.1, 100.2, 101.0])
class FakeClient:
def stream_sse(self, path, timeout=120.0):
calls.append((path, timeout))
def stream_sse(self, path, timeout=120.0, *, body=None):
calls.append((path, timeout, body))
return iter(())
monkeypatch.setattr(cli.time, "time", lambda: next(times))
@ -336,7 +336,7 @@ def test_cli_watch_caps_sse_wait_by_timeout(monkeypatch):
with pytest.raises(cli.TaskTimeoutCLIError):
cli._watch_task(FakeClient(), "abc123", jsonl=False, quiet=True, timeout_sec=0.5)
assert "wait=0" in calls[0][0]
assert calls[0][2]["wait"] == 0
assert calls[0][1] <= 1.5

View file

@ -0,0 +1,355 @@
"""Physical task-event continuation, without a second event store."""
from __future__ import annotations
import asyncio
import json
import os
from types import SimpleNamespace
import pytest
from starlette.requests import Request
from ouroboros.gateway.task_events import _TaskEventCursorFollower, api_task_events, iter_task_events
from ouroboros.gateway.task_list_scan import raw_result_facts
from ouroboros.task_results import write_task_result
def seed(root, task_id="root"):
(root / "logs").mkdir(parents=True)
write_task_result(root, task_id, "running", ts="2026-01-01T00:00:00Z")
return root
def append(root, text, *, source="progress", task_id="root", **fields):
path = root / "logs" / f"{source}.jsonl"
raw = (json.dumps({"task_id": task_id, "content": text, **fields}) + "\n").encode()
with path.open("ab") as handle:
handle.write(raw)
return len(raw)
def request(root, cursor=None, wait=0):
body = json.dumps({"v": 2, "wait": wait, "cursor": cursor}).encode()
async def receive():
return {"type": "http.request", "body": body, "more_body": False}
return Request({"type": "http", "method": "POST", "path": "/api/tasks/root/events",
"path_params": {"task_id": "root"}, "headers": [],
"app": SimpleNamespace(state=SimpleNamespace(drive_root=root))}, receive)
def stream(root, cursor=None, wait=0, on_event=None):
async def consume():
response = await api_task_events(request(root, cursor, wait))
assert response.status_code == 200
events = []
async for frame in response.body_iterator:
event = json.loads(next(line[6:] for line in frame.splitlines() if line.startswith("data: ")))
events.append(event)
if on_event:
on_event(event)
return events
return asyncio.run(consume())
def content(events):
return [row["data"]["content"] for row in events if row.get("data", {}).get("content")]
def test_reconnect_resumes_after_exact_row_and_keeps_backdated_append(tmp_path):
root = seed(tmp_path / "data")
first_bytes = append(root, "first", ts="2026-01-01T00:02:00Z")
append(root, "second", ts="2026-01-01T00:01:00Z")
first = stream(root)
assert content(first) == ["first", "second"]
first_cursor = first[0]["cursor"]
assert first_cursor["positions"][str(root)]["progress"] == first_bytes
# A cut connection after only the first frame must not acknowledge second.
resumed = stream(root, first_cursor)
assert content(resumed) == ["second"]
append(root, "backdated", ts="2025-01-01T00:00:00Z")
last = stream(root, resumed[-1]["cursor"])
assert content(last) == ["backdated"]
assert last[0]["seq"] > resumed[-1]["seq"]
assert first[1]["event_id"] == resumed[0]["event_id"]
@pytest.mark.parametrize("source", ["progress", "chat", "events", "tools", "supervisor"])
def test_every_source_survives_two_rotations_and_legacy_replay(tmp_path, source):
root = seed(tmp_path / "data")
append(root, "first", source=source)
initial = stream(root)
(root / "archive").mkdir()
live = root / "logs" / f"{source}.jsonl"
for index in [1, 2]:
append(root, f"before-{index}", source=source)
os.replace(live, root / "archive" / f"{source}_20260101T00000{index}.jsonl")
live.touch()
append(root, "last", source=source)
resumed = stream(root, initial[-1]["cursor"])
assert content(resumed) == ["before-1", "before-2", "last"]
assert content(iter_task_events(root, "root")) == ["first", "before-1", "before-2", "last"]
def test_late_child_and_changed_existing_lineage_replay_disclosed_view(tmp_path):
root = seed(tmp_path / "data")
append(root, "parent")
append(root, "before-discovery", task_id="", subagent_task_id="child")
write_task_result(root, "child", "running", delegation_role="subagent", parent_task_id="elsewhere")
first = stream(root)
assert content(first) == ["parent"]
child = tmp_path / "child"
(child / "logs").mkdir(parents=True)
append(child, "child", task_id="child")
write_task_result(root, "child", "running", parent_task_id="root", child_drive_root=str(child))
resumed = stream(root, first[-1]["cursor"])
assert resumed[0]["type"] == "cursor_replay"
assert resumed[0]["reason"] == "view_changed"
assert set(content(resumed)) == {"parent", "before-discovery", "child"}
assert next(r for r in resumed if r.get("data", {}).get("content") == "parent")["event_id"] == first[0]["event_id"]
def test_live_partial_waits_but_immutable_partial_discloses_and_advances(tmp_path):
root = seed(tmp_path / "data")
path = root / "logs" / "progress.jsonl"
partial = json.dumps({"task_id": "root", "content": "partial"}).encode()
path.write_bytes(partial)
first = stream(root)
assert content(first) == []
assert first[-1]["cursor"]["positions"][str(root)]["progress"] == 0
with path.open("ab") as handle:
handle.write(b"\n")
second = stream(root, first[-1]["cursor"])
assert content(second) == ["partial"]
with path.open("ab") as handle:
handle.write(b'{"torn":')
(root / "archive").mkdir()
os.replace(path, root / "archive" / "progress_20260101T000001.jsonl")
path.touch()
append(root, "after")
third = stream(root, second[-1]["cursor"])
assert [r["reason"] for r in third if r["type"] == "history_gap"] == ["invalid_archive_line"]
assert content(third) == ["after"]
assert content(stream(root, third[-1]["cursor"])) == []
def test_shortened_chain_refuses_without_reset(tmp_path):
root = seed(tmp_path / "data")
append(root, "first")
cursor = stream(root)[-1]["cursor"]
(root / "logs" / "progress.jsonl").write_bytes(b"")
events = stream(root, cursor)
assert len(events) == 1 and events[0]["type"] == "error"
assert events[0]["reason"] == "cursor_unavailable"
assert events[0]["cursor"]["positions"] == cursor["positions"]
def test_unreadable_chain_is_explicit(tmp_path, monkeypatch):
from ouroboros.gateway import task_events
root = seed(tmp_path / "data")
def broken(*args, **kwargs):
raise OSError("archive unreadable")
monkeypatch.setattr(task_events, "jsonl_chain_handles", broken)
events = stream(root)
assert events[-1]["type"] == "error"
assert events[-1]["reason"] == "cursor_unavailable"
def test_checkpoint_advances_unmatched_bytes_without_progress_or_sequence(tmp_path):
root = seed(tmp_path / "data")
first = stream(root)
appended = False
def after_initial(event):
nonlocal appended
if event["type"] == "task_result" and not appended:
append(root, "unrelated", task_id="other")
appended = True
events = stream(root, first[-1]["cursor"], wait=1, on_event=after_initial)
assert content(events) == []
assert events[-1]["type"] == "cursor_checkpoint"
assert events[-1]["seq"] == events[0]["seq"]
assert events[-1]["cursor"]["positions"][str(root)]["progress"] > events[0]["cursor"]["positions"][str(root)]["progress"]
assert content(stream(root, events[-1]["cursor"])) == []
def test_terminal_snapshot_materializes_once_and_legacy_ts_never_bounds_history(tmp_path, monkeypatch):
from ouroboros.gateway import tasks
root = seed(tmp_path / "data")
(root / "archive").mkdir()
(root / "archive" / "progress_20250101T000001.jsonl").write_text(
json.dumps({"task_id": "root", "content": "early"}) + "\n", encoding="utf-8")
write_task_result(root, "root", "completed", ts="2026-01-01T00:00:00Z", result="done")
real = tasks.load_effective_task_result
materialized = []
def read(*args, **kwargs):
materialized.append(kwargs.get("materialize_artifacts", True))
return real(*args, **kwargs)
monkeypatch.setattr(tasks, "load_effective_task_result", read)
events = stream(root)
assert content(events) == ["early"]
assert events[-1]["type"] == "task_result" and events[-1]["data"]["result"] == "done"
assert materialized.count(True) == 1
assert content(iter_task_events(root, "root")) == ["early"]
def test_result_memo_reconnect_and_same_size_rewrite_invalidation(tmp_path, monkeypatch):
from ouroboros.gateway import tasks
root = seed(tmp_path / "data")
write_task_result(root, "child", "running", delegation_role="subagent", parent_task_id="aaaa")
reads = []
real = tasks.read_json_dict
monkeypatch.setattr(tasks, "read_json_dict", lambda path: reads.append(path.name) or real(path))
stream(root)
reads.clear()
stream(root)
assert reads == []
path = root / "task_results" / "child.json"
before = path.stat()
path.write_bytes(path.read_bytes().replace(b'"aaaa"', b'"root"'))
os.utime(path, ns=(before.st_atime_ns, before.st_mtime_ns))
follower = _TaskEventCursorFollower(root, "root", {"v": 2, "seq": 0, "view": "", "positions": {}})
follower.refresh_view()
assert "child" in follower.task_filter_ids
assert reads == ["child.json"]
def test_result_memo_drops_deleted_and_never_caches_torn_or_changed_read(tmp_path):
root = seed(tmp_path / "data")
path = root / "task_results" / "child.json"
path.write_text('{"torn":', encoding="utf-8")
rows, errors = raw_result_facts(path.parent)
assert "child.json" in errors and "child.json" not in rows
path.unlink()
write_task_result(root, "child", "running", delegation_role="subagent", parent_task_id="root")
rows, errors = raw_result_facts(path.parent)
assert rows["child.json"]["parent_task_id"] == "root" and not errors
path.unlink()
assert "child.json" not in raw_result_facts(path.parent)[0]
def test_main_dialogue_uses_tail_and_preserves_original_ids(tmp_path, monkeypatch):
from ouroboros.gateway import _helpers
from ouroboros.server_routing_context import _main_routing_manifest
root = seed(tmp_path / "data")
path = root / "logs" / "chat.jsonl"
with path.open("w", encoding="utf-8") as handle:
for index in range(4000):
handle.write(json.dumps({"text": str(index) + "x" * 400,
"chat_id": 0, "client_message_id": f"m{index}", "task_id": "original"}) + "\n")
(root / "archive").mkdir()
(root / "archive" / "chat_20250101T000001.jsonl").write_text('{"text":"archive"}\n', encoding="utf-8")
real = _helpers.iter_jsonl_objects
reads = []
def read(path, *args, **kwargs):
reads.append((path, kwargs.get("tail_bytes")))
return real(path, *args, **kwargs)
monkeypatch.setattr(_helpers, "iter_jsonl_objects", read)
ctx = SimpleNamespace(DRIVE_ROOT=root, RUNNING={}, PENDING=[])
manifest = _main_routing_manifest(ctx)
dialogue = manifest["recent_canonical_dialogue"]
assert len(dialogue) == 20 and dialogue[0]["client_message_id"] == "m3980"
assert dialogue[-1]["chat_id"] == 0 and dialogue[-1]["task_id"] == "original"
assert len(reads) == 1 and reads[0][1] is not None
assert manifest["omissions"]["dialogue_rows"] is None
def test_suppression_change_replays_previously_filtered_done(tmp_path, monkeypatch):
from ouroboros.gateway import tasks
root = seed(tmp_path / 'data')
append(root, 'done-row', source='events', type='task_done')
result = {'task_id': 'root', 'status': 'running', 'workspace_root': '/project', 'artifact_status': 'pending'}
monkeypatch.setattr(tasks, 'load_effective_task_result', lambda *args, **kwargs: dict(result))
first = stream(root)
assert content(first) == []
result['artifact_status'] = 'ready'
second = stream(root, first[-1]['cursor'])
assert second[0]['type'] == 'cursor_replay'
assert content(second) == ['done-row']
def test_creation_floor_counts_skipped_bytes_in_cursor(tmp_path):
root = seed(tmp_path / 'data')
write_task_result(root, 'root', 'running', created_at='2026-01-01T00:01:00Z')
(root / 'archive').mkdir()
old = (json.dumps({'task_id': 'root', 'content': 'before-creation'}) + '\n').encode()
(root / 'archive' / 'progress_20250101T000001.jsonl').write_bytes(old)
recent_size = append(root, 'recent')
first = stream(root)
assert content(first) == ['recent']
assert first[-1]['cursor']['positions'][str(root)]['progress'] == len(old) + recent_size
append(root, 'later')
assert content(stream(root, first[-1]['cursor'])) == ['later']
def test_stream_disconnect_closes_open_chain_handles(tmp_path, monkeypatch):
import pathlib
root = seed(tmp_path / 'data')
append(root, 'first')
append(root, 'second')
handles = []
real = pathlib.Path.open
def opened(path, *args, **kwargs):
handle = real(path, *args, **kwargs)
if path.name == 'progress.jsonl':
handles.append(handle)
return handle
monkeypatch.setattr(pathlib.Path, 'open', opened)
async def consume_one():
response = await api_task_events(request(root))
await anext(response.body_iterator)
assert any(not handle.closed for handle in handles)
await response.body_iterator.aclose()
asyncio.run(consume_one())
assert handles and all(handle.closed for handle in handles)
def test_main_dialogue_preserves_two_archive_horizon(tmp_path):
from ouroboros.server_routing_context import _main_routing_manifest
root = seed(tmp_path / 'data')
(root / 'archive').mkdir()
for index in [1, 2, 3]:
(root / 'archive' / f'chat_20260101T00000{index}.jsonl').write_text(
json.dumps({'text': f'archive-{index}', 'client_message_id': f'm{index}', 'chat_id': -5}) + '\n',
encoding='utf-8')
(root / 'logs' / 'chat.jsonl').write_text('{"text": ""}\n', encoding='utf-8')
rows = _main_routing_manifest(SimpleNamespace(DRIVE_ROOT=root, RUNNING={}, PENDING=[]))['recent_canonical_dialogue']
assert [row['client_message_id'] for row in rows] == ['m2', 'm3']
assert [row['chat_id'] for row in rows] == [-5, -5]
def test_reconnect_gets_fresh_terminal_result_without_replaying_consumed_logs(tmp_path, monkeypatch):
from ouroboros.gateway import tasks
root = seed(tmp_path / 'data')
append(root, 'first')
first = stream(root)
assert first[0]['event_id'] and 'line' not in first[0]
assert first[-2]['data']['status'] == 'running'
write_task_result(root, 'root', 'completed', result='finished after reconnect')
real = tasks.load_effective_task_result
reads = []
def read(*args, **kwargs):
reads.append(kwargs.get('materialize_artifacts', True))
return real(*args, **kwargs)
monkeypatch.setattr(tasks, 'load_effective_task_result', read)
second = stream(root, first[-1]['cursor'])
assert content(second) == []
assert second[-1]['type'] == 'task_result'
assert second[-1]['data']['result'] == 'finished after reconnect'
assert second[-1]['data']['status'] == 'completed'
assert reads.count(True) == 1
def test_result_facts_cannot_admit_an_invalid_schema(tmp_path):
from ouroboros.task_results import load_task_result
root = seed(tmp_path / 'data')
path = root / 'task_results' / 'child.json'
raw = {'_schema_version': 999, 'task_id': 'child', 'status': 'running',
'delegation_role': 'subagent', 'parent_task_id': 'root'}
path.write_text(json.dumps(raw), encoding='utf-8')
rows, _ = raw_result_facts(path.parent)
assert rows['child.json']['schema_refusal'] == 'future_schema'
follower = _TaskEventCursorFollower(root, 'root', {'v': 2, 'seq': 0, 'view': '', 'positions': {}})
follower.refresh_view()
assert 'child' not in follower.task_filter_ids
with pytest.raises(ValueError, match='future_schema'):
load_task_result(root, 'child', strict=True)
assert json.loads(path.read_text()) == raw

View file

@ -0,0 +1,116 @@
"""CLI negotiation and consumption of additive task-event cursors."""
import io
import json
import urllib.error
import pytest
from ouroboros import cli
def cursor(seq):
return {"v": 2, "seq": seq, "view": "view", "positions": {"root": {"progress": seq}}}
def event(seq, text="progress", **fields):
return {"seq": seq, "type": "progress", "event_id": f"row-{seq}",
"cursor": cursor(seq), "data": {"content": text}, **fields}
def terminal(seq):
return event(seq, type="task_result", event_id="", data={"status": "completed", "result": "done"})
def http_error(code):
failure = cli.CLIError(f"HTTP {code}")
failure.__cause__ = urllib.error.HTTPError("http://test/events", code, "refused", {}, io.BytesIO())
return failure
def test_client_stream_posts_json_without_cursor_in_url(monkeypatch):
calls = []
frame = event(1)
def urlopen(request, timeout):
calls.append((request, timeout))
return io.BytesIO(("data: " + json.dumps(frame) + "\n\n").encode())
monkeypatch.setattr(cli.urllib.request, "urlopen", urlopen)
body = {"v": 2, "wait": 30, "cursor": cursor(0)}
assert list(cli.OuroborosHTTPClient("http://test").stream_sse("/events", body=body)) == [frame]
request, timeout = calls[0]
assert request.get_method() == "POST"
assert request.full_url == "http://test/events"
assert json.loads(request.data) == body
assert request.get_header("Content-type") == "application/json"
def test_watch_advances_consumed_checkpoint_and_deduplicates_view_replay(capsys):
calls = []
class Client:
def stream_sse(self, path, timeout, *, body):
calls.append(body)
if len(calls) == 1:
yield event(1, "once")
yield event(1, type="cursor_checkpoint", event_id="", cursor=cursor(2), data={})
else:
assert body["cursor"] == cursor(2)
yield event(2, type="cursor_replay", event_id="", data={})
yield event(3, "once", event_id="row-1")
yield event(4, "new")
yield terminal(5)
cli._watch_task(Client(), "task", jsonl=True, quiet=False, timeout_sec=0)
rows = [json.loads(line) for line in capsys.readouterr().out.splitlines()]
assert [row["data"].get("content") for row in rows if row["type"] == "progress"] == ["once", "new"]
assert rows[-1]["type"] == "task_result"
def test_watch_falls_back_only_after_first_connection_405():
calls = []
class Client:
def stream_sse(self, path, timeout, **kwargs):
calls.append((path, kwargs))
if len(calls) == 1:
raise http_error(405)
yield {"seq": 1, "type": "task_result", "data": {"status": "completed"}}
cli._watch_task(Client(), "task", jsonl=False, quiet=True, timeout_sec=0)
assert calls[0][0].endswith("/events") and calls[0][1]["body"]["v"] == 2
assert "cursor=0&wait=30" in calls[1][0] and not calls[1][1]
@pytest.mark.parametrize("code,after_frame", [(500, False), (405, True)])
def test_watch_does_not_fallback_on_other_failure_or_later_connection(code, after_frame):
calls = []
class Client:
def stream_sse(self, path, timeout, *, body):
calls.append(body)
if after_frame and len(calls) == 1:
yield event(1)
return
raise http_error(code)
with pytest.raises(cli.CLIError, match=f"HTTP {code}"):
cli._watch_task(Client(), "task", jsonl=False, quiet=True, timeout_sec=0)
assert len(calls) == (2 if after_frame else 1)
def test_watch_refuses_unavailable_cursor_without_restarting():
class Client:
calls = 0
def stream_sse(self, path, timeout, *, body):
self.calls += 1
yield event(1, type="error", error="shortened chain")
client = Client()
with pytest.raises(cli.CLIError, match="shortened chain"):
cli._watch_task(client, "task", jsonl=False, quiet=True, timeout_sec=0)
assert client.calls == 1
def test_watch_dedup_is_bounded_and_synthetic_results_are_always_consumed(capsys):
class Client:
def stream_sse(self, path, timeout, *, body):
for seq in range(1, 4098):
yield event(seq)
yield event(4098, event_id="row-1") # outside the recent window
yield terminal(4099)
cli._watch_task(Client(), "task", jsonl=True, quiet=False, timeout_sec=0)
rows = [json.loads(line) for line in capsys.readouterr().out.splitlines()]
assert len(rows) == 4099
assert rows[-2]["event_id"] == "row-1" and rows[-1]["type"] == "task_result"

View file

@ -281,8 +281,8 @@ def test_iter_task_events_reads_progress_archive_chain(tmp_path):
def test_sse_archive_floor_skips_archives_predating_task_creation(tmp_path):
"""Fix 4: an archive whose rotation stamp predates the watched task's raw
creation ts is never read (bounds the glob to the task lifetime); archives
"""An archive whose rotation stamp predates proven task creation is never
read (bounds the scan to the task lifetime); archives
stamped after creation are still consulted."""
data = tmp_path / "data"
(data / "logs").mkdir(parents=True)
@ -301,7 +301,8 @@ def test_sse_archive_floor_skips_archives_predating_task_creation(tmp_path):
json.dumps({"ts": "2026-01-01T00:02:30Z", "content": "live-row", "task_id": "t1"}) + "\n",
encoding="utf-8",
)
write_task_result(data, "t1", "running", result="working", ts="2026-01-01T00:01:00Z")
write_task_result(data, "t1", "running", result="working", ts="2026-01-01T00:01:00Z",
created_at="2026-01-01T00:01:00Z")
events = iter_task_events(data, "t1")
@ -500,6 +501,7 @@ def test_sse_torn_child_result_is_retried_and_discovered_next_tick(tmp_path):
child_drive.mkdir()
(data / "task_results" / "c1.json").write_text(
json.dumps({
"_schema_version": 1,
"task_id": "c1",
"status": "running",
"delegation_role": "subagent",
@ -528,7 +530,7 @@ def test_sse_nonlineage_result_names_read_once_then_committed(tmp_path, monkeypa
follower.full_merge()
(data / "task_results" / "other.json").write_text(
json.dumps({"task_id": "other", "status": "running", "delegation_role": "root", "ts": OLD_TS}),
json.dumps({"_schema_version": 1, "task_id": "other", "status": "running", "delegation_role": "root", "ts": OLD_TS}),
encoding="utf-8",
)
reads = []

View file

@ -936,6 +936,25 @@
* @property {string=} ts
* @property {string=} root
* @property {Object=} data
* @property {string=} event_id
* @property {TaskEventCursor=} cursor
* @property {string=} reason
* @property {string=} error
*/
/**
* @typedef {Object} TaskEventCursor
* @property {number} v
* @property {number} seq
* @property {string} view
* @property {Object<string, Object<string, number>>} positions
*/
/**
* @typedef {Object} TaskEventsRequest
* @property {number} v
* @property {number=} wait
* @property {?TaskEventCursor=} cursor
*/
/**