mirror of
https://github.com/razzant/ouroboros.git
synced 2026-10-03 12:18:39 +00:00
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:
parent
796f2e8708
commit
8aff64791d
24 changed files with 1125 additions and 137 deletions
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
||||
|
|
|
|||
|
|
@ -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) |
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
|
|
|||
|
|
@ -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 "")
|
||||
|
|
|
|||
354
ouroboros/delegate_custody_memo.py
Normal file
354
ouroboros/delegate_custody_memo.py
Normal 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",
|
||||
]
|
||||
|
|
@ -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:
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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": [
|
||||
|
|
|
|||
|
|
@ -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
202
ouroboros/jsonl_tail.py
Normal 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",
|
||||
]
|
||||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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 ""))
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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()``
|
||||
|
|
|
|||
266
tests/test_delegate_custody_memo.py
Normal file
266
tests/test_delegate_custody_memo.py
Normal 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"]
|
||||
129
tests/test_recent_sections_per_task.py
Normal file
129
tests/test_recent_sections_per_task.py
Normal 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"]
|
||||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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},
|
||||
)
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue