Bound the custody replay and the recent-activity windows per task

Every delegated-run custody question (delegate_wait, delegate_start, task
start, completion evidence, audits) re-read the whole rotated events chain
(160 segments, 329 MB on a long-lived install, ~0.6 s CPU per call). A new
process-local memo (ouroboros/delegate_custody_memo.py, the usage-memo shape)
keeps the custody rows behind delegate_custody.custody_rows: an ordered
(st_dev, st_ino, consumed, st_mtime_ns) fingerprint of the folded chain
prefix, only appended bytes folded on later reads, a refold on any doubt, a
bypass while the chain is unreadable, legacy inline request bodies replaced
by a locator delegate_pending.request_body re-reads. Every reader (replay,
run_timing, pending_invocations, invocation_record, unsettled_start_ids,
review custody recovery, payload holders, execution evidence and patch
dispositions) consumes it; replay() results are cloned per caller. Warm
reads drop from ~0.6 s to milliseconds on the same chain; the fold is
verified equal to a full replay across appends, rotations, torn tails and
chain anomalies.

The recent-activity context sections read a global 200-row tail of the
progress/tools/events logs, parsed whole, then filtered by task id, so under
many concurrent tasks a task saw whatever share of the shared suffix it
happened to occupy. They now read each task's own newest rows (progress 50,
tools 20, events 200) through the bounded rotation-aware reader, moved from
gateway/_helpers.py to the core leaf ouroboros/jsonl_tail.py (the gateway
keeps thin wrappers and its parser seam), and the header carries a coverage
line naming the window and any archives left unopened.

ARCHITECTURE and DEVELOPMENT describe both mechanisms; the events-chain
tripwire names the cold fold; chapter byte budgets are raised with reasons.

Co-authored-by: Ouroboros <311266734+ouroboros-agent@users.noreply.github.com>
This commit is contained in:
Ouroboros 2026-09-22 01:20:49 +03:00
parent 796f2e8708
commit 8aff64791d
24 changed files with 1125 additions and 137 deletions

View file

@ -224,6 +224,7 @@ server.py (Starlette+uvicorn) ← HTTP + WebSocket on configurable host:port (de
├── delegate_recovery.py ← Narrow exact-leaf recovery for proven crash + planned self-restart; vetoes every no-resume cause
├── delegate_registration_policy.py ← `persistent_registration` + the STARTED-row field tables
├── delegate_pending.py ← Durable pending-invocation replay preserving the original idempotency key + canonical start body
├── delegate_custody_memo.py ← Process-local memo of the custody rows (`custody_rows`): an ordered `(st_dev, st_ino, consumed, st_mtime_ns)` fingerprint of the rotated events chain prefix, advanced by folding only appended bytes, refolded on any doubt, bypassed (never cached) while the chain is unreadable; inline legacy request bodies replaced by a re-readable locator; a warm cache with an exact fallback, not a durable projection
├── delegate_terminal.py ← Terminal reconciliation + custody-audit persistence: counters stay a frozen snapshot while `actual_substrate` and the envelope mirror follow live custody; audit-only in both directions; the typed `terminal_custody_notice` card row; `refresh_recently_settled_terminals` over the byte-offset cursor `state/delegate_terminal_refresh_cursor.json` (5 MB per tick) (§6 Delegated subagents)
├── subagent_dispatch_notes.py ← Dispatch-time executor notes for delegated children (configured-nanny charter note); agent.py keeps re-exports
├── subagent_messages.py ← Bounded durable child-message identity shared by the final frame, recovery, persistence and replay; `executor_observation_meta` validates task-bound progress actor facts
@ -336,6 +337,7 @@ server.py (Starlette+uvicorn) ← HTTP + WebSocket on configurable host:port (de
├── browser_policy.py ← The browser tool's target and control-request policy: task-granted concrete origins, metadata/private/reserved refusals, the three-valued Ouroboros control-service identity (`runtime_service_kind`: proven kind / unknown / none); `tools/browser.py` keeps the Playwright lifecycle (§6 MCP and browser-facing external tools)
├── skill_payload_binding.py ← Skill payload targeting: `.seed-origin` distinguishes native vs external; read/list/search only for read profiles; bounded manifestless skill_publish recovery
├── utils.py ← SSOT for atomic JSON, timestamps, hashes, sanitization, subprocess helpers, `truncate_review_artifact`
├── jsonl_tail.py ← The one bounded rotation-aware filtered tail reader (window-doubling live tail, newest-first archive backfill bounded to three files, coverage facts + `coverage_line`) behind the history/logs/routing endpoints (`gateway/_helpers.py` wrappers keep the gateway parser seam) and the per-task recent-activity context sections
├── markdown_source.py ← `ouroboros/markdown_source.py`: byte-preserving Markdown structure shared by books and knowledge notes; physical LF/UTF-8 ranges tied to source SHA; `MarkdownSourceError` keeps missing grammars or malformed YAML visible without replacing original bytes
├── world_profiler.py ← Generates WORLD.md
├── contracts/ ← Frozen ABI package (§11)

View file

@ -279,7 +279,7 @@ An undisclosed spend contributes `0.0` to `accounted_usd` — inventing a conser
**Transport.** `gateways/claudexor.py` is pure transport (descriptor read, `/v2` handshake, the `config.CLAUDEXOR_MIN_VERSION` 3.2.0 floor). The daemon bearer token grants the entire `/v2` surface, so it never leaves this module, and the HTTP client runs `trust_env=False` so an ambient proxy cannot intercept the loopback control plane. Production starts obtain a handshaken owned gateway from `claudexor_daemon.ensure_owned_gateway` — exact reviewed engine/Node pins (`claudexor_runtime.py`; the reviewed pin IS the next-spawn selection — no mutable `current` pointer, no background updater — and `OUROBOROS_CLAUDEXOR_BIN` is the explicit operator override), never a PATH install. Keeping lifecycle above transport keeps account status side-effect-free and harness mechanisms out of Ouroboros (`claudexor_daemon.py`/`claudexor_runtime.py` docstrings; stop path, spawn latch and typed start failures: §9). Native child processes are contained by env token (`process_containment.py`, `OURO_PROC_CONTAINER_*`), because a surviving descendant can become invisible to parent-child traversal once its controller exits; an alive-or-undeterminable member is an honest hard-block answer, never a kill guarantee.
**Custody is durable, because the run is not ours to kill.** A delegated run lives inside the daemon and survives our worker, so the AUTHORITY is the durable `delegate_run_*` custody rows (`delegate_run_started` and friends) on the canonical/budget root (`ouroboros/delegate_custody.py`); one compact incident projection, `<drive_root>/logs/containment_faults.jsonl`, exists because the event log grows without bound and a tail-bounded scan can bury an unresolved fault. A lookup answers OWNED, FOREIGN or UNKNOWN — collapsing UNKNOWN into "not yours" made a restarted owner indistinguishable from an intruder. An ABSENT custody log is a positively established clean state; an EXISTING-but-unreadable one audits as `delegated_run_state_unknown:custody_log_unreadable`, never cleanly reconciled. Every INTENDED start mints a fresh per-intention invocation UUID as the wire `Idempotency-Key`; the content hash is only the LOOKUP identity — a content-stable wire key would hand a deliberate re-run the finished old run — and reuse happens only by explicit token (`pending_invocation_id`/`retry_of`, replaying the STORED canonical body under the SAME key). The complete replay envelope lives unredacted in the existing private observability CAS, written before `delegate_run_start_requested`; the event carries `request_ref` and `prompt_chars`, keeping a large reviewer packet out of every custody scan. `delegate_pending.request_body` resolves legacy inline first, then verifies the CAS reference for both invocation readers. JSON values and the canonical request digest stay unchanged, including the thread fields needed to reproduce its wire projection; a missing/corrupt blob leaves the request unknown while preserving pending custody identity for start blockers, terminal audits and snapshot retention; recovery retains that invocation without POSTing. Known pending review tokens take the same missing-request path, while absent or definitely refused skill-review records keep their existing fresh-retry behavior. Pending scans resolve only surviving invocations, while legacy event/archive bytes remain untouched. `reconcile_orphaned_runs` visits every open run whose owning task left the live set and settles the terminal ones, but it CANCELS only behind a deliberate verdict: a durable owner result that is readable, truly terminal and finished by the task's own decision. A custody row carries its owner's kind, and task association confers no lifecycle authority: a run a review surface registered (`RunCustody.review_owned`, durable `source` under the review substrate) belongs to its panel, bounded by the slot window and its own `maxSeconds`, so the sweep spares it unless the owner task was itself cancelled — a `left_live` row names the panel, and a task that consciously finished under a running acceptance panel keeps its reviewer alive; a pending review invocation is likewise retained, never re-posted by the generic recovery, because the review substrate owns its rejoin. The owner-cancel kill boundary STATES the verdict it is about to write, because it audits custody before that write, so an owner cancellation stops the paid run at the boundary rather than at the next sweep. A provider or transport death, a worker crash, a missing or unreadable result all SPARE the run, left live with a durable `left_live` reconciliation row: an undignified nanny death must not kill a healthy paid run, and "unknown" is exactly that case. `maxSeconds` is the damage limitation for a spared orphan, never custody. SETTLED is published before registration retirement, after the ledger row lands; settlement and registration retirement are separate durable duties, a failed retirement stays replayable on `project_owned` for the later sweep, and writing `settled` over a suppressed ledger failure would turn a lock timeout into a permanent leak. A start whose row did not land reports `started_uncustodied`: no supervision, no replacement, until the original run is proven absent or terminal.
**Custody is durable, because the run is not ours to kill.** A delegated run lives inside the daemon and survives our worker, so the AUTHORITY is the durable `delegate_run_*` custody rows (`delegate_run_started` and friends) on the canonical/budget root (`ouroboros/delegate_custody.py`), read by every consumer through `custody_rows` — the process-local memo in `delegate_custody_memo.py` (the `_usage_rows_memo` shape: an ordered `(st_dev, st_ino, consumed, st_mtime_ns)` fingerprint of the folded chain prefix, only appended bytes folded later, a refold on any doubt, a bypass and never a cache while the chain is unreadable; the rows stay the authority, each process pays one cold fold, a durable compact projection is the next step); one compact incident projection, `<drive_root>/logs/containment_faults.jsonl`, exists because the event log grows without bound and a tail-bounded scan can bury an unresolved fault. A lookup answers OWNED, FOREIGN or UNKNOWN — collapsing UNKNOWN into "not yours" made a restarted owner indistinguishable from an intruder. An ABSENT custody log is a positively established clean state; an EXISTING-but-unreadable one audits as `delegated_run_state_unknown:custody_log_unreadable`, never cleanly reconciled. Every INTENDED start mints a fresh per-intention invocation UUID as the wire `Idempotency-Key`; the content hash is only the LOOKUP identity — a content-stable wire key would hand a deliberate re-run the finished old run — and reuse happens only by explicit token (`pending_invocation_id`/`retry_of`, replaying the STORED canonical body under the SAME key). The complete replay envelope lives unredacted in the existing private observability CAS, written before `delegate_run_start_requested`; the event carries `request_ref` and `prompt_chars`, keeping a large reviewer packet out of every custody scan (the memo swaps a legacy inline body for a locator `delegate_pending.request_body` re-reads). `delegate_pending.request_body` resolves legacy inline first, then verifies the CAS reference for both invocation readers. JSON values and the canonical request digest stay unchanged, including the thread fields needed to reproduce its wire projection; a missing/corrupt blob leaves the request unknown while preserving pending custody identity for start blockers, terminal audits and snapshot retention; recovery retains that invocation without POSTing. Known pending review tokens take the same missing-request path, while absent or definitely refused skill-review records keep their existing fresh-retry behavior. Pending scans resolve only surviving invocations, while legacy event/archive bytes remain untouched. `reconcile_orphaned_runs` visits every open run whose owning task left the live set and settles the terminal ones, but it CANCELS only behind a deliberate verdict: a durable owner result that is readable, truly terminal and finished by the task's own decision. A custody row carries its owner's kind, and task association confers no lifecycle authority: a run a review surface registered (`RunCustody.review_owned`, durable `source` under the review substrate) belongs to its panel, bounded by the slot window and its own `maxSeconds`, so the sweep spares it unless the owner task was itself cancelled — a `left_live` row names the panel, and a task that consciously finished under a running acceptance panel keeps its reviewer alive; a pending review invocation is likewise retained, never re-posted by the generic recovery, because the review substrate owns its rejoin. The owner-cancel kill boundary STATES the verdict it is about to write, because it audits custody before that write, so an owner cancellation stops the paid run at the boundary rather than at the next sweep. A provider or transport death, a worker crash, a missing or unreadable result all SPARE the run, left live with a durable `left_live` reconciliation row: an undignified nanny death must not kill a healthy paid run, and "unknown" is exactly that case. `maxSeconds` is the damage limitation for a spared orphan, never custody. SETTLED is published before registration retirement, after the ledger row lands; settlement and registration retirement are separate durable duties, a failed retirement stays replayable on `project_owned` for the later sweep, and writing `settled` over a suppressed ledger failure would turn a lock timeout into a permanent leak. A start whose row did not land reports `started_uncustodied`: no supervision, no replacement, until the original run is proven absent or terminal.
**No terminal or cancel claim without a verified receipt.** `delegate_cancel` returns `confirmed` (read back terminal), `requested`, `failed` or `containment_fault_run_may_still_be_live`; the last two hold a durable CRITICAL containment fault until a receipt or settlement clears it — an overpowered run that may still be alive is an incident, not a reassuring string — and the state read decides, so a refused control is never a verdict about the RUN. One `daemon_says_absent` predicate decides everywhere that a 404 is the daemon ANSWERING that the resource is gone (scoped to the daemon that answered), never a failure to find out; custody closes such a run `delegate_run_closed_absent` (unreachable, not settled), inventing no terminal detail, usage or spend. Results are delivered, not severed: `delegate_wait` stages the whole terminal detail atomically under `task_drive/delegated_runs/<run>.json` with a typed `output_delivery` block, and cut fields are renamed `*_preview` so a partial read of head-truncated JSON fails loudly instead of looking like an answer.
@ -537,7 +537,7 @@ Every IMPLICIT claim — the UI conversion, that admission, the reaper's retry a
#### Durable memory and project focus
`context.py` assembles static governance, semi-stable memory, and dynamic task evidence without treating truncation as forgetting; the Development context matrix and `context_layout.py` own which reference form is resident. When the rendered scratchpad exceeds `SCRATCHPAD_SECTION_BUDGET_CHARS`, `context.py` keeps the newest whole blocks that fit and drops the oldest behind an in-band gap marker naming `memory/scratchpad.md` as the live source; no block is retired by a context build, and scratchpad replacement keeps its explicit summary and source-journal provenance.
`context.py` assembles static governance, semi-stable memory, and dynamic task evidence without treating truncation as forgetting; the recent-activity sections are each task's OWN newest rows (progress 50, tools 20, events 200: what the formatters render) through the bounded reader `jsonl_tail.py` (`Memory.read_task_recent`: a doubling live tail plus at most three newest archives), never a global tail filtered afterwards (issue #131), and their header's coverage line names the rows, the window and any unopened archives while `read_file` pages the rest; the Development context matrix and `context_layout.py` own which reference form is resident. When the rendered scratchpad exceeds `SCRATCHPAD_SECTION_BUDGET_CHARS`, `context.py` keeps the newest whole blocks that fit and drops the oldest behind an in-band gap marker naming `memory/scratchpad.md` as the live source; no block is retired by a context build, and scratchpad replacement keeps its explicit summary and source-journal provenance.
`consolidator.py` publishes dialogue summaries only after every part of a logical block succeeds, preserving raw chat generations and the captured generation cursor; an unfindable generation appends a loud durable `[MEMORY GAP]` block, never a silent offset reset. Each full Light request is measured with `context_fit` against fresh role/account capacity evidence and calibrated prompt density (local output reservation: `llm_local.local_context_limits`, the wire's normalization; missing or stale evidence stays unknown). Oversized source is split without clipping, including within one entry. A real context refusal requires strictly fewer input bytes on the same route: a genuine refusal immediately saves its source hash and route bound in `dialogue_meta.json` (`consolidation_retry`) so a later cycle starts smaller, and a changed source, route, capacity or output reserve invalidates that bound. Era compression remains a single aggregate request: it can still overflow, in which case the old blocks stay intact. Ordinary failures and empty summaries retain `last_consolidation_error` and an advance without a new failure clears it; a knowledge-nomination batch not fully published leaves `last_unpublished_nominations` in `dialogue_meta.json`; unknown spend stays nullable, and control/resource/unknown model errors preserve `propagate_model_error` semantics.

View file

@ -11,7 +11,7 @@ This chapter is the short list of properties the rest of the book must not contr
7. **This book is the present-tense map.** Structural owners, APIs, durable data, UI surfaces, and the rationale for non-obvious guards update in the owning chapter (entrypoint `docs/ARCHITECTURE.md`) in the same commit as the code (documentation contract: docs/DEVELOPMENT.md; residue ratchet: `tests/test_docs_sync.py`); release chronology lives in git and README.
8. **Skill gates do not collapse.** Discovery, deterministic preflight, content-hash-bound critic or qualified Advisory author authority, owner grants, dependency readiness, enablement, and execution remain separate. A PASS does not install dependencies, and `enabled=true` does not prove executable readiness.
9. **Startup rescue has one mutation owner.** Supervisor recovery writes rescue evidence before reset or blocks while preserving the tree. Worker or agent construction remains warning-only and never stages or commits inherited dirt.
10. **Projection over replay.** Interactive status, history, and cost reads are bounded, non-materializing projections; durable owners perform the one authoritative replay or terminal materialization.
10. **Projection over replay.** Interactive status, history, and cost reads are bounded, non-materializing projections; durable owners perform the one authoritative replay or terminal materialization. A process-local fingerprint memo (`_usage_rows_memo.py`, `delegate_custody_memo.py`) serves warm reads only while its store fingerprint holds, refolds on any doubt, and never touches disk.
11. **UI resources carry a disposer.** Every subscription, listener, observer, timer, stream, and live page instance has explicit teardown; navigation does not leave hidden instances mutating visible or durable state.
12. **Frozen contracts extend explicitly.** `ouroboros/contracts/` is a versioned, backward-compatible ABI — typed shapes together with their parsing/normalization/policy helpers (§11). New capability extends the frozen shape or ships an explicitly versioned successor; existing consumers keep working.
13. **Provider wire adaptation stays exact-route and success-confirmed.** Canonical history remains provider-neutral; typed physical projections may change values, fields, or a registered dialect on one provider/endpoint/API/model only. Failed candidates teach nothing durable, task-local cognition degradation never becomes future dispatch authority, and the physical-attempt ledger remains distinct from terminal request-wire history.

View file

@ -123,11 +123,15 @@ the answer.
- **House precedents — reuse these shapes:** archive-aware chat log rotation
(`supervisor/state.py::rotate_chat_log_if_needed`); the compact
`containment_faults.jsonl` projection maintained beside an unbounded event
log (`ouroboros/delegate_custody.py`); one shared custody replay per context
build and per terminal audit (`delegate_terminal.custody_audit_snapshot`,
consumed by `context_health.build_health_invariants` and the terminal
audit) — sharing ONE traversal bounds the multiplier, not the scan, so that
read stays O(history) until a compact projection replaces it; the
log (`ouroboros/delegate_custody.py`); the process-local custody row memo
behind `delegate_custody.custody_rows` (`ouroboros/delegate_custody_memo.py`:
an ordered inode/size/mtime fingerprint of the rotated chain prefix, only
appended bytes folded, a refold on any doubt, a bypass while unreadable — it
bounds the warm read, not the cold fold, so a durable compact projection
stays the next step); the bounded filtered tail reader
`ouroboros/jsonl_tail.py` (doubling live tail, three newest archives,
coverage facts) for history endpoints and the per-task recent-activity
sections alike; the
fingerprint-keyed render cache in `ouroboros/_usage_rows_memo.py`, held while
its input is unchanged and invalidated only by advance/refold, never by TTL;
the `gateway/task_list_scan.py` stat-invalidated result memo and the

View file

@ -2,7 +2,7 @@
Machine extraction of the `docs/ARCHITECTURE.md` "Data layout (`~/Ouroboros/`)" tree — the durable-file orientation carrier (this tree's counterpart of the reference PERSISTENCE_OWNERS derivation checklist) — regenerated by `python scripts/regenerate_inventories.py`. Do not edit. Every entry is probed against reality: repo entries must exist as tracked paths; data-plane entries must appear as a literal in the runtime sources that construct them. A durable file renamed or removed in code while its tree row survives = red (`tests/test_generated_inventories.py`).
Source: `docs/architecture/01-high-level-architecture.md`, physical LF lines 586-675; UTF-8 SHA-256 `c14e48d636e94f7ff143071ca4abaaf6e351cd37a97abfef15b9ae1eb5aa37e0`.
Source: `docs/architecture/01-high-level-architecture.md`, physical LF lines 588-677; UTF-8 SHA-256 `2534e9b319c114581a2c47cb226604f5cc7d2518f27580f9f9ecc7cf322217a8`.
- entries: **79** (code-ref: 72, repo-dir: 6, repo-path: 1)

View file

@ -2,7 +2,7 @@
AST-derived inventory of compatibility facades, regenerated by `python scripts/regenerate_inventories.py`. Do not edit. A facade row is any runtime module whose top-level `from <population module> import ...` statements carry the `noqa: F401` re-export marker — the codebase's declared "this binding exists for its binding, not for this module's own use" convention (reference FACADE_CONSUMERS method). Leaf domains come from `ouroboros/domains.toml`; a leaf outside the facade's domain is marked ✗ (that edge also appears in the manifest's pinned direction matrix). `tests/test_generated_inventories.py` pins byte-identity, so any re-export surface change must regenerate this file.
- facade modules: **58**; marked re-export bindings: **2292**; cross-domain facade→leaf pairs: **131**
- facade modules: **59**; marked re-export bindings: **2293**; cross-domain facade→leaf pairs: **132**
| facade | domain | bindings | leaves |
|---|---|---:|---|
@ -13,6 +13,7 @@ AST-derived inventory of compatibility facades, regenerated by `python scripts/r
| `ouroboros/context.py` | D03 | 4 | `ouroboros/context_runtime_facts.py` (4) |
| `ouroboros/delegate_custody.py` | D07 | 10 | `ouroboros/delegate_custody_reconcile.py` (9)<br>`ouroboros/delegate_evidence.py` (1) |
| `ouroboros/extension_loader.py` | D14 | 99 | `ouroboros/contracts/plugin_api.py` (7 ✗D19)<br>`ouroboros/extension_child_catalog.py` (8)<br>`ouroboros/extension_companion.py` (3)<br>`ouroboros/extension_import_staging.py` (6)<br>`ouroboros/extension_isolated_deps.py` (4)<br>`ouroboros/extension_liveness.py` (8)<br>`ouroboros/extension_plugin_api.py` (6)<br>`ouroboros/extension_registry_state.py` (20)<br>`ouroboros/extension_surface_names.py` (12)<br>`ouroboros/extension_ui_validation.py` (5)<br>`ouroboros/gateway/host_service.py` (1 ✗D11)<br>`ouroboros/provider_models.py` (1 ✗D02)<br>`ouroboros/skill_loader.py` (13)<br>`ouroboros/skill_token.py` (1)<br>`ouroboros/tools/skill_exec.py` (1)<br>`ouroboros/utils.py` (3 ✗D18) |
| `ouroboros/gateway/_helpers.py` | D11 | 1 | `ouroboros/__init__.py` (1 ✗D18) |
| `ouroboros/gateway/contracts.py` | D11 | 3 | `ouroboros/gateway/decision_contracts.py` (2)<br>`ouroboros/gateway/history_contracts.py` (1) |
| `ouroboros/gateway/history.py` | D11 | 1 | `ouroboros/gateway/cost_breakdown.py` (1) |
| `ouroboros/gateway/tasks.py` | D11 | 7 | `ouroboros/gateway/cost_breakdown.py` (1)<br>`ouroboros/gateway/task_decision.py` (1)<br>`ouroboros/gateway/task_events.py` (4)<br>`ouroboros/gateway/task_hurry.py` (1) |

View file

@ -849,11 +849,12 @@ def hot_store_growth_notes(env: Any) -> list:
notes.append(
"WARNING: HOT STORE GROWTH — the events chain (logs/events.jsonl + "
f"archive/events_*.jsonl) totals {events_chain_size / 1_000_000:.1f} MB "
f"(threshold {EVENTS_ARCHIVE_SCAN_WARN_BYTES // 1_000_000} MB). Custody "
"replay scans this chain on ownership questions. Legacy segments retain "
"inline delegated request bodies; new start rows reference the observability "
"store, without shrinking existing history. Investigate chain "
"indexing/compaction; archives are durable history and are never deleted."
f"(threshold {EVENTS_ARCHIVE_SCAN_WARN_BYTES // 1_000_000} MB). Each process's "
"first custody read folds this whole chain into its row memo (later reads fold "
"only appended bytes); forensic and retirement scans still walk it. Legacy "
"segments retain inline delegated request bodies; new start rows reference the "
"observability store, without shrinking existing history. Investigate a durable "
"compact custody projection; archives are durable history and are never deleted."
)
from ouroboros.context_budget import RETAINED_EXECUTION_DRIVES_WARN_COUNT
from ouroboros.headless import HEADLESS_TASKS_DIR, TASK_DRIVES_DIR

View file

@ -963,19 +963,25 @@ def build_recent_sections(
+ json.dumps(coverage_projection, ensure_ascii=False, sort_keys=True, default=str)
)
for log_name, header, formatter in (
("progress.jsonl", "## Recent progress", lambda rows: memory.summarize_progress(rows, limit=50)),
("tools.jsonl", "## Recent tools", memory.summarize_tools),
("events.jsonl", "## Recent events", memory.summarize_events),
# Each task reads ITS OWN newest rows through a bounded window (#131): a
# global tail filtered afterwards handed every task whatever share of the
# shared suffix it happened to occupy. Quotas are what the formatters render
# (progress 50; tools: 10 rendered + 20 scanned for review markers; events:
# type counts over the rows it is given, today's 200). The header states
# the window (BIBLE P1); a reader pages the rest with read_file on the log.
from ouroboros.jsonl_tail import coverage_line
for log_name, header, formatter, want in (
("progress.jsonl", "## Recent progress", lambda rows: memory.summarize_progress(rows, limit=50), 50),
("tools.jsonl", "## Recent tools", memory.summarize_tools, 20),
("events.jsonl", "## Recent events", memory.summarize_events, 200),
):
entries = memory.read_jsonl_tail(log_name, 200)
if task_id:
entries = [e for e in entries if str(e.get("task_id", "")).strip() == task_id]
entries, coverage = memory.read_task_recent(log_name, task_id, want if task_id else 200)
summary = formatter(entries)
if summary:
sections.append(f"{header}\n\n{summary}")
sections.append(f"{header} ({coverage_line(coverage)})\n\n{summary}")
supervisor_summary = memory.summarize_supervisor(memory.read_jsonl_tail("supervisor.jsonl", 200))
supervisor_summary = memory.summarize_supervisor(memory.read_task_recent("supervisor.jsonl", "", 200)[0])
if supervisor_summary:
sections.append("## Supervisor\n\n" + supervisor_summary)

View file

@ -287,11 +287,13 @@ SKILL_REVIEW_ROOT_TASKS_WARN_BYTES = 20_000_000
# explicit full-history read becomes seconds-scale; this is observability, not
# a retention gate and never shortens the memory horizon.
CHAT_ARCHIVE_SCAN_WARN_BYTES = 100_000_000
# Custody replay (delegate_custody) walks the WHOLE events chain — live file
# plus archive/events_*.jsonl — on ownership questions. This inherits the
# The FIRST custody read of each process folds the WHOLE events chain — live
# file plus archive/events_*.jsonl — into the process-local row memo
# (delegate_custody_memo); later reads fold only appended bytes. Explicit
# forensic and retirement scans still walk the chain. This inherits the
# pre-rotation 100MB replay-degradation signal, now measured over the chain;
# archives stay durable history (never GC'd), so the remediation is chain
# indexing/compaction, never deletion.
# archives stay durable history (never GC'd), so the remediation is a durable
# compact custody projection, never deletion.
EVENTS_ARCHIVE_SCAN_WARN_BYTES = 100_000_000
# Warn before the observed 242-of-253 retained-drive corpus becomes routine;
# count only direct children because startup health is an interactive path.

View file

@ -536,16 +536,40 @@ def _apply(state: Dict[str, RunCustody], row: Dict[str, Any]) -> None:
custody.settled = True
def _fold_rows(rows: Any) -> Dict[str, RunCustody]:
state: Dict[str, RunCustody] = {}
for row in rows:
_apply(state, row)
return state
def custody_rows(drive_root: Any) -> Tuple[Dict[str, Any], ...]:
"""Every custody row of the chain, served from the process-local memo.
The same rows ``_iter_rows`` yields (inline request bodies replaced by a
locator), advanced by the bytes appended since the last read and refolded
on any fingerprint doubt (``delegate_custody_memo``). Read-only.
"""
from ouroboros.delegate_custody_memo import custody_rows as _memo_rows
return _memo_rows(drive_root)
def replay(drive_root: Any,
rows: Optional[List[Dict[str, Any]]] = None) -> Dict[str, RunCustody]:
"""Rebuild every known run's custody from the durable rows (one pass).
``rows`` replays a pre-read snapshot so several projections can share ONE
consistent traversal (the atomic payload busy claim, gate fix 5a)."""
state: Dict[str, RunCustody] = {}
for row in rows if rows is not None else _iter_rows(event_log_path(drive_root)):
_apply(state, row)
return state
consistent traversal (the atomic payload busy claim, gate fix 5a). Without
``rows`` the fold runs over the memo's rows and is cached per memo
generation; the returned objects are always this caller's own copies."""
if rows is not None:
return _fold_rows(rows)
from ouroboros.delegate_custody_memo import clone_custody_state, folded_state
return folded_state(drive_root, _fold_rows, clone_custody_state)
def lookup(drive_root: Any, task_id: str, run_id: str) -> Tuple[str, Optional[RunCustody]]:
"""Answer OWNED / FOREIGN / UNKNOWN for ``run_id`` as seen by ``task_id``."""
@ -676,7 +700,7 @@ def run_timing(drive_root: Any, run_id: str) -> Tuple[str, int]:
started_ts, max_seconds = "", 0
if not rid:
return started_ts, max_seconds
for row in _iter_rows(event_log_path(drive_root)):
for row in custody_rows(drive_root):
if str(row.get("run_id") or "") != rid or str(row.get("type") or "") != STARTED:
continue
started_ts = started_ts or str(row.get("ts") or "")
@ -745,7 +769,7 @@ def invocation_record(drive_root: Any, invocation_id: str, *,
return None
found: Optional[Dict[str, Any]] = None
state, run_id = "pending", ""
for row in rows if rows is not None else _iter_rows(event_log_path(drive_root)):
for row in rows if rows is not None else custody_rows(drive_root):
if str(row.get("invocation_id") or "") != target:
continue
kind = str(row.get("type") or "")

View file

@ -0,0 +1,354 @@
"""Process-local memo of the delegated-run custody rows (razzant/ouroboros#804).
Every custody question — ownership, timing, pending invocations, open runs,
evidence — is answered from the rows ``delegate_custody.emit`` appended to the
rotated ``logs/events.jsonl`` chain. That chain is durable history and grows
without bound (hundreds of MB on a long-lived install), so re-reading it on
every ``delegate_wait``, ``delegate_start``, task start and completion is the
"full-history scan filtered down to the answer" DEVELOPMENT §03 forbids on an
interactive path.
This module is the warm-cache half of the fix, the shape of
``ouroboros/_usage_rows_memo.py``: an in-process copy of the custody rows plus
a fingerprint of the chain prefix they were read from, advanced by folding only
the bytes appended since the previous read and REFOLDED FROM SCRATCH on any
doubt. The durable rows stay the one authority (ARCHITECTURE §10, invariants 10
and 26): nothing here is written to disk, a refold costs exactly what one
``_iter_rows`` replay costs today, and an unreadable chain bypasses the memo
for that call instead of caching an "empty because unreadable" answer. It is
deliberately NOT the durable compact projection §03 also names
(``containment_faults.jsonl`` is that shape): each process pays one cold fold.
Fingerprint rule (the ledger's ``_read_new_records_locked``): the memo's
segments must be the chain's first ``k`` segments by ``(st_dev, st_ino)`` in
order — a rotated live file keeps its inode under its archive name, so the
prefix it consumed stays consumed; every consumed segment has
``size >= consumed``; one whose size did not move must keep its ``st_mtime_ns``
(a same-size rewrite refolds); a segment that grew is folded from ``consumed``
onward. A torn LIVE tail (no trailing newline) waits for the next call; a torn
line inside an immutable archive can never complete, so its bytes are consumed
and counted (the ``refresh_recently_settled_terminals`` rule). Disclosed
residual: a rewrite of an already-consumed prefix that preserves size and
mtime is invisible — the log is append-only by construction (``append_jsonl``
under the writer lock, rotation by ``os.replace``), and tests reset the memo
through ``reset_custody_memo`` (autouse fixture in ``tests/conftest.py``).
Inline request bodies (legacy ``delegate_run_start_requested`` rows, hundreds of
KB each) are not retained: the row carries a ``request_locator`` instead and
``delegate_pending.request_body`` re-reads that one line on demand, so the
memo holds a few MB of compact rows, never the legacy bodies.
"""
from __future__ import annotations
import copy
import json
import logging
import os
import pathlib
import threading
from dataclasses import dataclass, field
from typing import Any, Callable, Dict, List, Optional, Tuple
from ouroboros.utils import JsonlChainUnreadable, jsonl_archive_segments
log = logging.getLogger("ouroboros.delegate_custody")
REQUEST_LOCATOR_KEY = "request_locator"
def _custody():
"""The custody facade, resolved at call time (import cycle + test pins)."""
from ouroboros import delegate_custody
return delegate_custody
class _Refold(Exception):
"""A consumed segment changed in a way the prefix rule cannot advance over."""
@dataclass
class _Segment:
st_dev: int
st_ino: int
consumed: int
st_mtime_ns: int
@dataclass
class _ChainMemo:
segments: List[_Segment] = field(default_factory=list)
rows: List[Dict[str, Any]] = field(default_factory=list)
rows_view: Tuple[Dict[str, Any], ...] = ()
generation: int = 0
torn_archive_lines: int = 0
# (generation, folded state) for ``folded_state``; cloned on every return.
state_cache: Optional[Tuple[int, Any]] = None
_MEMOS: Dict[str, _ChainMemo] = {}
_LOCKS: Dict[str, threading.Lock] = {}
_REGISTRY_LOCK = threading.Lock()
def _key(path: pathlib.Path) -> str:
return str(pathlib.Path(path).resolve(strict=False))
def _lock_for(key: str) -> threading.Lock:
with _REGISTRY_LOCK:
lock = _LOCKS.get(key)
if lock is None:
lock = _LOCKS[key] = threading.Lock()
return lock
def reset_custody_memo(drive_root: Any = None) -> None:
"""Forget one drive root's memo (or every memo): the next read refolds."""
with _REGISTRY_LOCK:
if drive_root is None:
_MEMOS.clear()
return
_MEMOS.pop(_key(_custody().event_log_path(drive_root)), None)
def _enumerate(path: pathlib.Path) -> List[Tuple[pathlib.Path, os.stat_result, bool]]:
"""The chain as ``(segment, stat, is_live)`` in fold order; strict on archives."""
chain: List[Tuple[pathlib.Path, os.stat_result, bool]] = []
for segment in jsonl_archive_segments(path, strict=True):
try:
chain.append((segment, segment.stat(), False))
except FileNotFoundError:
continue # rotated away between enumeration and stat: not part of history
except OSError as exc:
raise JsonlChainUnreadable(f"cannot stat {segment}: {exc}") from exc
try:
chain.append((path, path.stat(), True))
except FileNotFoundError:
pass # an absent live file is a positively empty tail
return chain
def _prefix_intact(memo: _ChainMemo, chain: List[Tuple[pathlib.Path, os.stat_result, bool]]) -> bool:
if len(memo.segments) > len(chain):
return False
for known, (_segment, stat, _is_live) in zip(memo.segments, chain):
if (stat.st_dev, stat.st_ino) != (known.st_dev, known.st_ino):
return False
if stat.st_size < known.consumed:
return False
if stat.st_size == known.consumed and stat.st_mtime_ns != known.st_mtime_ns:
return False
return True
def _compact_row(row: Dict[str, Any], locator: Tuple[int, int, int, int]) -> Dict[str, Any]:
"""Drop an inline request body from a start row; keep a locator to re-read it."""
request = row.get("request")
if isinstance(request, dict) and request:
row = {k: v for k, v in row.items() if k != "request"}
row[REQUEST_LOCATOR_KEY] = {
"st_dev": locator[0], "st_ino": locator[1],
"offset": locator[2], "length": locator[3],
}
return row
def _fold_segment(
memo: _ChainMemo, segment: pathlib.Path, stat: os.stat_result, is_live: bool,
known: Optional[_Segment], *, inner: bool,
) -> _Segment:
"""Fold one segment's bytes past ``known.consumed`` into the memo's rows.
``inner`` marks a consumed segment that is NOT the memo's last one: rows
from later segments were already folded after it, so a complete line
appended to it cannot be folded in chain order and forces a refold. The
last consumed segment may grow freely — it was the live file when its
tail was held back, and nothing after it has been folded yet.
"""
marker = _custody()._ROW_MARKER.encode("ascii")
start = known.consumed if known is not None else 0
consumed = start
with segment.open("rb") as handle:
if start:
handle.seek(start)
for raw in handle:
if not raw.endswith(b"\n"):
if is_live:
break # a torn live tail completes on a later call
# An archive never completes its torn tail: consume it, count it.
memo.torn_archive_lines += 1
consumed += len(raw)
continue
if inner:
raise _Refold("inner archive segment grew after it was consumed")
offset, consumed = consumed, consumed + len(raw)
if marker not in raw:
continue
try:
row = json.loads(raw.decode("utf-8", errors="replace"))
except ValueError:
continue
if isinstance(row, dict) and str(row.get("type") or "").startswith(_custody()._ROW_MARKER):
memo.rows.append(_compact_row(row, (stat.st_dev, stat.st_ino, offset, len(raw))))
try:
after = os.fstat(handle.fileno())
except OSError:
after = stat
return _Segment(st_dev=stat.st_dev, st_ino=stat.st_ino, consumed=consumed, st_mtime_ns=after.st_mtime_ns)
def _advance(memo: _ChainMemo, chain: List[Tuple[pathlib.Path, os.stat_result, bool]]) -> None:
before = len(memo.rows)
last_known = len(memo.segments) - 1
for index, (segment, stat, is_live) in enumerate(chain):
known = memo.segments[index] if index < len(memo.segments) else None
if known is not None and stat.st_size == known.consumed:
continue
folded = _fold_segment(memo, segment, stat, is_live, known,
inner=known is not None and index < last_known)
if known is None:
memo.segments.append(folded)
else:
memo.segments[index] = folded
if len(memo.rows) != before:
memo.generation += 1
memo.rows_view = tuple(memo.rows)
memo.state_cache = None
def _refresh(drive_root: Any) -> Tuple[Optional[_ChainMemo], Tuple[Dict[str, Any], ...]]:
"""Advance (or refold) the memo for ``drive_root``; ``(None, rows)`` when bypassed.
Caller holds the per-path lock. A bypass serves today's lenient
``_iter_rows`` answer without caching it.
"""
custody = _custody()
path = custody.event_log_path(drive_root)
key = _key(path)
memo = _MEMOS.get(key)
try:
chain = _enumerate(path)
except JsonlChainUnreadable:
_MEMOS.pop(key, None)
return None, tuple(custody._iter_rows(path))
try:
if memo is not None and not _prefix_intact(memo, chain):
memo = None # the consumed prefix is not the chain's prefix any more: refold
if memo is not None:
try:
_advance(memo, chain)
except _Refold:
memo = None
if memo is None:
memo = _ChainMemo()
_advance(memo, chain)
memo.generation = 1
memo.rows_view = tuple(memo.rows)
except (OSError, _Refold):
_MEMOS.pop(key, None)
log.debug("custody memo bypassed for %s", path, exc_info=True)
return None, tuple(custody._iter_rows(path))
_MEMOS[key] = memo
return memo, memo.rows_view
def custody_rows(drive_root: Any) -> Tuple[Dict[str, Any], ...]:
"""Every custody row of ``drive_root``'s chain, in chain order (read-only).
The same dicts ``_iter_rows`` yields, except that an inline request body is
replaced by ``request_locator``; callers never mutate them.
"""
key = _key(_custody().event_log_path(drive_root))
with _lock_for(key):
_memo, rows = _refresh(drive_root)
return rows
def folded_state(
drive_root: Any,
fold: Callable[[Tuple[Dict[str, Any], ...]], Any],
clone: Callable[[Any], Any] = copy.deepcopy,
) -> Any:
"""``fold(rows)`` over the current rows, cached per memo generation, cloned on return.
``clone`` copies the cached fold for the caller (the custody fold supplies
its own container-aware copy; ``deepcopy`` is the safe default). A bypassed
memo folds the lenient rows uncached, exactly like a replay today.
"""
key = _key(_custody().event_log_path(drive_root))
with _lock_for(key):
memo, rows = _refresh(drive_root)
if memo is None:
return fold(rows)
if memo.state_cache is None or memo.state_cache[0] != memo.generation:
memo.state_cache = (memo.generation, fold(rows))
return clone(memo.state_cache[1])
def clone_custody_state(state: Dict[str, Any]) -> Dict[str, Any]:
"""A caller-owned copy of a folded custody state.
``copy.copy`` per entry plus a copy of every mutable container the fold
writes (``resource_ref``, ``work_order_source_request``,
``verified_source_ranges`` and the delivery confirmations attribute), so a
caller mutating its answer — ``lookup`` stores it in ``_CUSTODY`` and the
``record_*`` writers update it in place — never changes the cached fold.
"""
clones: Dict[str, Any] = {}
for run_id, entry in state.items():
clone = copy.copy(entry)
clone.resource_ref = dict(entry.resource_ref)
clone.work_order_source_request = dict(entry.work_order_source_request)
clone.verified_source_ranges = list(entry.verified_source_ranges)
confirmations = getattr(entry, "_source_delivery_confirmations", None)
if confirmations is not None:
setattr(clone, "_source_delivery_confirmations", copy.deepcopy(confirmations))
clones[run_id] = clone
return clones
def read_locator_request(drive_root: Any, locator: Any) -> Optional[Dict[str, Any]]:
"""Re-read the inline ``request`` body a compacted row points at, or None."""
if not isinstance(locator, dict):
return None
try:
identity = (int(locator["st_dev"]), int(locator["st_ino"]))
offset, length = int(locator["offset"]), int(locator["length"])
except (KeyError, TypeError, ValueError):
return None
path = _custody().event_log_path(drive_root)
try:
for segment, stat, _is_live in _enumerate(path):
if (stat.st_dev, stat.st_ino) != identity:
continue
with segment.open("rb") as handle:
handle.seek(offset)
raw = handle.read(length)
row = json.loads(raw.decode("utf-8", errors="replace"))
body = row.get("request") if isinstance(row, dict) else None
return body if isinstance(body, dict) and body else None
except (OSError, ValueError, JsonlChainUnreadable):
return None
return None
def memo_diagnostics(drive_root: Any) -> Dict[str, Any]:
"""Facts about one memo for tests and forensics (never authority)."""
key = _key(_custody().event_log_path(drive_root))
with _lock_for(key):
memo = _MEMOS.get(key)
if memo is None:
return {"cold": True}
return {
"cold": False, "generation": memo.generation, "rows": len(memo.rows),
"segments": [(s.st_ino, s.consumed) for s in memo.segments],
"torn_archive_lines": memo.torn_archive_lines,
}
__all__ = [
"REQUEST_LOCATOR_KEY", "clone_custody_state", "custody_rows", "folded_state",
"memo_diagnostics", "read_locator_request", "reset_custody_memo",
]

View file

@ -119,7 +119,7 @@ def task_execution_evidence(drive_root: Any, task_id: str) -> Dict[str, Any]:
pass
except OSError:
evidence_read_failed = True
for row in custody._iter_rows(_log_path):
for row in custody.custody_rows(drive_root):
if str(row.get("task_id") or "") != tid:
continue
if str(row.get("type") or "") == NANNY_NUDGE_STAMP:
@ -355,7 +355,7 @@ def acceptance_patch_dispositions(drive_root: Any, task_id: str) -> Dict[str, An
except OSError:
return {"evidence_read_failed": True}
rows: List[Dict[str, Any]] = []
for row in custody._iter_rows(log_path):
for row in custody.custody_rows(drive_root):
if str(row.get("type") or "") != "delegate_run_patch_verdict":
continue
if str(row.get("task_id") or "") != tid:

View file

@ -15,7 +15,7 @@ def pending_invocations(
found: Dict[str, Dict[str, Any]] = {}
state: Dict[str, str] = {}
source = rows if rows is not None else c._iter_rows(c.event_log_path(drive_root))
source = rows if rows is not None else c.custody_rows(drive_root)
for row in source:
invocation_id = str(row.get("invocation_id") or "")
if not invocation_id:
@ -30,6 +30,7 @@ def pending_invocations(
"operation_id": str(row.get("operation_id") or ""),
"request": row.get("request") if isinstance(row.get("request"), dict) else None,
"request_ref": row.get("request_ref"),
"request_locator": row.get("request_locator"),
"route": str(row.get("route") or ""),
"project_id": str(row.get("project_id") or ""),
"project_owned": bool(row.get("project_owned")),
@ -74,8 +75,9 @@ def pending_invocations(
# Resolve only survivors, not every historical start on each sweep.
body = request_body(drive_root, record)
ref = record.pop("request_ref")
locator = record.pop("request_locator")
# An unreadable stored body does not discharge the pending start.
if body or ref is not None:
if body or ref is not None or locator is not None:
record["request"] = body
pending.append(record)
return pending
@ -90,6 +92,14 @@ def request_body(drive_root: Any, row: Dict[str, Any]) -> Optional[Dict[str, Any
inline = row.get("request")
if isinstance(inline, dict) and inline:
return inline
locator = row.get("request_locator")
if locator is not None:
# A memo row carries the legacy inline body's location, not the body.
from ouroboros.delegate_custody_memo import read_locator_request
located = read_locator_request(drive_root, locator)
if located is not None:
return located
ref = row.get("request_ref")
if not isinstance(ref, dict) or not ref:
return None

View file

@ -219,9 +219,7 @@ def unsettled_start_ids(
"""
mine = str(task_id or "")
snapshot = list(rows) if rows is not None else list(
custody._iter_rows(custody.event_log_path(drive_root))
)
snapshot = list(rows) if rows is not None else list(custody.custody_rows(drive_root))
runs = custody.replay(drive_root, rows=snapshot)
return {
"open_run_ids": [

View file

@ -18,7 +18,7 @@ _TRUE_LITERALS = frozenset({"1", "true", "yes", "on"})
_FALSE_LITERALS = frozenset({"0", "false", "no", "off"})
_TAIL_WINDOW_START_BYTES = 512 * 1024
from ouroboros.jsonl_tail import TAIL_WINDOW_START_BYTES as _TAIL_WINDOW_START_BYTES # noqa: E402,F401 (re-exported for gateway/history.py)
async def run_sync_to_completion(function, /, *args, **kwargs):
@ -56,29 +56,15 @@ def _read_jsonl_segment_with_gaps(
*,
tail_bytes: int | None = None,
) -> tuple[list, set[str]]:
"""Read one JSONL segment while retaining truthful parse/read-gap facts.
"""Gateway wrapper over ``jsonl_tail.read_jsonl_segment_with_gaps``.
``iter_jsonl_objects`` intentionally skips malformed rows because most
callers are best-effort telemetry readers. Gateway history is different:
it publishes a completeness claim, so a skipped row must remain visible as
a bounded read-gap fact even when the valid rows can still be rendered.
The parser is THIS module's ``iter_jsonl_objects`` name, resolved at call
time, so the gateway tests that monkeypatch it keep governing every gateway
read (``tests/test_gateway_history.py``).
"""
path = pathlib.Path(path)
try:
path.stat()
except FileNotFoundError:
return [], set()
except OSError:
return [], {"unreadable_source"}
from ouroboros.jsonl_tail import read_jsonl_segment_with_gaps
gaps: set[str] = set()
try:
entries = list(
iter_jsonl_objects(path, tail_bytes=tail_bytes, gap_reasons=gaps)
)
except OSError:
gaps.add("unreadable_source")
return entries, gaps
return read_jsonl_segment_with_gaps(path, tail_bytes=tail_bytes, iter_objects=iter_jsonl_objects)
def read_rotated_jsonl_entries(
@ -91,71 +77,18 @@ def read_rotated_jsonl_entries(
*,
include_gaps: bool = False,
) -> list | tuple[list, set[str]]:
"""Bounded, rotation-aware read of one JSONL log (v6.90.x P2, built on the
``iter_jsonl_objects(tail_bytes=...)`` bounded-read SSOT).
"""Bounded, rotation-aware read of one JSONL log (v6.90.x P2).
The live file is read from a byte tail that starts at 512KB and DOUBLES until
the FILTERED quota is satisfied (rows for which ``counts_toward_quota`` is
true reach ``want``) or the window covers the whole file — the degenerate
case equals today's full read, so a quota the file cannot satisfy costs one
full pass, never an infinite loop. Rotated ``archive/<prefix>_*.jsonl``
segments are then backfilled newest-first until the quota is met, bounded to
``max_archives`` files, and everything is reassembled chronologically
(oldest chosen archive -> live window). The backfill is bounded to the
``max_archives`` newest archives: older segments are NOT consulted by this
reader (they stay durable on disk for full-history consumers)."""
live = pathlib.Path(live)
try:
size = live.stat().st_size
except FileNotFoundError:
size = 0
except OSError:
size = 0
window = _TAIL_WINDOW_START_BYTES
gaps: set[str] = set()
while True:
if window >= size:
if include_gaps:
live_entries, live_gaps = _read_jsonl_segment_with_gaps(live)
gaps.update(live_gaps)
else:
live_entries = list(iter_jsonl_objects(live))
break
if include_gaps:
live_entries, live_gaps = _read_jsonl_segment_with_gaps(live, tail_bytes=window)
gaps.update(live_gaps)
else:
live_entries = list(iter_jsonl_objects(live, tail_bytes=window))
if sum(1 for entry in live_entries if counts_toward_quota(entry)) >= want:
break
window *= 2
collected = sum(1 for entry in live_entries if counts_toward_quota(entry))
try:
archives = sorted(
archive_dir.glob(f"{archive_prefix}_*.jsonl"), key=lambda p: p.name, reverse=True
)
except Exception:
archives = []
chosen: list = []
for archive_path in archives:
if collected >= want or len(chosen) >= max_archives:
break
try:
if include_gaps:
archive_entries, archive_gaps = _read_jsonl_segment_with_gaps(archive_path)
gaps.update(archive_gaps)
else:
archive_entries = list(iter_jsonl_objects(archive_path))
except Exception:
gaps.add("unreadable_source")
continue
chosen.append(archive_entries)
collected += sum(1 for entry in archive_entries if counts_toward_quota(entry))
ordered: list = []
for archive_entries in reversed(chosen): # oldest chosen archive first
ordered.extend(archive_entries)
ordered.extend(live_entries)
return (ordered, gaps) if include_gaps else ordered
The reader itself lives in ``ouroboros/jsonl_tail.py`` (one bounded
filtered tail for the endpoints AND context assembly); this wrapper keeps
the gateway call shape and its parser seam (see above).
"""
from ouroboros.jsonl_tail import read_rotated_jsonl_entries as _read
return _read(
live, archive_dir, archive_prefix, want, counts_toward_quota, max_archives,
include_gaps=include_gaps, iter_objects=iter_jsonl_objects,
)
def request_drive_root(request: Request) -> pathlib.Path:

202
ouroboros/jsonl_tail.py Normal file
View file

@ -0,0 +1,202 @@
"""Bounded, rotation-aware tail reads of one JSONL log — the ONE reader for them.
Moved here from ``gateway/_helpers.py`` (v6.90.x P2) so that context assembly
(``memory.py``, razzant/ouroboros#131) can use the same window-doubling
filtered tail the history, logs and routing endpoints already use, without a
core module importing the gateway layer (DEVELOPMENT §13). The gateway module
keeps thin wrappers over these functions with its own parser seam.
Semantics: the live file is read from a byte tail that starts at
``TAIL_WINDOW_START_BYTES`` and DOUBLES until the FILTERED quota is satisfied
(rows for which ``counts_toward_quota`` is true reach ``want``) or the window
covers the whole file — the degenerate case is one full read of a
rotation-bounded file, never a loop. Rotated ``archive/<prefix>_*.jsonl``
segments are then backfilled newest-first until the quota is met, bounded to
``max_archives`` files, and everything is reassembled chronologically (oldest
chosen archive -> live window). Older segments are NOT consulted: they stay
durable history for full-history readers. ALL rows of the chosen window are
returned (the quota decides where to stop, the caller filters); the optional
gap set names every parse or read failure the window met, and
``archives_bounded`` tells a caller whether the quota went unmet while older
archives were left unopened — the fact a continuity disclosure needs.
"""
from __future__ import annotations
import os
import pathlib
from typing import Any, Callable, Iterable, Optional
from ouroboros.utils import iter_jsonl_objects
TAIL_WINDOW_START_BYTES = 512 * 1024
ARCHIVE_BACKFILL_MAX = 3
def archive_segments(archive_dir: pathlib.Path, archive_prefix: str, gaps: Optional[set] = None) -> list:
"""``archive/<prefix>_*.jsonl`` newest-first by name (chronological by construction).
An explicit ``scandir`` because ``Path.glob`` swallows a ``PermissionError``
on the directory and yields nothing, which reads like "never rotated";
unreadable enumeration is reported in ``gaps`` instead.
"""
prefix = f"{archive_prefix}_"
try:
with os.scandir(archive_dir) as entries:
found = [
pathlib.Path(entry.path) for entry in entries
if entry.name.startswith(prefix) and entry.name.endswith(".jsonl") and entry.is_file()
]
except FileNotFoundError:
return []
except OSError:
if gaps is not None:
gaps.add("unreadable_source")
return []
return sorted(found, key=lambda p: p.name, reverse=True)
def read_jsonl_segment_with_gaps(
path: pathlib.Path,
*,
tail_bytes: Optional[int] = None,
iter_objects: Optional[Callable[..., Iterable[Any]]] = None,
) -> tuple[list, set[str]]:
"""Read one JSONL segment while retaining truthful parse/read-gap facts.
``iter_jsonl_objects`` intentionally skips malformed rows because most
callers are best-effort telemetry readers. History readers publish a
completeness claim, so a skipped row must remain visible as a bounded
read-gap fact even when the valid rows can still be rendered.
"""
path = pathlib.Path(path)
parse = iter_objects or iter_jsonl_objects # module name resolved at call time (test seam)
try:
path.stat()
except FileNotFoundError:
return [], set()
except OSError:
return [], {"unreadable_source"}
gaps: set[str] = set()
try:
entries = list(parse(path, tail_bytes=tail_bytes, gap_reasons=gaps))
except OSError:
gaps.add("unreadable_source")
return entries, gaps
def read_rotated_jsonl_entries(
live: pathlib.Path,
archive_dir: pathlib.Path,
archive_prefix: str,
want: int,
counts_toward_quota,
max_archives: int = ARCHIVE_BACKFILL_MAX,
*,
include_gaps: bool = False,
iter_objects: Optional[Callable[..., Iterable[Any]]] = None,
coverage: Optional[dict] = None,
) -> list | tuple[list, set[str]]:
"""Bounded, rotation-aware read of one JSONL log (module docstring).
``iter_objects`` is the parser seam (the gateway wrapper passes its own
name so its tests keep governing it). ``coverage``, when given, is filled
with ``live_size``, ``live_window`` (bytes of the live file read), ``archives``
(consulted count), ``archives_available`` and ``archives_bounded``.
"""
live = pathlib.Path(live)
parse = iter_objects or iter_jsonl_objects # module name resolved at call time (test seam)
try:
size = live.stat().st_size
except OSError:
size = 0
window = TAIL_WINDOW_START_BYTES
gaps: set[str] = set()
collect = include_gaps or coverage is not None # a coverage claim needs the gap facts too
while True:
if window >= size:
if collect:
live_entries, live_gaps = read_jsonl_segment_with_gaps(live, iter_objects=parse)
gaps.update(live_gaps)
else:
live_entries = list(parse(live))
window = size
break
if collect:
live_entries, live_gaps = read_jsonl_segment_with_gaps(
live, tail_bytes=window, iter_objects=parse)
gaps.update(live_gaps)
else:
live_entries = list(parse(live, tail_bytes=window))
if sum(1 for entry in live_entries if counts_toward_quota(entry)) >= want:
break
window *= 2
collected = sum(1 for entry in live_entries if counts_toward_quota(entry))
archives = archive_segments(pathlib.Path(archive_dir), archive_prefix, gaps if collect else None)
chosen: list = []
for archive_path in archives:
if collected >= want or len(chosen) >= max_archives:
break
try:
if collect:
archive_entries, archive_gaps = read_jsonl_segment_with_gaps(
archive_path, iter_objects=parse)
gaps.update(archive_gaps)
else:
archive_entries = list(parse(archive_path))
except Exception:
gaps.add("unreadable_source")
continue
chosen.append(archive_entries)
collected += sum(1 for entry in archive_entries if counts_toward_quota(entry))
ordered: list = []
for archive_entries in reversed(chosen): # oldest chosen archive first
ordered.extend(archive_entries)
ordered.extend(live_entries)
if coverage is not None:
coverage.update({
"live_size": size, "live_window": min(window, size),
"archives": len(chosen), "archives_available": len(archives),
"archives_bounded": collected < want and len(chosen) < len(archives),
"matched": collected, "gaps": sorted(gaps),
})
return (ordered, gaps) if include_gaps else ordered
def _kb(n: int) -> str:
return f"{n / 1024:.0f} KB" if n < 1024 * 1024 else f"{n / (1024 * 1024):.1f} MB"
def coverage_line(coverage: dict) -> str:
"""One compact, human-and-model-readable line for a section header.
Says whose rows they are, how many were shown of how many matched inside
the read window, what the window was, and whether older archives were left
unopened with the quota unmet — the facts a continuity disclosure needs.
"""
task_id = str(coverage.get("task_id") or "")
shown, matched = int(coverage.get("shown") or 0), int(coverage.get("matched") or 0)
whose = f"task {task_id}" if task_id else "all tasks"
if shown and matched > shown:
rows = f"newest {shown} of {matched} matching rows"
else:
rows = f"all {shown} matching rows" if shown else "no matching rows"
live_size, live_window = int(coverage.get("live_size") or 0), int(coverage.get("live_window") or 0)
window = "live file" if live_window >= live_size else f"live tail {_kb(live_window)} of {_kb(live_size)}"
archives, available = int(coverage.get("archives") or 0), int(coverage.get("archives_available") or 0)
if archives:
window += f" + {archives} of {available} newest archives"
parts = [f"{whose}: {rows}", f"window: {window}"]
if coverage.get("archives_bounded"):
parts.append("older archives not opened")
gaps = coverage.get("gaps") or []
if gaps:
parts.append("gaps: " + ", ".join(str(g) for g in gaps))
return "; ".join(parts)
__all__ = [
"ARCHIVE_BACKFILL_MAX", "TAIL_WINDOW_START_BYTES", "archive_segments", "coverage_line",
"read_jsonl_segment_with_gaps", "read_rotated_jsonl_entries",
]

View file

@ -878,6 +878,40 @@ class Memory:
def read_jsonl_tail(self, log_name: str, max_entries: int = 100) -> List[Dict[str, Any]]:
return self._read_jsonl_entries(log_name, max_entries=max_entries)
def read_task_recent(
self, log_name: str, task_id: str, want: int,
) -> Tuple[List[Dict[str, Any]], Dict[str, Any]]:
"""The newest ``want`` rows of ONE task (or of the log when ``task_id`` is
empty) through the bounded rotation-aware reader (razzant/ouroboros#131).
The window is a doubling byte tail of the live file plus at most the
three newest archives, so a busy neighbour cannot push this task's own
rows out of a shared global suffix, and the whole file is never parsed
for its tail. ``coverage`` states what the window was and whether the
quota went unmet while older archives stayed unopened (BIBLE P1: the
section discloses it; ``read_file`` on the log pages the rest).
"""
from ouroboros.jsonl_tail import read_rotated_jsonl_entries
wanted = str(task_id or "").strip()
def counts(entry: Dict[str, Any]) -> bool:
return not wanted or str(entry.get("task_id", "")).strip() == wanted
stem = log_name[:-len(".jsonl")] if log_name.endswith(".jsonl") else log_name
coverage: Dict[str, Any] = {"task_id": wanted}
try:
rows = read_rotated_jsonl_entries(
self.logs_path(log_name), self.drive_root / "archive", stem,
max(1, int(want)), counts, coverage=coverage,
)
except Exception:
log.warning("Failed to read recent %s rows", log_name, exc_info=True)
return [], {**coverage, "shown": 0, "matched": 0, "quota_met": False, "gaps": ["read_failed"]}
shown = [row for row in rows if counts(row)][-max(1, int(want)):]
coverage.update({"shown": len(shown), "quota_met": int(coverage.get("matched") or 0) >= int(want)})
return shown, coverage
def read_jsonl_tail_after_offset(
self,
log_name: str,

View file

@ -26,7 +26,7 @@ def _recoverable_review_invocations(drive_root: pathlib.Path) -> Dict[str, str]:
"""Unique durable delegated tokens keyed by their reserved operation."""
from ouroboros import delegate_custody
rows = list(delegate_custody._iter_rows(delegate_custody.event_log_path(drive_root)))
rows = list(delegate_custody.custody_rows(drive_root))
candidates: Dict[str, Set[str]] = {}
records = [
(record, str(record.get("invocation_id") or ""))

View file

@ -703,7 +703,7 @@ def _payload_delegation_busy(drive: pathlib.Path, target: pathlib.Path) -> str:
from ouroboros.delegate_terminal import _task_is_terminal
resolved = _resolved(target)
rows = list(custody._iter_rows(custody.event_log_path(drive)))
rows = list(custody.custody_rows(drive))
for run in custody.replay(drive, rows=rows).values():
if (run.authority_source == "skill_payload"
and _resolved(run.target_root) == resolved

View file

@ -480,6 +480,18 @@ def _rebind_runtime_roots_between_tests():
yield
@pytest.fixture(autouse=True)
def _reset_custody_memo_between_tests():
"""The custody row memo is process-local and keyed by events-log path; a test
that rewrites its log in place (``write_text``) or reuses a path must never
inherit another test's consumed prefix (``delegate_custody_memo``)."""
from ouroboros.delegate_custody_memo import reset_custody_memo
reset_custody_memo()
yield
reset_custody_memo()
@pytest.fixture(autouse=True)
def _unlatch_supervisor_event_bus_between_tests():
"""A TestClient lifespan runs the server shutdown, whose ``workers.shutdown_event_q()``

View file

@ -0,0 +1,266 @@
"""The custody row memo answers exactly what a full chain replay answers (#804).
Every test compares the memo-backed readers with the same reader fed a full
``_iter_rows`` pass over the rotated chain, after appends, rotations, torn
tails and the anomalies the fingerprint rule must refuse to advance over.
"""
from __future__ import annotations
import dataclasses
import json
import os
import pathlib
import pytest
from ouroboros import delegate_custody as custody
from ouroboros import delegate_custody_memo as memo
from ouroboros.utils import append_jsonl
from supervisor.state import rotate_jsonl_log_if_needed
def _rotate(root: pathlib.Path) -> None:
rotate_jsonl_log_if_needed(root, "events.jsonl", "events", max_bytes=1)
def _strip(row):
return {k: v for k, v in row.items() if k not in ("request", memo.REQUEST_LOCATOR_KEY)}
def _full_rows(root):
return [_strip(r) for r in custody._iter_rows(custody.event_log_path(root))]
def _as_dicts(state):
return {run_id: dataclasses.asdict(entry) for run_id, entry in state.items()}
def _assert_equivalent(root: pathlib.Path) -> None:
full = list(custody._iter_rows(custody.event_log_path(root)))
assert [_strip(r) for r in custody.custody_rows(root)] == [_strip(r) for r in full]
assert _as_dicts(custody.replay(root)) == _as_dicts(custody.replay(root, rows=full))
assert custody.pending_invocations(root) == custody.pending_invocations(root, rows=full)
for row in full:
run_id = str(row.get("run_id") or "")
if run_id:
assert custody.run_timing(root, run_id) == _timing_from_rows(full, run_id)
token = str(row.get("invocation_id") or "")
if token:
assert custody.invocation_record(root, token) == custody.invocation_record(root, token, rows=full)
def _timing_from_rows(rows, run_id):
started_ts, max_seconds = "", 0
for row in rows:
if row.get("run_id") != run_id or row.get("type") != custody.STARTED:
continue
started_ts = started_ts or str(row.get("ts") or "")
max_seconds = max_seconds or int(row.get("max_seconds") or 0)
return started_ts, max_seconds
def _started(root, run_id, task_id, shape=None, **extra):
assert custody.record_started(root, custody.RunCustody(
run_id=run_id, task_id=task_id, route_id="codex", model="m", ledger_root=str(root), **extra),
shape=shape)
def _requested(root, token, task_id, request):
assert custody.emit(root, custody.START_REQUESTED, {
"invocation_id": token, "task_id": task_id, "route": "codex", "request": request})
def test_memo_matches_full_replay_across_appends_and_rotations(tmp_path):
root = tmp_path
assert custody.custody_rows(root) == ()
_started(root, "run-1", "task-a", shape={"max_seconds": 90})
_assert_equivalent(root)
_rotate(root)
_requested(root, "inv-a", "task-a", {"prompt": "legacy inline body"})
assert custody.emit(root, custody.SETTLED, {"run_id": "run-1", "task_id": "task-a", "state": "succeeded"})
_assert_equivalent(root)
_rotate(root)
_rotate(root) # an empty rotation: the touched live file holds nothing
_started(root, "run-2", "task-b", invocation_id="inv-a")
append_jsonl(custody.event_log_path(root), {"ts": "t", "type": "llm_usage", "task_id": "task-b"})
_assert_equivalent(root)
assert custody.replay(root)["run-1"].settled and not custody.replay(root)["run-2"].settled
assert custody.run_timing(root, "run-1")[1] == 90
assert [r["invocation_id"] for r in custody.pending_invocations(root)] == []
diagnostics = memo.memo_diagnostics(root)
assert diagnostics["rows"] == 4 and len(diagnostics["segments"]) >= 3
def test_rotation_between_calls_advances_without_refold(tmp_path):
root = tmp_path
_started(root, "run-1", "task-a")
custody.custody_rows(root)
_started(root, "run-2", "task-a")
custody.custody_rows(root)
generation = memo.memo_diagnostics(root)["generation"]
assert generation == 2
_rotate(root)
_started(root, "run-3", "task-a")
_assert_equivalent(root)
# An advance bumps the generation; a refold would start a fresh memo at 1.
assert memo.memo_diagnostics(root)["generation"] == generation + 1
assert len(memo.memo_diagnostics(root)["segments"]) == 2
def test_torn_live_tail_waits_then_folds_exactly_once(tmp_path):
root = tmp_path
_started(root, "run-1", "task-a")
custody.custody_rows(root)
row = json.dumps({"ts": "t", "type": custody.SETTLED, "run_id": "run-1", "task_id": "task-a", "state": "failed"})
path = custody.event_log_path(root)
with path.open("ab") as handle:
handle.write(row[:20].encode("utf-8"))
assert not custody.replay(root)["run-1"].settled
with path.open("ab") as handle:
handle.write(row[20:].encode("utf-8") + b"\n")
assert custody.replay(root)["run-1"].settled
assert [r["type"] for r in custody.custody_rows(root)].count(custody.SETTLED) == 1
_assert_equivalent(root)
def test_torn_archive_tail_is_consumed_and_counted(tmp_path):
root = tmp_path
_started(root, "run-1", "task-a")
_rotate(root)
archive = sorted((root / "archive").glob("events_*.jsonl"))[-1]
with archive.open("ab") as handle:
handle.write(b'{"type": "delegate_run_settled", "run_id": "run-1", "task_id": "task-a"') # never terminated
_started(root, "run-2", "task-a")
_assert_equivalent(root)
assert memo.memo_diagnostics(root)["torn_archive_lines"] == 1
assert not custody.replay(root)["run-1"].settled
# The next call does not re-read or re-count the torn bytes.
custody.custody_rows(root)
assert memo.memo_diagnostics(root)["torn_archive_lines"] == 1
def test_same_size_rewrite_of_consumed_bytes_refolds(tmp_path):
root = tmp_path
_started(root, "run-1", "task-a")
assert custody.replay(root)["run-1"].task_id == "task-a"
path = custody.event_log_path(root)
original = path.read_bytes()
rewritten = original.replace(b'"task-a"', b'"task-b"')
assert len(rewritten) == len(original)
path.write_bytes(rewritten)
stat = path.stat()
os.utime(path, ns=(stat.st_atime_ns, stat.st_mtime_ns + 1_000_000))
assert custody.replay(root)["run-1"].task_id == "task-b"
_assert_equivalent(root)
@pytest.mark.parametrize("anomaly", ["archive_deleted", "archive_inserted_before", "live_truncated"])
def test_chain_anomalies_refold_to_the_full_replay(tmp_path, anomaly):
root = tmp_path
_started(root, "run-1", "task-a")
_rotate(root)
_started(root, "run-2", "task-a")
_rotate(root)
_started(root, "run-3", "task-a")
_assert_equivalent(root)
archives = sorted((root / "archive").glob("events_*.jsonl"))
if anomaly == "archive_deleted":
archives[0].unlink()
elif anomaly == "archive_inserted_before":
early = root / "archive" / "events_20000101T000000.jsonl"
early.write_text(json.dumps({"ts": "t", "type": custody.STARTED, "run_id": "run-0",
"task_id": "task-z"}) + "\n", encoding="utf-8")
else:
custody.event_log_path(root).write_text("", encoding="utf-8")
_assert_equivalent(root)
if anomaly == "archive_inserted_before":
assert "run-0" in custody.replay(root)
if anomaly == "live_truncated":
assert "run-3" not in custody.replay(root)
def test_unreadable_archive_directory_bypasses_the_memo(tmp_path):
if hasattr(os, "geteuid") and os.geteuid() == 0:
pytest.skip("root ignores directory permissions")
root = tmp_path
_started(root, "run-1", "task-a")
_rotate(root)
_started(root, "run-2", "task-a")
archive_dir = root / "archive"
archive_dir.chmod(0o000)
try:
rows = custody.custody_rows(root)
lenient = list(custody._iter_rows(custody.event_log_path(root)))
assert [r["run_id"] for r in rows] == [r["run_id"] for r in lenient] == ["run-2"]
assert memo.memo_diagnostics(root)["cold"]
finally:
archive_dir.chmod(0o755)
_assert_equivalent(root)
assert {r["run_id"] for r in custody.custody_rows(root)} == {"run-1", "run-2"}
def test_warm_reads_open_no_archive_segment(tmp_path, monkeypatch):
root = tmp_path
_started(root, "run-1", "task-a")
_rotate(root)
_started(root, "run-2", "task-a")
opened = []
original_open = pathlib.Path.open
def counting_open(self, *args, **kwargs):
if self.parent.name == "archive":
opened.append(self.name)
return original_open(self, *args, **kwargs)
monkeypatch.setattr(pathlib.Path, "open", counting_open)
custody.custody_rows(root)
assert len(opened) == 1 # the cold fold reads the archive once
custody.custody_rows(root)
custody.run_timing(root, "run-1")
custody.replay(root)
custody.pending_invocations(root)
assert len(opened) == 1 # every warm read is served from the memo
def test_replay_returns_independent_copies(tmp_path):
root = tmp_path
_started(root, "run-1", "task-a", resource_ref={"root": "skill_payload"})
first = custody.replay(root)["run-1"]
first.resource_ref["root"] = "tampered"
first.verified_source_ranges.append((0, 1))
first.settled = True
second = custody.replay(root)["run-1"]
assert second.resource_ref == {"root": "skill_payload"}
assert second.verified_source_ranges == [] and not second.settled
assert second is not first
def test_legacy_inline_request_survives_compaction(tmp_path, monkeypatch):
from ouroboros import observability
root = tmp_path
inline = {"prompt": "legacy inline body", "instructions": "x" * 2000}
_requested(root, "inv-legacy", "task-a", inline)
monkeypatch.setattr(observability, "read_blob_ref",
lambda *a, **k: pytest.fail("an inline body must not read a blob"))
rows = custody.custody_rows(root)
assert "request" not in rows[0] and memo.REQUEST_LOCATOR_KEY in rows[0]
assert custody.invocation_record(root, "inv-legacy")["request"] == inline
pending = custody.pending_invocations(root)
assert pending[0]["request"] == inline
assert "request_locator" not in pending[0] and "request_ref" not in pending[0]
_rotate(root) # the locator follows the row into its archive by inode
assert custody.invocation_record(root, "inv-legacy")["request"] == inline
def test_reset_forgets_the_memo(tmp_path):
root = tmp_path
_started(root, "run-1", "task-a")
custody.custody_rows(root)
assert not memo.memo_diagnostics(root)["cold"]
memo.reset_custody_memo(root)
assert memo.memo_diagnostics(root)["cold"]
custody.custody_rows(root)
memo.reset_custody_memo()
assert memo.memo_diagnostics(root)["cold"]

View file

@ -0,0 +1,129 @@
"""Recent-activity sections are each task's OWN newest rows (razzant/ouroboros#131).
Two-sided pins: the interleaved case that the global-tail-then-filter reader
lost, the quiet single-task case that must render exactly as before, the
review-marker window the tools quota must keep, and the coverage line the
header must carry (BIBLE P1: a bounded window is disclosed, never silent).
"""
from __future__ import annotations
import json
import os
import pathlib
import pytest
from ouroboros.context import build_recent_sections
from ouroboros.memory import Memory
def _write(path: pathlib.Path, rows) -> None:
path.parent.mkdir(parents=True, exist_ok=True)
path.write_text("\n".join(json.dumps(r) for r in rows) + "\n", encoding="utf-8")
def _interleaved(root: pathlib.Path) -> None:
"""Task A's rows sit BEFORE a burst of 250 rows from a busy neighbour."""
tools = [{"ts": f"2026-09-21T20:00:{i % 60:02d}", "task_id": "task-a", "tool": "read_file",
"args": {"path": f"a-{i}.py"}, "result_preview": "ok"} for i in range(30)]
tools += [{"ts": "2026-09-21T21:00:00", "task_id": "task-b", "tool": "run_command",
"args": {"cmd": f"b-{i}"}, "result_preview": "ok"} for i in range(250)]
progress = [{"ts": "t", "task_id": "task-a", "text": f"a-step-{i}"} for i in range(60)]
progress += [{"ts": "t", "task_id": "task-b", "text": f"b-step-{i}"} for i in range(250)]
events = [{"ts": "t", "task_id": "task-a", "type": "llm_round"} for _ in range(40)]
events += [{"ts": "t", "task_id": "task-b", "type": "tool_error", "error": "boom"} for _ in range(250)]
_write(root / "logs" / "tools.jsonl", tools)
_write(root / "logs" / "progress.jsonl", progress)
_write(root / "logs" / "events.jsonl", events)
def _section(sections, header):
return next(s for s in sections if s.startswith(header))
def test_task_sees_its_own_newest_rows_behind_a_busy_neighbour(tmp_path):
_interleaved(tmp_path)
sections = build_recent_sections(Memory(drive_root=tmp_path), env=None, task_id="task-a")
tools = _section(sections, "## Recent tools")
assert "a-29.py" in tools and "a-20.py" in tools and "b-" not in tools
progress = _section(sections, "## Recent progress")
assert "a-step-59" in progress and "a-step-10" in progress and "b-step" not in progress
events = _section(sections, "## Recent events")
assert "llm_round: 40" in events and "tool_error" not in events
# The guard's other side: the retired global-tail reader misses every A row.
stale = [e for e in Memory(drive_root=tmp_path).read_jsonl_tail("tools.jsonl", 200)
if e.get("task_id") == "task-a"]
assert stale == []
def test_single_task_log_renders_exactly_as_before(tmp_path):
rows = [{"ts": "t", "task_id": "task-a", "tool": "shell", "args": {"cmd": f"c{i}"},
"result_preview": "ok"} for i in range(5)]
_write(tmp_path / "logs" / "tools.jsonl", rows)
memory = Memory(drive_root=tmp_path)
tools = _section(build_recent_sections(memory, env=None, task_id="task-a"), "## Recent tools")
assert tools.split("\n\n", 1)[1] == memory.summarize_tools(rows)
assert "task task-a: all 5 matching rows; window: live file" in tools.splitlines()[0]
def test_no_task_id_keeps_the_global_tail(tmp_path):
rows = [{"ts": "t", "task_id": f"task-{i % 3}", "text": f"row-{i}"} for i in range(30)]
_write(tmp_path / "logs" / "progress.jsonl", rows)
progress = _section(build_recent_sections(Memory(drive_root=tmp_path), env=None), "## Recent progress")
assert "row-29" in progress and "row-0" in progress
assert "all tasks: all 30 matching rows" in progress.splitlines()[0]
def test_review_marker_inside_rows_eleven_to_twenty_survives(tmp_path):
rows = [{"ts": "t", "task_id": "task-a", "tool": "commit_reviewed", "args": {},
"result_preview": "REVIEW_BLOCKED: tests red"}]
rows += [{"ts": "t", "task_id": "task-a", "tool": "read_file", "args": {"path": f"f{i}"},
"result_preview": "ok"} for i in range(15)]
rows += [{"ts": "t", "task_id": "task-b", "tool": "x", "args": {}, "result_preview": "ok"}
for _ in range(300)]
_write(tmp_path / "logs" / "tools.jsonl", rows)
tools = _section(build_recent_sections(Memory(drive_root=tmp_path), env=None, task_id="task-a"), "## Recent tools")
assert "REVIEW_FAIL commit_reviewed" in tools
def test_coverage_line_discloses_bounded_archives_and_gaps(tmp_path):
logs = tmp_path / "logs"
archive = tmp_path / "archive"
for i in range(5):
_write(archive / f"tools_2026090{i}T000000.jsonl",
[{"ts": "t", "task_id": "task-a", "tool": "t", "args": {}, "result_preview": "ok"}])
_write(logs / "tools.jsonl", [{"ts": "t", "task_id": "task-a", "tool": "live", "args": {}, "result_preview": "ok"}])
rows, coverage = Memory(drive_root=tmp_path).read_task_recent("tools.jsonl", "task-a", 20)
assert len(rows) == 4 and coverage["archives"] == 3 and coverage["archives_available"] == 5
assert coverage["archives_bounded"] and not coverage["quota_met"]
header = _section(build_recent_sections(Memory(drive_root=tmp_path), env=None, task_id="task-a"),
"## Recent tools").splitlines()[0]
assert "3 of 5 newest archives" in header and "older archives not opened" in header
if hasattr(os, "geteuid") and os.geteuid() == 0:
pytest.skip("root ignores directory permissions")
archive.chmod(0o000)
try:
rows, coverage = Memory(drive_root=tmp_path).read_task_recent("tools.jsonl", "task-a", 20)
finally:
archive.chmod(0o755)
assert [r["tool"] for r in rows] == ["live"] and coverage["gaps"] == ["unreadable_source"]
def test_reader_never_parses_the_whole_live_file_when_the_tail_suffices(tmp_path, monkeypatch):
from ouroboros import jsonl_tail
rows = [{"ts": "t", "task_id": "task-a", "text": "x" * 2000} for _ in range(1500)] # ~3 MB
_write(tmp_path / "logs" / "progress.jsonl", rows)
windows = []
original = jsonl_tail.iter_jsonl_objects
def spy(path, *args, **kwargs):
windows.append(kwargs.get("tail_bytes"))
return original(path, *args, **kwargs)
monkeypatch.setattr(jsonl_tail, "iter_jsonl_objects", spy)
shown, coverage = Memory(drive_root=tmp_path).read_task_recent("progress.jsonl", "task-a", 50)
assert len(shown) == 50 and windows == [jsonl_tail.TAIL_WINDOW_START_BYTES]
assert coverage["live_window"] == jsonl_tail.TAIL_WINDOW_START_BYTES < coverage["live_size"]

View file

@ -35,7 +35,10 @@ CHAPTER_BYTE_BUDGETS: dict[str, int] = {
# plan review and task acceptance (the slot census vocabulary, the `awaiting`
# projection, the only-awaited task outcome); the in-flight sentence it grew from is
# replaced, the rest has no older text to displace.
"docs/architecture/06-agent-core.md": 287600,
# +1000 (2026-09-22): the custody row memo and the per-task recent-activity
# windows are two new mechanisms described in the paragraphs they changed;
# the base sat 28 bytes under the previous budget.
"docs/architecture/06-agent-core.md": 288600,
"docs/architecture/07-configuration.md": 36991,
# 18947 -> 19287: CI failure collection now documents diagnostic desktop builds while release remains gated.
"docs/architecture/08-git-branching-ci-and-build.md": 19287,
@ -51,7 +54,9 @@ CHAPTER_BYTE_BUDGETS: dict[str, int] = {
# 20650 -> 21100: one more rule the chapter lacked, the usage ledger's reader contract
# ("money never reads a snapshot; a display never waits on money"). It REPLACES the
# residual sentence of the off-thread invariant; the rule itself has no older text.
"docs/architecture/10-key-invariants.md": 21100,
# +200 (2026-09-22): invariant 10 names the process-local fingerprint memos
# and their fallback rule; the base sat 23 bytes under the previous budget.
"docs/architecture/10-key-invariants.md": 21300,
"docs/architecture/11-frozen-contracts-v1.md": 24194,
"docs/architecture/12-host-service-companions-and-chat-ids.md": 11007,
"docs/architecture/13-external-skills-layer.md": 7764,
@ -60,7 +65,10 @@ CHAPTER_BYTE_BUDGETS: dict[str, int] = {
# 22873 -> 23100: one new invariant (notifications ring for live events
# only). Its text was compressed to the load-bearing facts first; the
# remainder is the cost of stating a rule that did not exist before.
"docs/development/03-module-size-and-complexity.md": 23100,
# +300 (2026-09-22): two new house precedents (custody row memo, bounded
# filtered tail reader) join the projection-over-replay list; the base sat
# 15 bytes under the previous budget.
"docs/development/03-module-size-and-complexity.md": 23400,
"docs/development/04-core-governance-artifacts.md": 16431,
"docs/development/05-review-and-commit-protocol.md": 12956,
# 94197 -> 94520: the usage-ledger lock rule gains its reader contract (a display read

View file

@ -66,13 +66,15 @@ def test_active_owner_batch_reuses_one_history_snapshot(tmp_path, monkeypatch):
surface="multi_model_review", request={"prompt": "duplicate"},
)
scans = []
original = delegate_custody._iter_rows
original = delegate_custody.custody_rows
def counted(path, *args, **kwargs):
scans.append(path)
yield from original(path, *args, **kwargs)
def counted(drive_root, *args, **kwargs):
# One custody read per reconcile: the batch shares that snapshot
# across replay, pending and invocation projections.
scans.append(drive_root)
return original(drive_root, *args, **kwargs)
monkeypatch.setattr(delegate_custody, "_iter_rows", counted)
monkeypatch.setattr(delegate_custody, "custody_rows", counted)
result = review_owner_custody.reconcile_review_custody_after_confirmed_process_deaths(
tmp_path, {101, 202},
)