mirror of
https://github.com/razzant/ouroboros.git
synced 2026-10-02 19:58:46 +00:00
Merge pull request #1311 from razzant/ouroboros-agent/tz3-pr1-era-provenance
Some checks are pending
CI / quick-test (push) Waiting to run
CI / benchmark-methodology (push) Waiting to run
CI / full-test (push) Waiting to run
CI / betterleaks-platform-smoke (macos-latest) (push) Waiting to run
CI / betterleaks-platform-smoke (ubuntu-latest) (push) Waiting to run
CI / betterleaks-platform-smoke (windows-latest) (push) Waiting to run
CI / Android emulator smoke (API 30) (push) Blocked by required conditions
CI / Android emulator smoke (API 33) (push) Blocked by required conditions
CI / Android emulator smoke (API 36) (push) Blocked by required conditions
CI / android-build (push) Blocked by required conditions
CI / release-preflight (push) Blocked by required conditions
CI / build (dmg, macos-latest, macos-arm64, syft_1.50.0_darwin_arm64.tar.gz, syft, e32fdb9d47823fa633748a1efca2528fd77c37469ea93c9e40ab835da44e4cce) (push) Blocked by required conditions
CI / build (tar.gz, ubuntu-latest, linux-x86_64, syft_1.50.0_linux_amd64.tar.gz, syft, bf7b29ff57f06da30918266a0e1c2885a8f99784798d1bdb1628886aa015d788) (push) Blocked by required conditions
CI / build (zip, windows-latest, windows-x64, syft_1.50.0_windows_amd64.zip, syft.exe, 815ee6973ec5dff6a671d7f41b0e78835a8c45b91d5a39f4743ea1cee833d3be) (push) Blocked by required conditions
CI / vendor-package-smoke (push) Blocked by required conditions
CI / release (push) Blocked by required conditions
CI / integration-test (push) Waiting to run
CI / skill-smoke (macos-latest) (push) Waiting to run
CI / skill-smoke (ubuntu-latest) (push) Waiting to run
CI / skill-smoke (windows-latest) (push) Waiting to run
CI / marker-guards (push) Waiting to run
CI / ui-smoke (push) Waiting to run
CI / docker-ui-smoke (push) Waiting to run
CI / docker-portable-test (push) Waiting to run
CI / system-e2e-mock (push) Waiting to run
CI / e2e-live (SM1 x1 — largest subset feasible under the $30 cap) (push) Waiting to run
CI / android-test (push) Waiting to run
CI / Android emulator smoke (API 26) (push) Blocked by required conditions
CI / Android emulator smoke (API 29) (push) Blocked by required conditions
Claudexor platform gate (API keys — subscription auth NOT covered) / fixture · macos-latest · exact managed runtime, fake harness, no model (push) Waiting to run
Claudexor platform gate (API keys — subscription auth NOT covered) / fixture · ubuntu-latest · exact managed runtime, fake harness, no model (push) Waiting to run
Claudexor platform gate (API keys — subscription auth NOT covered) / fixture · windows-latest · exact managed runtime, fake harness, no model (push) Waiting to run
Claudexor platform gate (API keys — subscription auth NOT covered) / live · macos-latest · claude · API key only, subscription NOT covered (push) Waiting to run
Claudexor platform gate (API keys — subscription auth NOT covered) / live · ubuntu-latest · claude · API key only, subscription NOT covered (push) Waiting to run
Claudexor platform gate (API keys — subscription auth NOT covered) / live · windows-latest · claude · API key only, subscription NOT covered (push) Waiting to run
Claudexor platform gate (API keys — subscription auth NOT covered) / live · macos-latest · codex · API key only, subscription NOT covered (push) Waiting to run
UI browser (ouroboros push) / ui-smoke (push) Waiting to run
Some checks are pending
CI / quick-test (push) Waiting to run
CI / benchmark-methodology (push) Waiting to run
CI / full-test (push) Waiting to run
CI / betterleaks-platform-smoke (macos-latest) (push) Waiting to run
CI / betterleaks-platform-smoke (ubuntu-latest) (push) Waiting to run
CI / betterleaks-platform-smoke (windows-latest) (push) Waiting to run
CI / Android emulator smoke (API 30) (push) Blocked by required conditions
CI / Android emulator smoke (API 33) (push) Blocked by required conditions
CI / Android emulator smoke (API 36) (push) Blocked by required conditions
CI / android-build (push) Blocked by required conditions
CI / release-preflight (push) Blocked by required conditions
CI / build (dmg, macos-latest, macos-arm64, syft_1.50.0_darwin_arm64.tar.gz, syft, e32fdb9d47823fa633748a1efca2528fd77c37469ea93c9e40ab835da44e4cce) (push) Blocked by required conditions
CI / build (tar.gz, ubuntu-latest, linux-x86_64, syft_1.50.0_linux_amd64.tar.gz, syft, bf7b29ff57f06da30918266a0e1c2885a8f99784798d1bdb1628886aa015d788) (push) Blocked by required conditions
CI / build (zip, windows-latest, windows-x64, syft_1.50.0_windows_amd64.zip, syft.exe, 815ee6973ec5dff6a671d7f41b0e78835a8c45b91d5a39f4743ea1cee833d3be) (push) Blocked by required conditions
CI / vendor-package-smoke (push) Blocked by required conditions
CI / release (push) Blocked by required conditions
CI / integration-test (push) Waiting to run
CI / skill-smoke (macos-latest) (push) Waiting to run
CI / skill-smoke (ubuntu-latest) (push) Waiting to run
CI / skill-smoke (windows-latest) (push) Waiting to run
CI / marker-guards (push) Waiting to run
CI / ui-smoke (push) Waiting to run
CI / docker-ui-smoke (push) Waiting to run
CI / docker-portable-test (push) Waiting to run
CI / system-e2e-mock (push) Waiting to run
CI / e2e-live (SM1 x1 — largest subset feasible under the $30 cap) (push) Waiting to run
CI / android-test (push) Waiting to run
CI / Android emulator smoke (API 26) (push) Blocked by required conditions
CI / Android emulator smoke (API 29) (push) Blocked by required conditions
Claudexor platform gate (API keys — subscription auth NOT covered) / fixture · macos-latest · exact managed runtime, fake harness, no model (push) Waiting to run
Claudexor platform gate (API keys — subscription auth NOT covered) / fixture · ubuntu-latest · exact managed runtime, fake harness, no model (push) Waiting to run
Claudexor platform gate (API keys — subscription auth NOT covered) / fixture · windows-latest · exact managed runtime, fake harness, no model (push) Waiting to run
Claudexor platform gate (API keys — subscription auth NOT covered) / live · macos-latest · claude · API key only, subscription NOT covered (push) Waiting to run
Claudexor platform gate (API keys — subscription auth NOT covered) / live · ubuntu-latest · claude · API key only, subscription NOT covered (push) Waiting to run
Claudexor platform gate (API keys — subscription auth NOT covered) / live · windows-latest · claude · API key only, subscription NOT covered (push) Waiting to run
Claudexor platform gate (API keys — subscription auth NOT covered) / live · macos-latest · codex · API key only, subscription NOT covered (push) Waiting to run
UI browser (ouroboros push) / ui-smoke (push) Waiting to run
memory: era-run boundaries, typed maintenance events, answering-route provenance (TZ-3 PR-1, writer-neutral)
This commit is contained in:
commit
f0af68e824
14 changed files with 1258 additions and 163 deletions
|
|
@ -144,10 +144,10 @@ scanned data-relative path to be covered by a row here (count-anchored both ways
|
|||
| `memory/scratchpad.md` + `scratchpad_blocks.json` | `ouroboros/memory.py` (derived, regenerated from blocks under lock) | none | bounded: 10 blocks, eviction journaled first (fail-closed) | regenerated; evicted history in journal |
|
||||
| `memory/WORLD.md` | `ouroboros/world_profiler.py` (write-once) | none | fixed | regenerates on restart — deletion IS the refresh mechanism |
|
||||
| `memory/registry.md`, `memory/deep_review.md` | `ouroboros/tools/memory_tools.py` (section RMW), `ouroboros/agent.py` (overwrite) | none | unbounded / last-wins — accepted | recreated lazily |
|
||||
| `memory/dialogue_blocks.json` + `dialogue_meta.json` | `ouroboros/consolidator.py`, `memory_nomination_receipts.py` (locked atomic) | `pending_knowledge_nominations` source-entry IDs; legacy `last_unpublished_nominations` preserved | blocks bounded by era compression (10 blocks, oldest 4); unresolved nomination index unbounded; no tool-level resolver yet, later success never retires old debt | blocks: compressed biography irreproducible; meta: cursor and unpublished-obligation evidence lost |
|
||||
| `memory/dialogue_blocks.json` + `dialogue_meta.json` | `ouroboros/consolidator.py`, `memory_nomination_receipts.py` (locked atomic) | `pending_knowledge_nominations` source-entry IDs; legacy `last_unpublished_nominations` preserved; `era_retry` (source hash + effective Light dispatch binding as key + observed route of a not-shorter era) | blocks reduced by era compression (ordinary pass: over 10 blocks, the oldest run of up to 4 SUMMARY blocks before the newest; pressure pass: every complete run, uncapped; an era is never re-compressed — calendar-era recompression is deferred, not implemented — so eras accumulate unbounded); unresolved nomination index unbounded; no tool-level resolver yet, later success never retires old debt | blocks: compressed biography irreproducible; meta: cursor and unpublished-obligation evidence lost |
|
||||
| `memory/dialogue_summary.md` | none — legacy read-only (reader in context.py) | none | frozen | legacy artifact; nothing writes it |
|
||||
| `memory/knowledge/**` (topic .md + `index-full.md` + `patterns.md`) | `ouroboros/tools/knowledge.py`, `consolidator.py` (index rebuild), `reflection.py` (patterns CAS rewrite) | none | topic files unbounded — accepted (curated by consolidation); backlog topic merge-only fail-closed | recreated lazily; knowledge lost |
|
||||
| `memory/*_journal.jsonl`, `memory/knowledge_history.jsonl`, `memory/knowledge/patterns_history.jsonl` | `ouroboros/memory.py`, `tools/control_runtime.py`, `tools/knowledge.py`, `reflection.py` — every append through the `append_jsonl` sidecar-lock seam | scratchpad journal: `type` rows; others unversioned full-text snapshots; historical digested rows retain `content_digested: true` | complete new old+new snapshots are retained indefinitely; `memory_journal_compaction.py` is a read-only compatibility entry point, not a source rewriter; existing digest-only rows cannot be restored; the `memory_journal_observation` startup event gives byte sizes (or missing/unreadable) for the three named journals; scratchpad keeps its eviction journal | deleting the journals loses undo/provenance; eviction/rewrite paths fail closed when journal append fails; historically digested content remains irrecoverable |
|
||||
| `memory/*_journal.jsonl`, `memory/knowledge_history.jsonl`, `memory/knowledge/patterns_history.jsonl` | `ouroboros/memory.py`, `tools/control_runtime.py`, `tools/knowledge.py`, `reflection.py` — every append through the `append_jsonl` sidecar-lock seam | scratchpad journal: `type` rows; others unversioned full-text snapshots; `knowledge_history.jsonl` `source_capture` rows carry the host stamp `writer`/`route`/`writer_input_ref`/`old_chars`/`new_chars` (rows older than the stamp read `unknown`); historical digested rows retain `content_digested: true` | complete new old+new snapshots are retained indefinitely; `memory_journal_compaction.py` is a read-only compatibility entry point, not a source rewriter; existing digest-only rows cannot be restored; the `memory_journal_observation` startup event gives byte sizes (or missing/unreadable) for the three named journals; scratchpad keeps its eviction journal | deleting the journals loses undo/provenance; eviction/rewrite paths fail closed when journal append fails; historically digested content remains irrecoverable |
|
||||
| `memory/owner_mailbox/<task>.jsonl` + `.acks.jsonl` | `ouroboros/owner_mailbox.py` (append-only; revocation appends, reader resolves) | `kind` discriminator | lifecycle-bounded: unlinked at task terminal; a startup sweep unlinks mailboxes whose task has a SETTLED durable result (no result / non-terminal keeps the mailbox fail-closed) | undelivered owner directives + restart-surviving hurry latch lost; acks lost ⇒ re-delivery |
|
||||
|
||||
## 7. Skills payloads, tasks, uploads, projects, services
|
||||
|
|
@ -168,7 +168,7 @@ scanned data-relative path to be covered by a row here (count-anchored both ways
|
|||
| `uploads/**` (+`screenshots/`, `views/`, atomic-copy `.tmp` files) | `ouroboros/gateway/files.py` through `artifacts.copy_artifact_file` for chat uploads; `tools/browser.py`, `tools/vision.py`, `server_owner_routing.py` | raw owner bytes; chat ingestion measures size and SHA256 while copying | owner attachments in the `uploads/` root: NO retention, owner-explicit delete only; chat upload no longer applies the former 50 MiB cap, while Files-browser and downstream transport limits retain their own contracts; agent-generated `screenshots/`/`views/` age out past GC retention at startup | chat attachments dangle (readers skip missing); staged task copies survive |
|
||||
| `services/<task>/*.log` | `ouroboros/workspace_executor.py`, `tools/services.py` | none — raw text; archived content becomes observability blob with events.jsonl receipt | GC-retention prune at startup (archive-then-unlink); per-task terminal archive; oversize logs retained live | live tails lost; archived blobs survive |
|
||||
| `projects/<id>/**` (knowledge, journal, workpad, reflections, `.project.json` marker) | `ouroboros/project_facts.py`, `projects_registry.py` | marker unversioned; registry stamped (§2) | NEVER age-pruned (owner curates; delete = durable tombstone) | per-project memory lost; registry row survives, room reappears empty |
|
||||
| `projects/<id>/knowledge_history.jsonl`, `projects/<id>/knowledge_journal.jsonl` | `ouroboros/knowledge.py` through the shared knowledge write lock | append rows retain source topic, revision and operation facts | follows the owning project shelf; retained with the project until explicit deletion | history/provenance lost while authored project notes remain |
|
||||
| `projects/<id>/knowledge_history.jsonl` | `ouroboros/knowledge.py` through the shared knowledge write lock | append rows retain source topic, revision, operation facts and the host stamp; the former `knowledge_journal.jsonl` size telemetry beside it (global and project) is no longer written — an existing file is inert | follows the owning project shelf; retained with the project until explicit deletion | history/provenance lost while authored project notes remain |
|
||||
| `archive/**` (rotated segments, `rescue/`, `usage_import/`, `managed_repo/`) | rotation + `supervisor/git_ops_rescue.py`, `usage_legacy_import.py`, `launcher_bootstrap.py` | segments inherit source shape; usage_import carries sha256 sidecar | UNBOUNDED BY DESIGN — durable history, never GC'd (P1) — accepted | memory horizon truncated; rescue copies of uncommitted work destroyed |
|
||||
| `archive/usage_ledger/segment_*.jsonl` | `ouroboros/usage_compaction.py` (exact pre-compaction ledger bytes, written + fsync'd BEFORE the live swap) | each segment is a whole valid ledger generation; hash-pinned by the live `usage_baseline` header (`source_sha256`), chained recursively through each segment's own leading header | UNBOUNDED BY DESIGN — the folded monetary history, never GC'd (P1); read via `archived_attempt_ids` (tamper-evident, per-attempt joins for the model-send reverse sweep) | folded per-attempt monetary history unrecoverable; live aggregates (baseline block) survive, but seal/attempt joins for folded ids break — the mirror of the ledger row: deleting the ARCHIVE alone under a stamped ledger is a typed chain break on every history question, deleting the LEDGER alone raises `generation newer` for surviving newer, non-prefix segments; once fresh compactions reach those generations, old unreferenced segments are skipped and their attempt IDs are absent ([history readers](USAGE_COMPACTION.md#10-history-readers-model-send-reconciliation-and-audits)); reset both together |
|
||||
| `observability/{calls,blobs,salvaged}/**` | `ouroboros/observability.py` (private 0700/0600, CAS gzip); model-send records beside the call manifests: `model_send_seal` block + write-once `<attempt>.model_send_violation.json` typed facts (`ouroboros/model_send_seal.py`) | call manifests `schema_version: 1` + custody/redaction honesty markers; blob refs sha-verified on read; `model_send_seal.seal_version: 1` with the `canonical_json_v1` basis string | preserved indefinitely BY CONTRACT (the startup census counts, never deletes); the inert `OUROBOROS_OBSERVABILITY_RETENTION_DAYS` knob is RETIRED (key in `RETIRED_SETTING_KEYS`) | every recorded `result_ref`/`manifest_ref` dangles (strict readers raise); pending delegated request bodies become unknown, while custody identity still blocks replacement and protects snapshots; replay refuses without the recorded body; salvaged outputs unrecoverable; a lost seal on a seam-dispatched attempt surfaces as a typed `unlogged_attempt` fact at the next startup sweep |
|
||||
|
|
|
|||
|
|
@ -527,7 +527,7 @@ Availability follows `deep_review_route`, never a window floor. An API row needs
|
|||
|
||||
Typed root post-task triggers decide whether a run warrants Experience Review. `reflection.generate_reflection` sends the Light route one open prompt with the EXACT initial text (never a prefix) and its host-recorded `task_inputs.run_origin` beside it (provenance, never by itself the accepted requirement) plus tool-use, error, review and child projections and the same frozen non-final cost snapshot the task summary uses; it runs outside the tool loop, records its own usage, and its failure never erases the delivered result or changes a review verdict. Its execution trace is the ALL-CALLS listing (`build_trace_summary(all_calls=True)`): every call in order with every argument, identical consecutive calls folded into one `×N` row with their rounds, the first line of a failed or repeated call's result, and one header count of rounds whose every call was non-ok — no positional window and no literal cut, because the consolidation seam fits the call to the Light route whenever that route's window is known (an unknown window sends the prompt unchecked — the accepted residual of enlarging it); the STORED `trace_summary` (task card, parents, children) stays the bounded two-argument preview. The trace row carries the `round_id` of the model round that issued the call (absent when unknown), when the listing really cut an argument value or a failed/repeated call's answer, the redacted per-call record is retained through `retain_memory_source` and named in the prompt as OPTIONAL reading (never a required source), and its claim states the STORED bounds on both axes, not completeness: an argument already passed `sanitize_tool_args_for_log` (an oversized value carries a marker with its length and sha) and a result is the stored actor-visible cap — more than the listing, which shows only the first line of a failed or repeated answer — a partial one naming its own `FULL_RESULT_SOURCE_JSON` or `FULL_RESULT_SOURCE_UNAVAILABLE`, with a call's recorded manifest named only when it has one — claiming results "in full" or an unconditional manifest overstated a cognitive artifact. Cut detection reads the ONE shared marker list (`artifacts.SANITIZER_OMISSION_MARKERS`), because a width test over already-sanitized args measured the widest argument in the task as a small one and retained nothing at all, while a hand-rolled subset missed the `_repr` and `_error` shapes whose arguments survive only in the call blob. Unavailable source retention is disclosed, error details group by full redacted content before display clipping, and post-task synthesis — the reflection, its Pattern Register update and the episodic summary — thinks at the owner's Task / Chat effort (`settings_scales.resolve_effort("task")`), never a literal. Admission to the Pattern Register is typed, not a word scan: it opens on a call the loop recorded as errored (its stamped `tool_result_code`, or the recorded status for a legacy row), on a producer fact naming a failure the ok status cannot carry (a preserved commit whose post-commit tests failed publishes `post_commit_tests`), on typed codes already stored with an entry, or on a genuinely FAILED child — a cancelled, soft-landed best-effort or degraded child is not a failure. Those failed-child classes reach the root through the child evidence the synthesis walk already collects and make the run error-bearing for both the trigger and the prompt's error details: children do not reflect, so a short clean root that delegated the work is the only place its child's failure can be learned from at all. Deliberately not "a reason code exists", which would open the register on every terminal.
|
||||
|
||||
A reflection lands where it durably belongs: a non-project root appends the full entry to the canonical `logs/task_reflections.jsonl`; a project-scoped root appends the full entry to its project drive and the canonical log receives only a bounded pointer row — full project text never enters the canonical log, which feeds future global context. A project-bound task's context includes a bounded labeled tail of its own project's reflections; the headless mirror drive of a split root is never the reflection home, and the Pattern Register update stays canonical in both cases. Every entry carries task identity, evidence, lessons, backlog candidates, and validated memory actions. `MEMORY_ACTIONS_JSON` permits only `scratchpad_append`, `knowledge_write`, and `identity_update_candidate`, at bounded count and size. `apply_memory_actions` routes accepted actions through provenance-preserving memory and knowledge APIs. An `identity_update_candidate` is recorded in the scratchpad for review and is never auto-written to `identity.md`. For a project-scoped task, reflection applies knowledge actions only (the project store by default, explicit global allowed); its scratchpad and identity-candidate actions are skipped because this automatic Light pass lacks the conversation's full view — the conversation writes identity and scratchpad from any room through its own tools. Reflection may propose a future campaign or backlog item, but it cannot enqueue, review, commit, or enable one.
|
||||
A reflection lands where it durably belongs: a non-project root appends the full entry to the canonical `logs/task_reflections.jsonl`, a project-scoped root to its project drive with only a bounded pointer row in the canonical log. A project-bound task's context holds a bounded labeled tail of its project's reflections; a split root's headless mirror drive is never the reflection home; the Pattern Register update stays canonical. Entries carry task identity, evidence, lessons, backlog candidates and validated memory actions. `MEMORY_ACTIONS_JSON` permits only `scratchpad_append`, `knowledge_write` and `identity_update_candidate`, bounded in count and size; `apply_memory_actions` applies them via provenance-preserving memory and knowledge APIs; an identity candidate lands in the scratchpad for review, never auto-written to `identity.md`. A project-scoped task applies knowledge actions only (project store by default, explicit global allowed), skipping scratchpad and identity candidates, which the conversation writes from any room with the full view this Light pass lacks; each skipped, empty or topic-less action is a `reflection_memory_action_skipped` event with its reason and `input_ref`: the retained exact task-input prompt, or an action's reflection-entry copy, else the canonical log pointer or `source_unavailable`. Reflection may propose a campaign or backlog item, never enqueue, review, commit or enable one.
|
||||
|
||||
The Pattern Register writer (`reflection._update_patterns`) REPLACES the whole document, so every decision input it reads is complete — the full current register, the exact initial text (`goal_exact`, beside the bounded `goal` display) and its `run_origin` and the whole reflection text; a prefix of a decision input can never authorize the rewrite, because a clipped clause can record the inverse of what the reflection concluded. `patterns` is a reserved global-only topic alongside `improvement-backlog` and `overview`: whichever room writes it, the register has ONE home on the canonical drive, the only one its writer and its readers (context assembly, deep self-review, the headless copy) address. Re-read, exact compare, history append and atomic replacement are one critical section outside the Light call; a register that moved under a losing writer is preserved and that task's learning is DROPPED with a warning naming the task, never retried against a source it did not decide from. A row's count is bumped once per observed episode, so two roots of one owner request bump it twice — the number counts episodes, not distinct requests.
|
||||
|
||||
|
|
@ -555,17 +555,17 @@ 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 recent-activity sections are each task's OWN newest rows (progress 50 rendered; tools 20 selected, 10 rendered and 20 scanned for review markers; events 200 counted by type) 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; a subagent child gets the same three windows beside its `## Working sources` block, its tools and events read from its own execution drive (its worker rows; host-side rows such as waits stay in the canonical log, as the header says) and progress from the canonical log; 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 rendered; tools 20 selected, 10 rendered and 20 scanned for review markers; events 200 counted by type) via the bounded reader `jsonl_tail.py` (`Memory.read_task_recent`: a doubling live tail plus at most three newest archives), never a filtered global tail (issue #131), and their header's coverage line names the rows, the window and any unopened archives, `read_file` paging the rest; a subagent child gets the same three windows beside its `## Working sources` block, its tools and events read from its own execution drive (its worker rows; host-side rows like waits stay in the canonical log, as the header says) and progress from the canonical log; 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 context build retires a block; scratchpad replacement keeps its explicit summary and source-journal provenance.
|
||||
|
||||
`consolidator.py` publishes a block and advances its generation-aware cursor only after complete room draft and correction; a missing generation appends `[MEMORY GAP]` instead of resetting. `context_fit` measures Light against fresh route/account capacity and calibrated density; `llm_local` owns the local output reserve, and absent evidence stays unknown. `room_consolidation.py` processes each room separately and assembles sections deterministically. Both knowledge stages receive the entire current note and source episode; the corrector's complete read, not the draft's, binds revised entries. Older episodes cannot negate newer facts; model judgment governs supported corrections. Source range reads and CAS preserve old/new history. Startup compaction is a no-op; earlier digests cannot be reversed. A failed correction withholds the chunk and cursor.
|
||||
`consolidator.py` publishes a block and advances its generation-aware cursor only after complete room draft and correction; a missing generation appends `[MEMORY GAP]`, never resets. `context_fit` measures Light against fresh route/account capacity and calibrated density; `llm_local` owns the local output reserve; absent evidence stays unknown. `room_consolidation.py` handles rooms separately, assembling sections deterministically. Both knowledge stages get the whole current note and source episode; the corrector's complete read, not the draft's, binds revised entries. Older episodes never negate newer facts; model judgment governs supported corrections. Source range reads and CAS preserve old/new history. Startup compaction is a no-op; earlier digests are irreversible. A failed correction withholds the chunk and cursor.
|
||||
|
||||
Oversized sources split without clipping, including within an entry. `consolidation_retry` records source hash and a smaller same-route bound, invalidated by source/route/capacity/reserve changes. Era compression regroups rooms deterministically; failed eras retain blocks and legacy provenance stays unknown. `last_consolidation_error` clears after a failure-free advance. `pending_knowledge_nominations` records each source entry BEFORE note publication; an unrelated successful batch cannot remove one. Legacy `last_unpublished_nominations` persists. Health shows three distinct abbreviated source+position IDs and omitted count; full proposals live in `knowledge_history.jsonl`. Unreadable meta preserves debts and warns in Health; Nano pressure records no-progress instead of aborting Main. Invalid legacy receipts warn separately. Old digests and debts need explicit resolution. Spend remains nullable; model-control errors follow `propagate_model_error`.
|
||||
Oversized sources split, never clip, even mid-entry. `consolidation_retry`: source hash + a smaller same-route bound, void on source/route/capacity/reserve change. An era is ONE summary run bounded by gaps/eras (never an era of an era; calendar-era recompression deferred, not implemented): ordinary pass: the OLDEST run of up to four before the newest block; pressure pass: EVERY complete run, newest too, uncapped. A failed or not-shorter era keeps its blocks, the latter as `era_retry` (PER RUN by source hash; key = the EXECUTED Light dispatch binding `_light_dispatch_binding` (read post-call): lane, account pin, model-wait override; answering route a separate stamp; newest `ERA_RETRY_MAX_RUNS` kept, a legacy record reads), honoured by both passes (no paid repeat while run and key hold), an `era_not_shorter` event in Health; a lock skip is `consolidation_skipped_locked`, each scratchpad pass ends in `scratchpad_consolidation`, legacy provenance stays unknown. `last_consolidation_error` clears after a failure-free advance. `pending_knowledge_nominations` holds each source entry BEFORE note publication; no unrelated success removes one; legacy `last_unpublished_nominations` persists. Health shows three short source+position IDs and the omitted count; full proposals: `knowledge_history.jsonl`. Unreadable meta keeps debts, warning in Health; Nano pressure records no-progress, never aborts Main; invalid legacy receipts warn separately; old digests and debts need explicit resolution. Spend stays nullable; model-control errors follow `propagate_model_error`.
|
||||
|
||||
Consolidation labels every chronological source message through `dialogue_provenance.RoomLabelResolver`, using the actual `chat_id`, never lineage `project_id`; one read-only registry snapshot supplies the window. Main is named only for the actual Main id, a resolved project uses its current registry name and stable chat id, and missing, unknown or ambiguous rooms stay explicit. Ephemeral formatter offsets carry the original room/author/direction/transport header into split continuations without parsing message bodies or duplicating their bytes. Room draft and correction prompts require meaningful decisions, approvals, outcomes and unresolved commitments of that room, retaining source distinctions (who decided, what was authorized, what stays owed) and one first-person Ouroboros voice. Length adapts to content within the existing output-token ceiling; no per-room word quota, semantic gate or absent room is imposed. Labels establish provenance, not summary success. The mixed Main recent view opts into the same labels; focused Project rendering, membership and explicit `chat_history` retain their existing behavior and bytes.
|
||||
Consolidation labels each chronological source message via `dialogue_provenance.RoomLabelResolver`, using the actual `chat_id`, never lineage `project_id`; one read-only registry snapshot supplies the window. Main is named only for the actual Main id, a resolved project uses its current registry name and stable chat id; missing, unknown or ambiguous rooms stay explicit. Ephemeral formatter offsets carry the original room/author/direction/transport header into split continuations without parsing or duplicating message bodies. Room draft and correction prompts require that room's meaningful decisions, approvals, outcomes and unresolved commitments, keeping source distinctions (who decided, what was authorized, what stays owed) and one first-person Ouroboros voice. Length adapts to content under the output-token ceiling; no per-room word quota, semantic gate or absent room is imposed. Labels establish provenance, not summary success. The mixed Main recent view opts into the same labels; focused Project rendering, membership and explicit `chat_history` keep their behavior and bytes.
|
||||
|
||||
`knowledge.py` owns note reads, revision-checked writes and indexing; `tools/knowledge.py` exposes them. `knowledge_write(mode="edit", old_str=..., content=...)` replaces ONE body occurrence under the current revision, preserving other bytes and frontmatter. Missing/ambiguous anchors, absent notes and stale revisions refuse before writing. Character/heading deltas reach results, bounded there but complete in history. Edit adds storage capability, not a semantic writer policy; overwrite/append stay unchanged. Blank revision creates only a missing note. Malformed legacy preambles stay readable with uncertain metadata; missing linked notes remain unwritten.
|
||||
`knowledge.py` owns note reads, revision-checked writes and indexing; `tools/knowledge.py` exposes them. `knowledge_write(mode="edit", old_str=..., content=...)` replaces ONE body occurrence under the current revision, other bytes and frontmatter intact; a missing/ambiguous anchor, absent note or stale revision refuses before writing. Character/heading deltas reach results, bounded there but complete in history, whose `source_capture` rows carry `writer`/`route`/`writer_input_ref` and `old_chars`/`new_chars`; `route` is what ANSWERED: `knowledge.observed_route_stamp` reads only the returned usage's physical facts (provider, resolved model, Claudexor profile/account fingerprint), else `unknown` (old rows, fakes); the configured/requested route fills no gap; a merge forwards only the LAST call's stamp; a nomination takes its correction call's route. The reader-less `knowledge_journal.jsonl` telemetry is gone. Edit adds storage capability, not a semantic writer policy; overwrite/append unchanged; a blank revision creates only a missing note; malformed legacy preambles read with uncertain metadata; missing linked notes stay unwritten.
|
||||
|
||||
Ouroboros remains one identity across Main, project rooms, and Background Consciousness: unified dialogue memory remains available to the one agent, while an executing project task preferentially receives its own thread, journal, workpad, and project knowledge. `project_facts.py` routes project facts to `projects/<id>/knowledge`; subagents inherit the root's resolved project id and never derive a new one; identity and the scratchpad are one canonical pair written from every room, with no per-project copy. Project `journal.jsonl` records curated milestones and `workpad.md` retains active working context (`tools/project_journal.py`); focused context includes the workpad in full and recent journal rows with a visible pointer to older entries. On root completion, only high-signal blockers, questions, and interface contracts are mirrored once from the ephemeral task-tree ledger into the durable journal, and a finished root whose effective working tree is not the registered `working_dir` writes one typed "work lives at <path> @ <sha>" journal row from facts the task record already holds. A project digest gives consciousness a concise completion signal without pretending to be the raw project memory.
|
||||
Ouroboros stays one identity across Main, project rooms and Background Consciousness: unified dialogue memory stays available to the one agent; an executing project task prefers its own thread, journal, workpad and project knowledge. `project_facts.py` routes project facts to `projects/<id>/knowledge`; subagents inherit the root's resolved project id, never deriving a new one; identity and the scratchpad are one canonical pair written from every room, no per-project copy. Project `journal.jsonl` records curated milestones and `workpad.md` keeps active working context (`tools/project_journal.py`); focused context holds the whole workpad and recent journal rows with a visible pointer to older entries. On root completion, only high-signal blockers, questions and interface contracts are mirrored once from the ephemeral task-tree ledger into the durable journal, and a finished root whose effective working tree is not the registered `working_dir` writes one typed "work lives at <path> @ <sha>" journal row from facts the task record holds. A project digest gives consciousness a concise completion signal, never posing as the raw project memory.
|
||||
|
||||
### Skills and extensions
|
||||
|
||||
|
|
|
|||
|
|
@ -7,18 +7,10 @@ from typing import Any, Callable, Dict, List, Optional, Tuple
|
|||
|
||||
from ouroboros.contracts.chat_id_policy import is_a2a_chat_id
|
||||
from ouroboros import room_consolidation
|
||||
from ouroboros.utils import (
|
||||
append_jsonl,
|
||||
atomic_write_json,
|
||||
replace_atomic,
|
||||
utc_now_iso,
|
||||
read_text,
|
||||
)
|
||||
from ouroboros.utils import append_jsonl, atomic_write_json, replace_atomic, utc_now_iso, read_text
|
||||
|
||||
from ouroboros.platform_layer import (
|
||||
file_lock_exclusive as _lock_ex,
|
||||
file_lock_exclusive_nb as _lock_nb,
|
||||
file_unlock as _unlock,
|
||||
file_lock_exclusive as _lock_ex, file_lock_exclusive_nb as _lock_nb, file_unlock as _unlock,
|
||||
)
|
||||
|
||||
log = logging.getLogger(__name__)
|
||||
|
|
@ -45,6 +37,56 @@ def _consolidation_route() -> Tuple[str, bool]:
|
|||
return resolve_credentialed_model(lane.model), False
|
||||
|
||||
|
||||
def _light_dispatch_binding() -> Dict[str, Any]:
|
||||
"""The EFFECTIVE Light binding a call dispatches on NOW, in dispatch's own field names.
|
||||
|
||||
The configured lane, then the Light account pin, then the live model-wait
|
||||
override for the role — exactly what ``_call_consolidation_llm`` sends, so a
|
||||
key derived here changes whenever the physical dispatch would."""
|
||||
from ouroboros.model_slots import MODEL_ACCOUNTS_KEY, model_role_option
|
||||
from ouroboros.model_wait import current_model_wait
|
||||
|
||||
model, use_local = _consolidation_route()
|
||||
binding: Dict[str, Any] = {"model": model, "use_local": use_local,
|
||||
"model_account_override": model_role_option(MODEL_ACCOUNTS_KEY, "light")}
|
||||
waiter = current_model_wait()
|
||||
if waiter:
|
||||
binding.update(waiter.overrides.get("light", {}))
|
||||
return binding
|
||||
|
||||
|
||||
def _light_route() -> Any:
|
||||
"""The ``era_retry`` dispatch key: the effective binding of ``_light_dispatch_binding``.
|
||||
|
||||
A dispatch key, never a provenance stamp: what actually answered is
|
||||
``knowledge.observed_route_stamp`` over the returned usage. An account pin or
|
||||
a model-wait override changes this key (the paid retry is allowed) exactly
|
||||
because it changes the dispatch; an empty (Auto) account is omitted so a
|
||||
record keyed before the account joined still holds on an unchanged route."""
|
||||
try:
|
||||
binding = _light_dispatch_binding()
|
||||
except Exception:
|
||||
return "unknown"
|
||||
account = binding.get("model_account_override") or ""
|
||||
return {"model": binding["model"], "use_local": binding["use_local"],
|
||||
**({"model_account_override": account} if account else {})}
|
||||
|
||||
|
||||
def _route_stamp(usage: Any) -> Any:
|
||||
"""The route the usage says answered, for a history stamp; unknown without a physical fact."""
|
||||
from ouroboros.knowledge import observed_route_stamp
|
||||
|
||||
return observed_route_stamp(usage)
|
||||
|
||||
|
||||
def _emit_event(logs_dir: pathlib.Path, kind: str, **fields: Any) -> None:
|
||||
"""A memory-maintenance outcome is a typed fact beside the chat it concerns, never silence (I4)."""
|
||||
try:
|
||||
append_jsonl(logs_dir / "events.jsonl", {"ts": utc_now_iso(), "type": kind, **fields})
|
||||
except Exception:
|
||||
log.debug("Failed to emit %s event", kind, exc_info=True)
|
||||
|
||||
|
||||
def retain_memory_source(context: Any, source_id: str, data: bytes, extension: str = "md") -> Dict[str, Any]:
|
||||
"""Use existing immutable source storage with a reader valid after this task."""
|
||||
from ouroboros.artifacts import store_actor_source_bytes, task_artifact_dir_path
|
||||
|
|
@ -123,13 +165,9 @@ def should_consolidate(
|
|||
|
||||
|
||||
def consolidate(
|
||||
chat_path: pathlib.Path,
|
||||
blocks_path: pathlib.Path,
|
||||
meta_path: pathlib.Path,
|
||||
llm_client: Any,
|
||||
identity_text: str = "",
|
||||
*, knowledge_context: Any = None, force_tail: bool = False, compact_chronicle: bool = False,
|
||||
pressure_fits: Optional[Callable[[], bool]] = None,
|
||||
chat_path: pathlib.Path, blocks_path: pathlib.Path, meta_path: pathlib.Path, llm_client: Any,
|
||||
identity_text: str = "", *, knowledge_context: Any = None, force_tail: bool = False,
|
||||
compact_chronicle: bool = False, pressure_fits: Optional[Callable[[], bool]] = None,
|
||||
room_registry_root: Any = None,
|
||||
) -> Optional[Dict[str, Any]]:
|
||||
lock_path = meta_path.parent / ".consolidation.lock"
|
||||
|
|
@ -141,21 +179,18 @@ def consolidate(
|
|||
_lock_nb(lock_fd)
|
||||
except (OSError, BlockingIOError):
|
||||
log.info("Chat block consolidation already running, skipping")
|
||||
_emit_event(chat_path.parent, "consolidation_skipped_locked", lock_path=str(lock_path),
|
||||
task_id=str(getattr(knowledge_context, "task_id", "") or ""))
|
||||
return None
|
||||
|
||||
usage = _run_block_consolidation(
|
||||
source_path=chat_path,
|
||||
blocks_path=blocks_path,
|
||||
meta_path=meta_path,
|
||||
llm_client=llm_client,
|
||||
identity_text=identity_text,
|
||||
knowledge_context=knowledge_context,
|
||||
force_tail=force_tail,
|
||||
room_registry_root=room_registry_root,
|
||||
)
|
||||
source_path=chat_path, blocks_path=blocks_path, meta_path=meta_path, llm_client=llm_client,
|
||||
identity_text=identity_text, knowledge_context=knowledge_context, force_tail=force_tail,
|
||||
room_registry_root=room_registry_root)
|
||||
if (compact_chronicle and not (usage or {}).get("_consolidation_errors")
|
||||
and not (pressure_fits is not None and pressure_fits())):
|
||||
reduced = _compact_chronicle(blocks_path, llm_client, identity_text, knowledge_context)
|
||||
reduced = _compact_chronicle(blocks_path, llm_client, identity_text, knowledge_context,
|
||||
meta_path=meta_path)
|
||||
merged = _merge_consolidation_usage(*([usage] if usage else []), reduced)
|
||||
# A fixed-key merge would drop this receipt; a chronicle-only pass wrote no block.
|
||||
merged["_blocks_written"] = (usage or {}).get("_blocks_written", 0)
|
||||
|
|
@ -281,9 +316,7 @@ def _run_block_consolidation(
|
|||
if not new_entries or (len(new_entries) < BLOCK_SIZE and not force_tail):
|
||||
return None
|
||||
|
||||
total_usage: Dict[str, Any] = {
|
||||
"prompt_tokens": 0, "completion_tokens": 0, "total_tokens": 0, "cost": 0.0,
|
||||
}
|
||||
total_usage: Dict[str, Any] = {"prompt_tokens": 0, "completion_tokens": 0, "total_tokens": 0, "cost": 0.0}
|
||||
# A failed chunk withholds only itself (its earlier sibling chunks are
|
||||
# complete units); the stale-error clear below must know whether THIS run
|
||||
# recorded a failure it would otherwise erase.
|
||||
|
|
@ -342,16 +375,15 @@ def _run_block_consolidation(
|
|||
run_failed = True
|
||||
meta["last_consolidation_error"] = dict(
|
||||
usage["_consolidation_errors"][-1], cursor_offset=last_offset + processed,
|
||||
chat_log_signature=segment_sigs[0], message_count=len(chunk),
|
||||
)
|
||||
chat_log_signature=segment_sigs[0], message_count=len(chunk))
|
||||
|
||||
if block is None:
|
||||
log.warning("Block summary withheld for chunk %d, will retry next cycle", i)
|
||||
break
|
||||
new_blocks.append({
|
||||
"ts": utc_now_iso(), "type": "summary", "message_count": len(chunk), **block,
|
||||
**({"knowledge_entries": usage["_knowledge_entries"]} if usage.get("_knowledge_entries") else {}),
|
||||
})
|
||||
**({"knowledge_entries": usage["_knowledge_entries"], "_nomination_route": _route_stamp(usage)}
|
||||
if usage.get("_knowledge_entries") else {})})
|
||||
processed += len(chunk)
|
||||
|
||||
# Set after the last merge of this stretch: _merge_consolidation_usage forwards
|
||||
|
|
@ -367,13 +399,15 @@ def _run_block_consolidation(
|
|||
atomic_write_json(meta_path, meta)
|
||||
return total_usage
|
||||
|
||||
# The route that nominated a block's entries is a history stamp for the
|
||||
# writes below, never a persisted block field; receipts keep their pairs.
|
||||
nomination_routes = {id(block): block.pop("_nomination_route", "unknown") for block in new_blocks}
|
||||
pending_knowledge = [(block, block.pop("knowledge_entries")) for block in new_blocks
|
||||
if block.get("knowledge_entries")]
|
||||
if pending_knowledge:
|
||||
source = {"path": str(source_path), "generations": segment_sigs,
|
||||
"start_offset": last_offset, "end_offset": last_offset + processed,
|
||||
"nominations": [{"range": block["range"], "entries": entries}
|
||||
for block, entries in pending_knowledge]}
|
||||
"nominations": [{"range": block["range"], "entries": entries} for block, entries in pending_knowledge]}
|
||||
source_id = hashlib.sha256(json.dumps(source, ensure_ascii=False, sort_keys=True).encode()).hexdigest()
|
||||
ref = {"read": {"tool": "read_file", "arguments": {
|
||||
"root": "runtime_data", "path": "memory/knowledge_history.jsonl"}},
|
||||
|
|
@ -396,32 +430,27 @@ def _run_block_consolidation(
|
|||
all_blocks = existing_blocks + new_blocks
|
||||
|
||||
if len(all_blocks) > MAX_SUMMARY_BLOCKS and block is not None:
|
||||
compress_count = min(ERA_COMPRESS_COUNT, len(all_blocks) - 1)
|
||||
old_blocks = all_blocks[:compress_count]
|
||||
remaining = all_blocks[compress_count:]
|
||||
# Gap markers are DURABLE discontinuity facts (BIBLE P1): they keep
|
||||
# their exact chronological positions, and an era may only compress ONE
|
||||
# CONTIGUOUS run of ordinary summary blocks — never a span that bridges
|
||||
# a known discontinuity.
|
||||
run_start = next((i for i, b in enumerate(old_blocks) if not _is_gap_block(b)), None)
|
||||
# Gap markers are DURABLE discontinuity facts (BIBLE P1) that keep their
|
||||
# chronological positions, and an earlier era is a boundary too: an era
|
||||
# compresses ONE CONTIGUOUS run of ordinary summary blocks, never a span
|
||||
# bridging a discontinuity and never a summary of its own summary. The
|
||||
# run is the OLDEST run of up to ERA_COMPRESS_COUNT summary blocks anywhere
|
||||
# before the newest block — eras and gaps ahead of it are skipped, not a
|
||||
# reason to stop compressing (a window of the first four blocks went blind
|
||||
# once those four were eras).
|
||||
run_start = next((i for i, b in enumerate(all_blocks[:-1]) if not _is_run_boundary(b)), None)
|
||||
era = None
|
||||
if run_start is not None:
|
||||
run_end = run_start
|
||||
while run_end < len(old_blocks) and not _is_gap_block(old_blocks[run_end]):
|
||||
while (run_end < len(all_blocks) - 1 and run_end - run_start < ERA_COMPRESS_COUNT
|
||||
and not _is_run_boundary(all_blocks[run_end])):
|
||||
run_end += 1
|
||||
era, era_usage = _compress_blocks_to_era(
|
||||
old_blocks[run_start:run_end], llm_client, identity_text,
|
||||
**({"knowledge_context": knowledge_context} if knowledge_context is not None else {}),
|
||||
)
|
||||
total_usage = _merge_consolidation_usage(total_usage, era_usage)
|
||||
# An era is a COMPRESSION: replace the run only when it is shorter, as
|
||||
# _compact_chronicle already requires. Per-room sections and the
|
||||
# length-adaptive correction pass can make an era longer than the
|
||||
# blocks it summarizes; keeping those blocks loses nothing.
|
||||
if era is not None and len(era.get("content", "")) < sum(len(b.get("content", "")) for b in old_blocks[run_start:run_end]):
|
||||
all_blocks = [
|
||||
*old_blocks[:run_start], era, *old_blocks[run_end:], *remaining,
|
||||
]
|
||||
era, era_usage = _era_for_run(all_blocks[run_start:run_end], meta, source_path.parent,
|
||||
llm_client, identity_text, knowledge_context)
|
||||
if era_usage is not None:
|
||||
total_usage = _merge_consolidation_usage(total_usage, era_usage)
|
||||
if era is not None:
|
||||
all_blocks = [*all_blocks[:run_start], era, *all_blocks[run_end:]]
|
||||
|
||||
_write_locked_json(blocks_path, all_blocks)
|
||||
|
||||
|
|
@ -430,7 +459,9 @@ def _run_block_consolidation(
|
|||
for block, entries in pending_knowledge:
|
||||
block["knowledge_writes"] = _write_knowledge_entries(
|
||||
pathlib.Path(knowledge_context.drive_root) / "memory" / "knowledge",
|
||||
entries, context=knowledge_context)
|
||||
entries, context=knowledge_context, stamp={
|
||||
"writer": "consolidation", "route": nomination_routes.get(id(block), "unknown"),
|
||||
"writer_input_ref": block["knowledge_source_ref"]})
|
||||
published.extend(block["knowledge_writes"])
|
||||
if any(not outcome["ok"] for outcome in block["knowledge_writes"]):
|
||||
append_jsonl(pathlib.Path(knowledge_context.drive_root) / "memory" / "knowledge_history.jsonl", {
|
||||
|
|
@ -470,12 +501,22 @@ def _light_call(llm_client: Any, knowledge_context: Any, model_route: Dict[str,
|
|||
|
||||
def _merge_consolidation_usage(*usages: Dict[str, Any]) -> Dict[str, Any]:
|
||||
"""Combine helper usage without turning absent spend/counters into zero."""
|
||||
from ouroboros.knowledge import observed_route_stamp
|
||||
|
||||
merged: Dict[str, Any] = {}
|
||||
for key in ("prompt_tokens", "completion_tokens", "total_tokens", "cost"):
|
||||
values = [usage.get(key) for usage in usages]
|
||||
merged[key] = None if None in values else sum(values)
|
||||
for key in ("ledger_attempt_ids", "_consolidation_errors"):
|
||||
merged[key] = [value for usage in usages for value in usage.get(key, [])]
|
||||
# The route that answered the LAST send of this unit, and only that one: a
|
||||
# physical usage carries provider/resolved_model, a merged one its forwarded
|
||||
# stamp, and a final call without a physical fact reads unknown — an earlier
|
||||
# call's stamp never masquerades as the final call's.
|
||||
if usages:
|
||||
last = observed_route_stamp(usages[-1])
|
||||
if isinstance(last, dict):
|
||||
merged["_observed_route"] = last
|
||||
return merged
|
||||
|
||||
|
||||
|
|
@ -659,7 +700,11 @@ class KnowledgeReadContext:
|
|||
# The host records what THIS operation actually read. Model-supplied
|
||||
# revision text cannot attest an unread note; absent read permits
|
||||
# creation only, because the common writer requires existing CAS.
|
||||
bound.append({**entry, "scope": address.scope,
|
||||
# Underscore keys are host facts (``_nomination_route`` and any later
|
||||
# one): a model-supplied value is dropped here, so a forged route can
|
||||
# never reach history — the host stamps only after this binding.
|
||||
bound.append({**{key: value for key, value in entry.items() if not str(key).startswith("_")},
|
||||
"scope": address.scope,
|
||||
"expected_revision": self.reads.get((address.scope, address.topic)),
|
||||
"canonical_root": str(address.canonical_root),
|
||||
"task_id": str(getattr(self.context, "task_id", "") or "")})
|
||||
|
|
@ -715,7 +760,6 @@ def _call_consolidation_llm(
|
|||
_failed_route_evidence, _route_calibration_ratio, estimate_context_prompt_tokens, resolve_context_fit_route,
|
||||
)
|
||||
from ouroboros.tools.compact_context import record_context_view
|
||||
from ouroboros.model_slots import MODEL_ACCOUNTS_KEY, model_role_option
|
||||
from ouroboros.model_wait import current_model_wait
|
||||
from ouroboros.provider_models import parse_claudexor_model, provider_for_model
|
||||
|
||||
|
|
@ -812,14 +856,12 @@ def _call_consolidation_llm(
|
|||
"measurement_basis": "canonical_visible_estimate", "strict_bound_proven": False}
|
||||
|
||||
try:
|
||||
model, use_local = _consolidation_route()
|
||||
values = dict(messages=[{"role": "user", "content": prompt}], model=model,
|
||||
# The same effective binding keys era_retry (``_light_route``): a refusal
|
||||
# recorded under one binding never suppresses the retry under another.
|
||||
values = dict(messages=[{"role": "user", "content": prompt}],
|
||||
model_role="light", tools=knowledge.tools if knowledge else None,
|
||||
reasoning_effort=reasoning_effort, max_tokens=16384,
|
||||
use_local=use_local,
|
||||
model_account_override=model_role_option(MODEL_ACCOUNTS_KEY, "light"))
|
||||
if waiter:
|
||||
values.update(waiter.overrides.get("light", {}))
|
||||
**_light_dispatch_binding())
|
||||
# Carry part-to-part evidence only on initial preparation. A wait's
|
||||
# reprepare without an observed receipt rediscovers Auto after rotation.
|
||||
try:
|
||||
|
|
@ -949,33 +991,107 @@ def _is_gap_block(block: Any) -> bool:
|
|||
return isinstance(block, dict) and bool(block.get("gap_id") or "[MEMORY GAP]" in str(block.get("content") or ""))
|
||||
|
||||
|
||||
def _is_run_boundary(block: Any) -> bool:
|
||||
"""Gaps and earlier eras bound a run: an era is built from summary blocks, never from an era."""
|
||||
return _is_gap_block(block) or (isinstance(block, dict) and block.get("type") == "era")
|
||||
|
||||
|
||||
ERA_RETRY_MAX_RUNS = 16 # refusals remembered per meta; the oldest run's record ages out first
|
||||
|
||||
|
||||
def _era_retry_runs(meta: Dict[str, Any]) -> Dict[str, Dict[str, Any]]:
|
||||
"""``era_retry`` as {source_sha256: {route, observed_route}}; a legacy single record reads as one entry."""
|
||||
retry = meta.get("era_retry")
|
||||
if not isinstance(retry, dict):
|
||||
return {}
|
||||
if "source_sha256" in retry: # legacy single-record shape
|
||||
return {str(retry["source_sha256"]): {key: value for key, value in retry.items() if key != "source_sha256"}}
|
||||
return {str(key): dict(value) for key, value in retry.items() if isinstance(value, dict)}
|
||||
|
||||
|
||||
def _era_for_run(run: List[Dict[str, Any]], meta: Dict[str, Any], logs_dir: pathlib.Path, llm_client: Any,
|
||||
identity_text: str, context: Any) -> Tuple[Optional[Dict[str, Any]], Optional[Dict[str, Any]]]:
|
||||
"""The era of one run when it is a COMPRESSION; otherwise a recorded, visible refusal.
|
||||
|
||||
Per-room sections and the length-adaptive correction can make an era longer than
|
||||
its blocks; keeping the blocks loses nothing. The refusal is ``era_retry`` in meta
|
||||
(keyed by source hash, PER RUN — a chronicle holds several runs between gaps and
|
||||
eras, and one run's refusal or success must not erase another's; each record keeps the
|
||||
dispatch key, read AFTER the call because an owner switch during a wait inside it rebinds
|
||||
the role before the paid send, and the route that ANSWERED): the same source is not paid for
|
||||
again while the binding a call would dispatch on now is the one that refused, and every
|
||||
refusal, attempted or not, is an ``era_not_shorter`` event. Returns ``(era or None, usage or None without a call)``."""
|
||||
fact = {"source_sha256": hashlib.sha256(json.dumps(run, ensure_ascii=False, sort_keys=True).encode()).hexdigest(),
|
||||
"route": _light_route(), "blocks": len(run), "source_chars": sum(len(b.get("content", "")) for b in run)}
|
||||
runs = _era_retry_runs(meta)
|
||||
prior = runs.get(fact["source_sha256"])
|
||||
if prior is not None and prior.get("route") == fact["route"]:
|
||||
_emit_event(logs_dir, "era_not_shorter", attempted=False, **fact)
|
||||
return None, None
|
||||
era, usage = _compress_blocks_to_era(run, llm_client, identity_text,
|
||||
**({"knowledge_context": context} if context is not None else {}))
|
||||
if era is not None and len(era.get("content", "")) >= fact["source_chars"]:
|
||||
# The dispatch key the next attempt compares against, read AFTER the call: a switch inside it
|
||||
# rebound the role, so a pre-call key would suppress what never answered and repay what did.
|
||||
fact["route"] = _light_route()
|
||||
runs.pop(fact["source_sha256"], None)
|
||||
runs[fact["source_sha256"]] = {"route": fact["route"], "observed_route": _route_stamp(usage)}
|
||||
while len(runs) > ERA_RETRY_MAX_RUNS:
|
||||
runs.pop(next(iter(runs)))
|
||||
meta["era_retry"] = runs
|
||||
_emit_event(logs_dir, "era_not_shorter", attempted=True, era_chars=len(era["content"]),
|
||||
observed_route=_route_stamp(usage), **fact)
|
||||
return None, usage
|
||||
if era is not None and runs.pop(fact["source_sha256"], None) is not None:
|
||||
if runs:
|
||||
meta["era_retry"] = runs
|
||||
else:
|
||||
meta.pop("era_retry", None)
|
||||
return era, usage
|
||||
|
||||
|
||||
def _compact_chronicle(blocks_path: pathlib.Path, llm_client: Any,
|
||||
identity_text: str, context: Any) -> Dict[str, Any]:
|
||||
"""Reduce every contiguous historical span, preserving gaps and exact sources."""
|
||||
identity_text: str, context: Any, *, meta_path: Optional[pathlib.Path] = None) -> Dict[str, Any]:
|
||||
"""Reduce every contiguous run of summary blocks, preserving gaps, earlier eras and exact sources.
|
||||
|
||||
The pressure pass consults and records the SAME ``era_retry`` as the ordinary
|
||||
run (``meta_path``): a run that was not shorter on this route is not paid for
|
||||
again by the next pressure pass on unchanged input. Without a meta path (a caller
|
||||
that has none) it reads and writes no durable retry metadata: every run is still paid for."""
|
||||
blocks = _load_blocks(blocks_path)
|
||||
meta: Dict[str, Any] = {}
|
||||
if meta_path is not None:
|
||||
try:
|
||||
meta = _load_meta(meta_path)
|
||||
except Exception:
|
||||
# An unreadable meta is the caller's typed maintenance gap; the
|
||||
# chronicle pass must not write a rebuilt meta over it.
|
||||
log.warning("Chronicle pass cannot read dialogue meta; era_retry not consulted", exc_info=True)
|
||||
meta_path = None
|
||||
retry_before = json.dumps(meta.get("era_retry"), sort_keys=True)
|
||||
reduced, usages, start = [], [], 0
|
||||
while start < len(blocks):
|
||||
if _is_gap_block(blocks[start]):
|
||||
if _is_run_boundary(blocks[start]):
|
||||
reduced.append(blocks[start])
|
||||
start += 1
|
||||
continue
|
||||
end = start + 1
|
||||
while end < len(blocks) and not _is_gap_block(blocks[end]):
|
||||
while end < len(blocks) and not _is_run_boundary(blocks[end]):
|
||||
end += 1
|
||||
run = blocks[start:end]
|
||||
era, usage = _compress_blocks_to_era(run, llm_client, identity_text, context)
|
||||
usages.append(usage)
|
||||
if era and len(era["content"]) < sum(len(b.get("content", "")) for b in run):
|
||||
reduced.append(era)
|
||||
else:
|
||||
reduced.extend(run)
|
||||
if usage.get("_consolidation_errors"):
|
||||
era, usage = _era_for_run(run, meta, blocks_path.parent.parent / "logs", llm_client, identity_text, context)
|
||||
if usage is not None: # a recorded refusal makes no call and has no usage
|
||||
usages.append(usage)
|
||||
reduced.extend([era] if era is not None else run)
|
||||
if (usage or {}).get("_consolidation_errors"):
|
||||
reduced.extend(blocks[end:])
|
||||
break
|
||||
start = end
|
||||
if reduced != blocks:
|
||||
_mutate_locked_json_list(blocks_path, lambda live:
|
||||
reduced + live[len(blocks):] if live[:len(blocks)] == blocks else live)
|
||||
if meta_path is not None and json.dumps(meta.get("era_retry"), sort_keys=True) != retry_before:
|
||||
atomic_write_json(meta_path, meta)
|
||||
return _merge_consolidation_usage(*usages)
|
||||
|
||||
|
||||
|
|
@ -1054,7 +1170,8 @@ def maintain_memory_pressure(memory: Any, llm_client: Any, context: Any, *,
|
|||
if raw.strip():
|
||||
try:
|
||||
entries = knowledge.bind_entries(json.loads(raw).get("knowledge_entries"))
|
||||
action["writes"] = _write_knowledge_entries(shelf, entries, context=context)
|
||||
action["writes"] = _write_knowledge_entries(shelf, entries, context=context, stamp={
|
||||
"writer": "knowledge_maintenance", "route": _route_stamp(usage), "writer_input_ref": source_ref})
|
||||
except (ValueError, TypeError, AttributeError) as exc:
|
||||
action["error"] = str(exc)
|
||||
actions.append(action)
|
||||
|
|
@ -1121,17 +1238,7 @@ def _load_blocks(path: pathlib.Path) -> List[Dict[str, Any]]:
|
|||
log.error("Corrupt blocks file %s — quarantined to %s, starting fresh", path, quarantine)
|
||||
except OSError:
|
||||
log.error("Corrupt blocks file %s — quarantine failed, starting fresh", path, exc_info=True)
|
||||
try:
|
||||
from ouroboros.utils import append_jsonl
|
||||
|
||||
append_jsonl(path.parent.parent / "logs" / "events.jsonl", {
|
||||
"ts": utc_now_iso(),
|
||||
"type": "memory_store_corrupt",
|
||||
"path": str(path),
|
||||
"quarantine": str(quarantine),
|
||||
})
|
||||
except Exception:
|
||||
log.debug("Failed to emit memory_store_corrupt event", exc_info=True)
|
||||
_emit_event(path.parent.parent / "logs", "memory_store_corrupt", path=str(path), quarantine=str(quarantine))
|
||||
return []
|
||||
|
||||
|
||||
|
|
@ -1300,10 +1407,7 @@ def should_consolidate_scratchpad(memory: Any) -> bool:
|
|||
|
||||
|
||||
def consolidate_scratchpad(
|
||||
memory: Any,
|
||||
knowledge_dir: pathlib.Path,
|
||||
llm_client: Any,
|
||||
identity_text: str = "",
|
||||
memory: Any, knowledge_dir: pathlib.Path, llm_client: Any, identity_text: str = "",
|
||||
*, pressure: bool = False, knowledge_context: Any = None,
|
||||
) -> Optional[Dict[str, Any]]:
|
||||
blocks = memory.load_scratchpad_blocks()
|
||||
|
|
@ -1315,12 +1419,8 @@ def consolidate_scratchpad(
|
|||
|
||||
|
||||
def _consolidate_scratchpad_blocks(
|
||||
memory: Any,
|
||||
blocks: List[Dict[str, Any]],
|
||||
knowledge_dir: pathlib.Path,
|
||||
llm_client: Any,
|
||||
identity_text: str,
|
||||
*, pressure: bool = False, knowledge_context: Any = None,
|
||||
memory: Any, blocks: List[Dict[str, Any]], knowledge_dir: pathlib.Path, llm_client: Any,
|
||||
identity_text: str, *, pressure: bool = False, knowledge_context: Any = None,
|
||||
) -> Optional[Dict[str, Any]]:
|
||||
total_chars = sum(len(b.get("content", "")) for b in blocks)
|
||||
if total_chars <= SCRATCHPAD_CONSOLIDATION_THRESHOLD and not pressure:
|
||||
|
|
@ -1362,6 +1462,7 @@ Respond with JSON only (no fences), after any useful knowledge reads:
|
|||
"""
|
||||
|
||||
usage: Dict[str, Any] = {}
|
||||
outcome, source_entry_id, writes, new_blocks = "failed", "", [], blocks
|
||||
try:
|
||||
from ouroboros.tools.registry import ToolContext
|
||||
|
||||
|
|
@ -1373,6 +1474,7 @@ Respond with JSON only (no fences), after any useful knowledge reads:
|
|||
knowledge=knowledge)
|
||||
raw = raw.strip()
|
||||
if not raw:
|
||||
outcome = "call_failed" if usage.get("_consolidation_errors") else "empty_response"
|
||||
return usage
|
||||
if raw.startswith("```"):
|
||||
raw = raw.split("\n", 1)[-1].rsplit("```", 1)[0].strip()
|
||||
|
|
@ -1382,45 +1484,34 @@ Respond with JSON only (no fences), after any useful knowledge reads:
|
|||
compressed_text = result.get("compressed_block", "")
|
||||
if not compressed_text or not compressed_text.strip():
|
||||
log.warning("Scratchpad block consolidation returned empty, skipping")
|
||||
outcome = "empty_block"
|
||||
return usage
|
||||
if pressure and len(compressed_text) >= sum(len(b.get("content", "")) for b in old_blocks):
|
||||
outcome = "not_shorter"
|
||||
return usage # an authored expansion is not pressure relief
|
||||
|
||||
entries = knowledge.bind_entries(result.get("knowledge_entries"))
|
||||
|
||||
compressed_block = {
|
||||
"ts": utc_now_iso(),
|
||||
"source": "consolidation",
|
||||
"content": compressed_text.strip(),
|
||||
}
|
||||
|
||||
source_bytes = json.dumps(
|
||||
old_blocks, ensure_ascii=False, sort_keys=True, separators=(",", ":"),
|
||||
).encode("utf-8")
|
||||
compressed_block = {"ts": utc_now_iso(), "source": "consolidation", "content": compressed_text.strip()}
|
||||
source_bytes = json.dumps(old_blocks, ensure_ascii=False, sort_keys=True, separators=(",", ":")).encode("utf-8")
|
||||
source_entry_id = "scratchpad-consolidation:" + hashlib.sha256(source_bytes).hexdigest()
|
||||
source_ref = memory.scratchpad_journal_source_ref(source_entry_id)
|
||||
if not append_jsonl(memory.journal_path(), {
|
||||
"ts": utc_now_iso(),
|
||||
"type": "blocks_consolidated",
|
||||
"entry_id": source_entry_id,
|
||||
"source_blocks": old_blocks,
|
||||
"source_ref": source_ref,
|
||||
"knowledge_entries": entries,
|
||||
}):
|
||||
"ts": utc_now_iso(), "type": "blocks_consolidated", "entry_id": source_entry_id,
|
||||
"source_blocks": old_blocks, "source_ref": source_ref, "knowledge_entries": entries}):
|
||||
log.error("Scratchpad consolidation source journal write failed; preserving blocks")
|
||||
outcome = "journal_unavailable"
|
||||
return usage
|
||||
compressed_block["metadata"] = {"source_ref": source_ref}
|
||||
writes = _write_knowledge_entries(knowledge_dir, entries, context=context)
|
||||
writes = _write_knowledge_entries(knowledge_dir, entries, context=context, stamp={
|
||||
"writer": "scratchpad_consolidation", "route": _route_stamp(usage), "writer_input_ref": source_ref})
|
||||
if writes:
|
||||
compressed_block["metadata"]["knowledge_writes"] = writes
|
||||
if any(not row["ok"] for row in writes):
|
||||
compressed_block["content"] += (
|
||||
"\n\nSome nominated knowledge updates were not published; their complete "
|
||||
"proposals and original episode remain in the source journal referenced by this block.")
|
||||
append_jsonl(memory.journal_path(), {
|
||||
"ts": utc_now_iso(), "type": "knowledge_writes_incomplete",
|
||||
"source_ref": source_ref, "knowledge_writes": writes,
|
||||
})
|
||||
append_jsonl(memory.journal_path(), {"ts": utc_now_iso(), "type": "knowledge_writes_incomplete",
|
||||
"source_ref": source_ref, "knowledge_writes": writes})
|
||||
|
||||
# Merge-aware replace UNDER the write lock: blocks appended DURING the
|
||||
# slow LLM call live only on disk — building the new list from the
|
||||
|
|
@ -1432,6 +1523,7 @@ Respond with JSON only (no fences), after any useful knowledge reads:
|
|||
return [compressed_block] + live_blocks[len(old_blocks):]
|
||||
|
||||
new_blocks = memory.mutate_scratchpad_blocks(_merge_survivors)
|
||||
outcome = "replaced" if new_blocks[:1] == [compressed_block] else "source_changed"
|
||||
|
||||
log.info("Scratchpad blocks consolidated: %d blocks (%d chars) -> %d blocks (%d chars)",
|
||||
len(blocks), total_chars,
|
||||
|
|
@ -1442,14 +1534,34 @@ Respond with JSON only (no fences), after any useful knowledge reads:
|
|||
from ouroboros.llm_claudexor import propagate_model_error
|
||||
propagate_model_error(e)
|
||||
log.error("Scratchpad block consolidation failed: %s", e, exc_info=True)
|
||||
return {**usage, "_consolidation_errors": [*usage.get("_consolidation_errors", []), {
|
||||
usage = {**usage, "_consolidation_errors": [*usage.get("_consolidation_errors", []), {
|
||||
"kind": "scratchpad_consolidation_failed", "message": f"{type(e).__name__}: {e}"}]}
|
||||
return usage
|
||||
finally:
|
||||
# Every exit above names its outcome; the chat_block_consolidation row is the model.
|
||||
errors = usage.get("_consolidation_errors") or []
|
||||
_emit_event(pathlib.Path(memory.drive_root) / "logs", "scratchpad_consolidation", outcome=outcome,
|
||||
pressure=pressure, blocks_before=len(blocks), chars_before=total_chars,
|
||||
compressed_blocks=len(old_blocks), blocks_after=len(new_blocks),
|
||||
chars_after=sum(len(b.get("content", "")) for b in new_blocks), source_entry_id=source_entry_id,
|
||||
knowledge_writes={"ok": sum(w["ok"] for w in writes), "failed": sum(not w["ok"] for w in writes)},
|
||||
last_error_kind=(errors[-1] or {}).get("kind") if errors else None,
|
||||
accounted_upper_bound_usd=round(float(usage["cost"]), 6) if usage.get("cost") is not None else None)
|
||||
|
||||
|
||||
def _write_knowledge_entries(
|
||||
knowledge_dir: pathlib.Path, entries: List[Dict[str, Any]], *, context: Any = None,
|
||||
stamp: Optional[Dict[str, Any]] = None,
|
||||
) -> List[Dict[str, Any]]:
|
||||
"""Publish only source-aware nominations through the common note writer."""
|
||||
"""Publish only source-aware nominations through the common note writer.
|
||||
|
||||
``stamp`` is the caller's ``writer``/``route``/``writer_input_ref`` history
|
||||
stamp (see ``write_knowledge_note``); an unnamed caller leaves ``unknown``. An
|
||||
entry carrying its own ``_nomination_route`` outranks the caller's block-level
|
||||
``route``: provenance is per nomination. That field is HOST-authored only —
|
||||
``KnowledgeReadContext.bind_entries`` strips every underscore key a model
|
||||
supplied, and the room seam stamps it after binding from the correction
|
||||
call's own usage — so the writer never trusts model output for it."""
|
||||
from ouroboros.knowledge import KnowledgeAddress, sanitize_topic, write_knowledge_note
|
||||
from ouroboros.tools.knowledge import _address, _record_backlog_history
|
||||
|
||||
|
|
@ -1475,8 +1587,11 @@ def _write_knowledge_entries(
|
|||
outcomes.append({"topic": topic, "scope": "global", "ok": merged >= 0,
|
||||
"reason": "backlog_merge" if merged >= 0 else "unparseable_backlog"})
|
||||
continue
|
||||
entry_stamp = dict(stamp or {})
|
||||
if entry.get("_nomination_route") is not None:
|
||||
entry_stamp["route"] = entry["_nomination_route"]
|
||||
result = write_knowledge_note(address, content, expected_revision=entry.get("expected_revision"),
|
||||
task_id=str(entry.get("task_id") or ""))
|
||||
task_id=str(entry.get("task_id") or ""), **entry_stamp)
|
||||
outcomes.append({"topic": topic, "scope": address.scope, "ok": result.ok,
|
||||
"reason": result.reason,
|
||||
"source_ref": result.current.source_ref() if result.current else None})
|
||||
|
|
|
|||
|
|
@ -273,6 +273,19 @@ def _memory_health_lines(env: Any) -> List[str]:
|
|||
f"WARNING: LAST DIALOGUE CONSOLIDATION FAILED — kind={error.get('kind') or 'unknown'} "
|
||||
f"at cursor {error.get('cursor_offset')}"
|
||||
)
|
||||
from ouroboros.consolidator import _era_retry_runs
|
||||
runs = _era_retry_runs(meta)
|
||||
for shown, (source_sha256, record) in enumerate(runs.items()):
|
||||
if shown == 3:
|
||||
lines.append(f"WARNING: DIALOGUE ERA COMPRESSION WITHHELD — {len(runs) - 3} more run(s) recorded in era_retry")
|
||||
break
|
||||
route = record.get("route")
|
||||
lines.append(
|
||||
f"WARNING: DIALOGUE ERA COMPRESSION WITHHELD — the era for source run "
|
||||
f"{source_sha256[:12]} on route "
|
||||
f"{route.get('model') if isinstance(route, dict) else route} was not shorter than its blocks; "
|
||||
"blocks retained, no paid repeat until that run or the route changes"
|
||||
)
|
||||
return lines
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -13,7 +13,7 @@ from collections import Counter
|
|||
from contextlib import contextmanager
|
||||
from dataclasses import dataclass, replace
|
||||
from pathlib import Path, PurePosixPath
|
||||
from typing import Any, Mapping
|
||||
from typing import Any, Dict, Mapping
|
||||
from urllib.parse import quote, unquote, urlsplit
|
||||
|
||||
import yaml
|
||||
|
|
@ -23,11 +23,49 @@ from ouroboros.platform_layer import file_lock_exclusive, file_unlock
|
|||
from ouroboros.utils import append_jsonl, utc_now_iso, write_bytes_atomic
|
||||
|
||||
INDEX_FILE = "index-full.md"
|
||||
UNKNOWN_STAMP = "unknown" # a history stamp the writer could not name; legacy rows read the same way
|
||||
OVERVIEW_TOPIC = "overview"
|
||||
_INDEX_HEADER = "# Knowledge Base Index\n<!-- ouroboros:knowledge-index:1 -->\n\n"
|
||||
_LEGACY_INDEX_MARKER = "\n<!-- ouroboros:legacy-knowledge-index -->\n"
|
||||
|
||||
|
||||
def observed_route_stamp(usage: Any) -> Any:
|
||||
"""The route a physical usage row says ANSWERED, as the history ``route`` stamp.
|
||||
|
||||
Every wire lane stamps ``provider`` and ``resolved_model`` on the usage it
|
||||
returns (a model-wait override or account rotation changes them, the
|
||||
configured route does not), and the Claudexor lane adds its ``route`` with
|
||||
the serving ``source``/``account``. Only those physical facts are read: a
|
||||
usage without any — a released send, a fake in a test — is the honest
|
||||
``unknown``, and a partial one leaves the missing field ``unknown``; the
|
||||
configured or requested route never fills a gap, so no caller argument can.
|
||||
An already derived stamp (``_observed_route``, forwarded by usage merges as
|
||||
the LAST call's stamp only) is returned as is.
|
||||
"""
|
||||
if not isinstance(usage, dict):
|
||||
return UNKNOWN_STAMP
|
||||
prior = usage.get("_observed_route")
|
||||
if isinstance(prior, dict):
|
||||
return prior
|
||||
provider, resolved = usage.get("provider"), usage.get("resolved_model")
|
||||
if not provider and not resolved:
|
||||
return UNKNOWN_STAMP
|
||||
stamp: Dict[str, Any] = {"provider": str(provider or UNKNOWN_STAMP), "model": str(resolved or UNKNOWN_STAMP)}
|
||||
served = usage.get("claudexor")
|
||||
if isinstance(served, dict) and isinstance(served.get("route"), dict):
|
||||
route = served["route"]
|
||||
if route.get("source"):
|
||||
stamp["source"] = route["source"]
|
||||
# The engine's served route names the account as ``credentialProfileId``
|
||||
# (+ ``accountFingerprint``); ``account`` is the legacy/request spelling.
|
||||
account = route.get("credentialProfileId") or route.get("account")
|
||||
if account:
|
||||
stamp["account"] = str(account)
|
||||
if route.get("accountFingerprint"):
|
||||
stamp["account_fingerprint"] = str(route["accountFingerprint"])
|
||||
return stamp
|
||||
|
||||
|
||||
def sanitize_topic(topic: str) -> str:
|
||||
"""Keep a shelf-relative topic identity, including useful nested names."""
|
||||
if not isinstance(topic, str) or not topic.strip():
|
||||
|
|
@ -328,8 +366,17 @@ class KnowledgeWriteResult:
|
|||
def write_knowledge_note(
|
||||
address: KnowledgeAddress, content: str, mode: str = "overwrite",
|
||||
expected_revision: str | None = None, task_id: str = "", old_str: str | None = None,
|
||||
*, writer: str = "", route: Any = None, writer_input_ref: Any = None,
|
||||
) -> KnowledgeWriteResult:
|
||||
"""Publish a note against the actual current source, with no inference lock."""
|
||||
"""Publish a note against the actual current source, with no inference lock.
|
||||
|
||||
``writer`` names the seam that authored ``content`` (turn, consolidation,
|
||||
scratchpad_consolidation, reflection, knowledge_maintenance), ``route`` the
|
||||
model route it ran on and ``writer_input_ref`` what it saw. They are host
|
||||
facts stamped on the history row, never on the note body; a caller that
|
||||
cannot name one leaves the honest ``unknown``, which is also how rows written
|
||||
before the stamp existed read.
|
||||
"""
|
||||
if mode not in {"overwrite", "append", "edit"} or not isinstance(content, str):
|
||||
raise ValueError("content must be Markdown text; mode must be overwrite, append or edit")
|
||||
if mode != "edit" and old_str is not None:
|
||||
|
|
@ -380,6 +427,9 @@ def write_knowledge_note(
|
|||
# failed publication returns its actual current source, never success.
|
||||
history = {"ts": utc_now_iso(), "task_id": task_id, "topic": address.topic, "mode": mode,
|
||||
"address": address.as_dict(), "publication": "source_capture",
|
||||
"writer": writer or UNKNOWN_STAMP, "route": route or UNKNOWN_STAMP,
|
||||
"writer_input_ref": writer_input_ref or UNKNOWN_STAMP,
|
||||
"old_chars": len(old_text), "new_chars": len(updated.text),
|
||||
"old_sha256": hashlib.sha256(current.raw).hexdigest() if current and current.raw else "",
|
||||
"new_sha256": updated.revision if raw else "", "old_content": old_text,
|
||||
"new_content": updated.text, "source_ref": updated.source_ref(), "delta": delta}
|
||||
|
|
@ -397,13 +447,4 @@ def write_knowledge_note(
|
|||
observed = None
|
||||
return KnowledgeWriteResult(False, "publication_incomplete", observed, revision,
|
||||
delta if observed is not None and observed.raw == raw else None)
|
||||
try:
|
||||
append_jsonl(address.shelf.parent / "knowledge_journal.jsonl", {
|
||||
"ts": utc_now_iso(), "task_id": task_id, "topic": address.topic, "mode": mode,
|
||||
"address": address.as_dict(), "revision": updated.revision,
|
||||
"file_kb": len(raw) / 1024,
|
||||
"total_knowledge_kb": round(sum(p.stat().st_size for p in address.shelf.rglob("*.md")) / 1024, 2),
|
||||
}, ensure_record_boundary=True)
|
||||
except OSError:
|
||||
pass # Size telemetry is not source/history publication authority.
|
||||
return KnowledgeWriteResult(True, "saved", updated, revision, delta)
|
||||
|
|
|
|||
|
|
@ -22,6 +22,7 @@ from ouroboros.deadline_utils import (
|
|||
seconds_until,
|
||||
main_transport_timeout_sec as _main_transport_timeout,
|
||||
)
|
||||
from ouroboros.knowledge import observed_route_stamp
|
||||
from ouroboros.llm import LLMClient, LocalContextTooLargeError, add_usage
|
||||
from ouroboros.llm_claudexor import propagate_model_error, presence_refusal_unstarted
|
||||
from ouroboros.llm_substitution import same_route_refusal, stamp_substitutions
|
||||
|
|
@ -1437,9 +1438,8 @@ def call_llm_with_retry(
|
|||
for stale in ("_last_llm_error", "_last_llm_error_kind", "_last_llm_retry_same_request",
|
||||
"_last_llm_status_code", "_last_llm_provider_code", "_last_llm_resource_refusal"):
|
||||
accumulated_usage.pop(stale, None)
|
||||
cost, display_model, provider, cost_estimated = _normalize_usage_cost(
|
||||
usage, model=model, use_local=use_local,
|
||||
)
|
||||
cost, display_model, provider, cost_estimated = _normalize_usage_cost(usage, model=model, use_local=use_local)
|
||||
accumulated_usage["_observed_route"] = observed_route_stamp(usage)
|
||||
add_usage(accumulated_usage, usage)
|
||||
fold_retrieval_usage(accumulated_usage, usage)
|
||||
response_ref = persist_observed_call(
|
||||
|
|
|
|||
|
|
@ -2,10 +2,11 @@
|
|||
|
||||
from __future__ import annotations
|
||||
|
||||
import functools
|
||||
import json
|
||||
import logging
|
||||
import pathlib
|
||||
from typing import Any, Dict, List, Optional
|
||||
from typing import Any, Callable, Dict, List, Optional
|
||||
|
||||
from ouroboros._outcome_tool_errors import _OK_TOOL_STATUSES
|
||||
from ouroboros.utils import utc_now_iso, append_jsonl, write_text_atomic
|
||||
|
|
@ -350,11 +351,46 @@ def _extract_trailing_json(text: str, marker: str) -> tuple[str, Optional[list]]
|
|||
return remainder, value if isinstance(value, list) else None
|
||||
|
||||
|
||||
def _validate_memory_actions(raw: Any, task_id: str) -> List[Dict[str, Any]]:
|
||||
"""Keep only well-formed, allowed-type memory actions (max 3)."""
|
||||
def record_memory_action_skip(events: pathlib.Path, action: Dict[str, Any], reason: str, *,
|
||||
project_id: str = "", input_ref: Any = None) -> None:
|
||||
"""A lesson the host declines is a fact, not silence (I4).
|
||||
|
||||
One writer for every rejection seam — the validator on the model's raw output
|
||||
and ``apply_memory_actions`` on a bound action — so the event names the reason
|
||||
and, as ``input_ref``, what the seam that dropped the lesson had retained: the
|
||||
validator's exact task-input prompt (the rejected reflection text itself is not
|
||||
retained), a bound action's exact task-source copy of its reflection entry, or
|
||||
the canonical log pointer. It can warn, never raise: an audit-write failure
|
||||
must not discard the independent later lessons of the same batch."""
|
||||
try:
|
||||
recorded = append_jsonl(events, {"ts": utc_now_iso(), "type": "reflection_memory_action_skipped",
|
||||
"task_id": str(action.get("task_id") or ""), "project_id": project_id,
|
||||
"action_type": str(action.get("type") or "")[:80], "reason": reason,
|
||||
"content_chars": len(str(action.get("content") or "")),
|
||||
"input_ref": input_ref})
|
||||
if not recorded:
|
||||
log.warning("Reflection memory skip event was not recorded: task=%s reason=%s",
|
||||
action.get("task_id"), reason)
|
||||
except Exception:
|
||||
log.warning("Reflection memory skip event could not be written: task=%s reason=%s",
|
||||
action.get("task_id"), reason, exc_info=True)
|
||||
|
||||
|
||||
def _validate_memory_actions(raw: Any, task_id: str, *,
|
||||
on_skip: Optional[Callable[[Dict[str, Any], str], None]] = None) -> List[Dict[str, Any]]:
|
||||
"""Keep only well-formed, allowed-type memory actions (max 3).
|
||||
|
||||
``on_skip(action, reason)`` hears every dropped dict item (``unsupported_type``,
|
||||
``empty_content``, ``missing_topic``): this is the seam the model's output
|
||||
actually crosses, so the skip event fires here, before any action is bound."""
|
||||
out: List[Dict[str, Any]] = []
|
||||
if not isinstance(raw, list):
|
||||
return out
|
||||
|
||||
def skip(item: Dict[str, Any], action_type: str, reason: str) -> None:
|
||||
if on_skip is not None:
|
||||
on_skip({"type": action_type, "content": str(item.get("content") or ""), "task_id": task_id}, reason)
|
||||
|
||||
for item in raw[:10]:
|
||||
if len(out) >= 3:
|
||||
break
|
||||
|
|
@ -362,15 +398,18 @@ def _validate_memory_actions(raw: Any, task_id: str) -> List[Dict[str, Any]]:
|
|||
continue
|
||||
action_type = str(item.get("type") or "").strip()
|
||||
if action_type not in _ALLOWED_MEMORY_ACTION_TYPES:
|
||||
skip(item, action_type, "unsupported_type")
|
||||
continue
|
||||
content = (str(item.get("content") or "") if action_type == "knowledge_write"
|
||||
else _truncate_with_notice(item.get("content", ""), 1200)).strip()
|
||||
if not content:
|
||||
skip(item, action_type, "empty_content")
|
||||
continue
|
||||
action: Dict[str, Any] = {"type": action_type, "content": content, "task_id": task_id}
|
||||
if action_type == "knowledge_write":
|
||||
topic = str(item.get("topic") or "").strip()
|
||||
if not topic:
|
||||
skip(item, action_type, "missing_topic")
|
||||
continue
|
||||
action["topic"] = topic
|
||||
if item.get("scope") is not None:
|
||||
|
|
@ -554,6 +593,8 @@ def generate_reflection(
|
|||
reasoning_effort=resolve_effort("task")) # the owner's Task / Chat level: one SSOT, no literal
|
||||
raw_reflection_text = raw_reflection_text.strip()
|
||||
memory_operation_errors = refl_usage.get("_consolidation_errors") or []
|
||||
from ouroboros.knowledge import observed_route_stamp
|
||||
reflection_route = observed_route_stamp(refl_usage)
|
||||
if not raw_reflection_text and memory_operation_errors:
|
||||
raw_reflection_text = "(reflection generation failed: " + str(memory_operation_errors[-1].get("message") or "unknown") + ")"
|
||||
task_id_str = str(task.get("id", "") or "")
|
||||
|
|
@ -596,7 +637,11 @@ def generate_reflection(
|
|||
"priority": _truncate_with_notice(raw.get("priority", "med"), 10).strip().lower() or "med",
|
||||
"kind": _truncate_with_notice(raw.get("kind", "improvement"), 40).strip() or "improvement",
|
||||
})
|
||||
memory_actions = _validate_memory_actions(raw_memory_actions, task_id_str)
|
||||
# A rejected raw action is a typed event HERE, where production drops it
|
||||
# (apply_memory_actions never sees it); the retained task input is its source.
|
||||
memory_actions = _validate_memory_actions(raw_memory_actions, task_id_str, on_skip=functools.partial(
|
||||
record_memory_action_skip, pathlib.Path(knowledge_context.drive_root) / "logs" / "events.jsonl",
|
||||
project_id=str(getattr(knowledge_context, "project_id", "") or ""), input_ref=source_ref))
|
||||
memory_actions = [bound for action in memory_actions for bound in (
|
||||
knowledge.bind_entries([action]) if action["type"] == "knowledge_write" else [action])]
|
||||
|
||||
|
|
@ -614,6 +659,7 @@ def generate_reflection(
|
|||
reflection_text = f"(reflection generation failed: {e})"
|
||||
backlog_candidates = []
|
||||
memory_actions = []
|
||||
reflection_route = "unknown"
|
||||
|
||||
return {
|
||||
"ts": utc_now_iso(),
|
||||
|
|
@ -644,6 +690,10 @@ def generate_reflection(
|
|||
"reflection": reflection_text,
|
||||
"backlog_candidates": backlog_candidates,
|
||||
"memory_actions": memory_actions,
|
||||
# The route that ANSWERED the Light call (provider/resolved model, account
|
||||
# when served by Claudexor), never the configured route: a model-wait
|
||||
# override rebinds the call, and the history stamp must name what wrote.
|
||||
"route": reflection_route,
|
||||
**({"source_ref": source_ref} if source_ref else {}),
|
||||
**({"memory_operation_errors": memory_operation_errors} if memory_operation_errors else {}),
|
||||
}
|
||||
|
|
@ -665,12 +715,29 @@ def apply_memory_actions(env: Any, actions: List[Dict[str, Any]], *, project_id:
|
|||
"""
|
||||
pid = str(project_id or "").strip()
|
||||
applied = 0
|
||||
events = pathlib.Path(env.drive_root) / "logs" / "events.jsonl"
|
||||
# Project reflections live in a protected project store. Generic read_file
|
||||
# cannot open it, while the canonical log contains only a bounded pointer.
|
||||
# append_reflection_routed attaches an exact task-source copy to each action.
|
||||
fallback_ref = ({"status": "source_unavailable", "project_id": pid} if pid else
|
||||
{"read": {"tool": "read_file", "arguments": {
|
||||
"root": "runtime_data", "path": f"logs/{REFLECTIONS_FILENAME}"}}})
|
||||
|
||||
def retained_input(action: Dict[str, Any]) -> Dict[str, Any]:
|
||||
ref = action.get("_reflection_source_ref")
|
||||
return ref if isinstance(ref, dict) and ref.get("kind") == "task_source" else fallback_ref
|
||||
|
||||
def skipped(action: Dict[str, Any], reason: str) -> None:
|
||||
record_memory_action_skip(events, action, reason, project_id=pid, input_ref=retained_input(action))
|
||||
|
||||
for action in (actions or [])[:3]:
|
||||
atype = str(action.get("type") or "")
|
||||
content = str(action.get("content") or "").strip()
|
||||
if not content:
|
||||
skipped(action, "empty_content")
|
||||
continue
|
||||
if pid and atype in ("scratchpad_append", "identity_update_candidate"):
|
||||
skipped(action, "project_scoped_task")
|
||||
continue
|
||||
try:
|
||||
if atype == "scratchpad_append":
|
||||
|
|
@ -685,6 +752,7 @@ def apply_memory_actions(env: Any, actions: List[Dict[str, Any]], *, project_id:
|
|||
elif atype == "knowledge_write":
|
||||
topic = str(action.get("topic") or "").strip()
|
||||
if not topic:
|
||||
skipped(action, "missing_topic")
|
||||
continue
|
||||
from ouroboros.consolidator import _write_knowledge_entries
|
||||
from ouroboros.tools.registry import ToolContext
|
||||
|
|
@ -694,7 +762,10 @@ def apply_memory_actions(env: Any, actions: List[Dict[str, Any]], *, project_id:
|
|||
ctx = ToolContext(repo_dir=getattr(env, "repo_dir", env.drive_root), drive_root=root,
|
||||
budget_drive_root=canonical,
|
||||
project_id=pid, task_id=str(action.get("task_id") or ""))
|
||||
outcomes = _write_knowledge_entries(root / "memory" / "knowledge", [action], context=ctx)
|
||||
outcomes = _write_knowledge_entries(
|
||||
root / "memory" / "knowledge", [action], context=ctx,
|
||||
stamp={"writer": "reflection", "route": action.get("_reflection_route") or "unknown",
|
||||
"writer_input_ref": {**retained_input(action), "task_id": ctx.task_id}})
|
||||
applied += sum(row["ok"] for row in outcomes)
|
||||
if any(not row["ok"] for row in outcomes):
|
||||
log.warning("Reflection knowledge update was not published: %s", outcomes)
|
||||
|
|
@ -778,6 +849,7 @@ def append_reflection_routed(env: Any, task: Dict[str, Any], entry: Dict[str, An
|
|||
pid = ""
|
||||
if not pid:
|
||||
append_reflection(canonical, entry)
|
||||
_bind_reflection_action_source(canonical, entry)
|
||||
return
|
||||
from ouroboros.project_facts import project_reflections_path
|
||||
|
||||
|
|
@ -811,6 +883,31 @@ def append_reflection_routed(env: Any, task: Dict[str, Any], entry: Dict[str, An
|
|||
})
|
||||
except Exception:
|
||||
log.warning("Failed to write canonical reflection pointer", exc_info=True)
|
||||
_bind_reflection_action_source(canonical, entry)
|
||||
|
||||
|
||||
def _bind_reflection_action_source(canonical: pathlib.Path, entry: Dict[str, Any]) -> None:
|
||||
"""Give the later action writer an exact actor-readable source and the route that nominated it."""
|
||||
actions = entry.get("memory_actions") or []
|
||||
if not actions:
|
||||
return
|
||||
for action in actions:
|
||||
if isinstance(action, dict):
|
||||
action["_reflection_route"] = entry.get("route") or "unknown"
|
||||
try:
|
||||
from types import SimpleNamespace
|
||||
|
||||
from ouroboros.consolidator import retain_memory_source
|
||||
|
||||
ref = retain_memory_source(
|
||||
SimpleNamespace(drive_root=canonical, task_id=str(entry.get("task_id") or "reflection")),
|
||||
"reflection_memory_actions", json.dumps(entry, ensure_ascii=False).encode("utf-8"), "json")
|
||||
for action in actions:
|
||||
if isinstance(action, dict):
|
||||
action["_reflection_source_ref"] = ref
|
||||
except Exception:
|
||||
log.warning("Reflection action source retention failed for task %s", entry.get("task_id"), exc_info=True)
|
||||
|
||||
|
||||
_PATTERNS_PROMPT = """\
|
||||
You maintain a Pattern Register for Ouroboros, a self-modifying AI agent.
|
||||
|
|
|
|||
|
|
@ -259,7 +259,13 @@ def summarize_source(
|
|||
corrected, raw = _extract_trailing_json(corrected, "KNOWLEDGE_ENTRIES_JSON:")
|
||||
kept = [e for e in (raw if isinstance(raw, list) else []) if isinstance(e, dict)
|
||||
and (str(e.get("topic") or ""), str(e.get("scope") or "")) in draft_topics]
|
||||
entries.extend(knowledge.bind_entries(kept) if knowledge is not None and kept else [])
|
||||
# Provenance is per nomination: the route that answered THIS part's
|
||||
# correction wrote these entries, whatever route the block's other
|
||||
# rooms or parts ran on (a wait may rebind between parts).
|
||||
from ouroboros.knowledge import observed_route_stamp
|
||||
route = observed_route_stamp(usage)
|
||||
entries.extend({**entry, "_nomination_route": route}
|
||||
for entry in (knowledge.bind_entries(kept) if knowledge is not None and kept else []))
|
||||
summaries.append(corrected.strip())
|
||||
continue
|
||||
failure = usage["_consolidation_errors"][-1]
|
||||
|
|
|
|||
|
|
@ -114,7 +114,6 @@ BAND_PATHS = {
|
|||
"ouroboros/claudexor_daemon.py": "Installation daemon lifecycle owns marker and authenticated endpoint stop authority, confirmed self-started handles, and duplicate-start refusal; process signal and ledger mechanics remain in process_custody. No new lifecycle store or scheduler.",
|
||||
"ouroboros/claudexor_runtime.py": "Exact byte verification and delivery now have a shared owner for engine and skill resources; this module retains engine pin, archive installation and platform-specific contracts.",
|
||||
"ouroboros/cli.py": "The existing command-line transport keeps task-event negotiation, bounded replay deduplication and result rendering together; the additive cursor does not introduce a second CLI or task engine.",
|
||||
"ouroboros/consolidator.py": "Shrunk from the 1501-1600 giant band after per-room consolidation moved room draft/correction into room_consolidation.py; the block/era orchestration, chunk atomicity and route-fit machinery still share this owner.",
|
||||
"ouroboros/context.py": "Entered the band from the 1501-1600 zone (1590 lines) by the v7 D03 extraction of the runtime-section fact builders into ouroboros/context_runtime_facts.py; shrink-only residue of the split, not new growth.",
|
||||
"ouroboros/context_compaction.py": "Existing compaction owns propagation of typed model outcomes; unchanged semantic compaction policy.",
|
||||
"ouroboros/delegate_custody.py": "D07 DEL1 split brought the custody monolith DOWN from the 1600 hard cap into the band (1600->1305); reconcile family extracted to delegate_custody_reconcile.py, shrink-only direction",
|
||||
|
|
@ -141,6 +140,7 @@ BAND_PATHS = {
|
|||
"ouroboros/preflight_runner.py": None,
|
||||
"ouroboros/presence_runner.py": "Presence turn admission, durable retry identity and transport custody remain one owner; separating them now would duplicate the gate and receipt seam.",
|
||||
"ouroboros/projects_registry.py": "Entered the band from 999 lines: the stuck-Working liveness sprint homed the project-thread membership lens (mtime-cached) and its broadcast-choke marker here \u2014 registry semantics belong to the registry, not to message_bus.",
|
||||
"ouroboros/reflection.py": "TZ-3 PR-1: reflection now stamps typed skip events, writer/route provenance and the project-vs-canonical reflection locator on its memory actions (948->1017); one owner for reflection generation and its memory-action application, no new subsystem",
|
||||
"ouroboros/request_wire_receipts.py": "Wire candidates and semantic-success receipts share one exact serializer digest owner.",
|
||||
"ouroboros/request_wire_recovery.py": "E4 (#447): typed CustomToolProjectionError fallback keeps the wire-recovery ladder alive; includes the one-site-sufficient decision record at both retry catch sites",
|
||||
"ouroboros/review.py": "Entered the band from 952 lines: re-anchoring the size ratchet on the official line added the candidate and pairwise base-vs-tip transition validators (validate_size_ratchet_candidate/validate_size_ratchet_transition_against_base) with merge-aware previous-manifest resolution, replacing the retired first-parent history audit (update-flow-redesign sprint, Q7-C/Q18-A/Q19-A owner decisions).",
|
||||
|
|
|
|||
|
|
@ -164,9 +164,13 @@ def _knowledge_write(
|
|||
raise ValueError("The improvement-backlog requires parseable ### ibl-<id> blocks with - summary: lines; the global backlog was preserved")
|
||||
_record_backlog_history(backlog_path(root), sanitized, mode, str(getattr(ctx, "task_id", "") or ""))
|
||||
return f"✅ Knowledge '{sanitized}' merged into the global backlog ({merged} item(s))."
|
||||
# The turn is the writer; the route stamp is the route that ANSWERED the
|
||||
# loop's last round (provider + resolved model, account when Claudexor
|
||||
# served it), recorded by the loop, otherwise honestly unknown.
|
||||
result = knowledge_store.write_knowledge_note(
|
||||
_address(ctx, sanitized, scope), content, mode, expected_revision,
|
||||
str(getattr(ctx, "task_id", "") or ""), old_str,
|
||||
str(getattr(ctx, "task_id", "") or ""), old_str, writer="turn",
|
||||
route=(getattr(ctx, "_accumulated_usage", None) or {}).get("_observed_route") or None,
|
||||
)
|
||||
except ValueError as exc:
|
||||
return _publish_tool_result(ctx, ToolResult(
|
||||
|
|
|
|||
|
|
@ -1,7 +1,8 @@
|
|||
"""CPL4-C17 pins: knowledge journals append through the sidecar-lock seam.
|
||||
"""CPL4-C17 pins: knowledge history appends through the sidecar-lock seam.
|
||||
|
||||
``knowledge_history.jsonl`` / ``knowledge_journal.jsonl`` used raw
|
||||
``open("a")`` — the one torn-line hazard left among the memory journals.
|
||||
``knowledge_history.jsonl`` (and the since-removed ``knowledge_journal.jsonl``
|
||||
size telemetry) used raw ``open("a")`` — the one torn-line hazard left among
|
||||
the memory journals.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
|
|
|||
812
tests/test_memory_maintenance_visibility.py
Normal file
812
tests/test_memory_maintenance_visibility.py
Normal file
|
|
@ -0,0 +1,812 @@
|
|||
"""Nothing in memory maintenance leaves silently (TZ-3 PR-1, invariant I4).
|
||||
|
||||
Every host decision that used to vanish is a typed fact, each proven in both
|
||||
directions: a consolidation skipped on the lock, an era withheld because it was
|
||||
not shorter (with its ``era_retry`` record), a scratchpad pass and its outcome,
|
||||
a reflection lesson the host declined, and the ``writer``/``route``/
|
||||
``writer_input_ref``/``old_chars``/``new_chars`` stamp on every
|
||||
``source_capture`` history row. An era is built from summary blocks, never from
|
||||
an earlier era. The reader-less ``knowledge_journal.jsonl`` writer is gone.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import inspect
|
||||
import json
|
||||
import os
|
||||
from types import SimpleNamespace
|
||||
|
||||
import pytest
|
||||
|
||||
from ouroboros import consolidator as c
|
||||
from ouroboros import context_health
|
||||
from ouroboros import knowledge as store
|
||||
from ouroboros import reflection
|
||||
from ouroboros.memory import Memory
|
||||
from ouroboros.tools import knowledge as knowledge_tools
|
||||
from ouroboros.tools.registry import ToolContext
|
||||
from tests import test_consolidator_context_fit as fit_helpers
|
||||
from tests.test_consolidator_context_fit import _LLM, _paths, _write_chat
|
||||
|
||||
fit = fit_helpers.fit
|
||||
|
||||
|
||||
def _events(root, kind):
|
||||
path = root / "logs" / "events.jsonl"
|
||||
if not path.exists():
|
||||
return []
|
||||
rows = [json.loads(line) for line in path.read_text(encoding="utf-8").splitlines() if line.strip()]
|
||||
return [row for row in rows if row.get("type") == kind]
|
||||
|
||||
|
||||
def _history(root):
|
||||
path = root / "memory" / "knowledge_history.jsonl"
|
||||
return [json.loads(line) for line in path.read_text(encoding="utf-8").splitlines() if line.strip()]
|
||||
|
||||
|
||||
def _summary_blocks(count, start=0):
|
||||
return [{"ts": "2026-01-01T00:00:00Z", "type": "summary", "range": f"2026-01-01 {i:02d}:00 - {i:02d}:59",
|
||||
"message_count": 1, "content": f"block-{i} " + "x" * 40} for i in range(start, start + count)]
|
||||
|
||||
|
||||
def _era_block(label="old"):
|
||||
return {"ts": "2025-12-31T00:00:00Z", "type": "era", "range": "2025-12-01 to 2025-12-31",
|
||||
"message_count": 4, "content": f"### Era: {label}\n" + "e" * 30}
|
||||
|
||||
|
||||
# --- the consolidation lock ---------------------------------------------------------
|
||||
|
||||
|
||||
def test_a_lock_skip_is_a_typed_event_and_a_free_run_is_not(tmp_path, fit):
|
||||
chat, blocks, meta = _paths(tmp_path)
|
||||
_write_chat(chat, text_size=0)
|
||||
meta.parent.mkdir(parents=True, exist_ok=True)
|
||||
holder = os.open(str(meta.parent / ".consolidation.lock"), os.O_CREAT | os.O_WRONLY, 0o644)
|
||||
c._lock_nb(holder)
|
||||
try:
|
||||
ctx = ToolContext(repo_dir=tmp_path, drive_root=tmp_path, task_id="held")
|
||||
assert c.consolidate(chat, blocks, meta, _LLM(), knowledge_context=ctx) is None
|
||||
finally:
|
||||
c._unlock(holder)
|
||||
os.close(holder)
|
||||
skipped = _events(tmp_path, "consolidation_skipped_locked")
|
||||
assert len(skipped) == 1 and skipped[0]["task_id"] == "held"
|
||||
assert skipped[0]["lock_path"].endswith(".consolidation.lock")
|
||||
assert not blocks.exists() # the holder owned the run; nothing was consolidated twice
|
||||
|
||||
assert c.consolidate(chat, blocks, meta, _LLM())["_blocks_written"] == 1
|
||||
assert len(_events(tmp_path, "consolidation_skipped_locked")) == 1
|
||||
|
||||
|
||||
# --- an era is a compression of summary blocks, never of an era ---------------------
|
||||
|
||||
|
||||
def _seed_run(tmp_path, old_blocks, *, chat_count=c.BLOCK_SIZE):
|
||||
chat, blocks_path, meta_path = _paths(tmp_path)
|
||||
_write_chat(chat, count=chat_count, text_size=2)
|
||||
blocks_path.parent.mkdir(parents=True, exist_ok=True)
|
||||
blocks_path.write_text(json.dumps(old_blocks), encoding="utf-8")
|
||||
return chat, blocks_path, meta_path
|
||||
|
||||
|
||||
def _fake_era(monkeypatch, *, shorter):
|
||||
seen = []
|
||||
|
||||
def fake(run, *_args, **_kwargs):
|
||||
seen.append(list(run))
|
||||
source_len = sum(len(b["content"]) for b in run)
|
||||
content = "e" * (max(1, source_len // 4) if shorter else source_len + 10)
|
||||
return {"type": "era", "range": "era", "message_count": len(run), "content": content}, {
|
||||
"prompt_tokens": 1, "completion_tokens": 1, "total_tokens": 2, "cost": 0.01}
|
||||
|
||||
monkeypatch.setattr(c, "_compress_blocks_to_era", fake)
|
||||
return seen
|
||||
|
||||
|
||||
def test_an_earlier_era_bounds_the_run_and_is_never_recompressed(tmp_path, fit, monkeypatch):
|
||||
old = [_era_block(), *_summary_blocks(9)]
|
||||
chat, blocks_path, meta_path = _seed_run(tmp_path, old)
|
||||
seen = _fake_era(monkeypatch, shorter=True)
|
||||
assert c._run_block_consolidation(chat, blocks_path, meta_path, _LLM(), "", force_tail=True)
|
||||
run = old[1:1 + c.ERA_COMPRESS_COUNT]
|
||||
assert seen == [run] # the oldest run of summary blocks, the era ahead of it excluded
|
||||
stored = json.loads(blocks_path.read_text(encoding="utf-8"))
|
||||
assert stored[0] == old[0] # the old era keeps its place and bytes
|
||||
assert stored[1]["type"] == "era" and stored[1]["content"] != old[0]["content"]
|
||||
assert stored[2:2 + len(old) - 1 - c.ERA_COMPRESS_COUNT] == old[1 + c.ERA_COMPRESS_COUNT:] # the rest untouched
|
||||
|
||||
|
||||
def test_eras_ahead_of_the_run_never_hide_the_later_summaries(tmp_path, fit, monkeypatch):
|
||||
# Once the oldest four blocks were eras, a fixed four-block window found no
|
||||
# summary to compress and the later summaries were never compressed (review F1).
|
||||
old = [_era_block(str(i)) for i in range(c.ERA_COMPRESS_COUNT)] + _summary_blocks(6)
|
||||
chat, blocks_path, meta_path = _seed_run(tmp_path, old)
|
||||
seen = _fake_era(monkeypatch, shorter=True)
|
||||
assert c._run_block_consolidation(chat, blocks_path, meta_path, _LLM(), "", force_tail=True)
|
||||
run = old[c.ERA_COMPRESS_COUNT:2 * c.ERA_COMPRESS_COUNT]
|
||||
assert seen == [run]
|
||||
stored = json.loads(blocks_path.read_text(encoding="utf-8"))
|
||||
assert stored[:c.ERA_COMPRESS_COUNT] == old[:c.ERA_COMPRESS_COUNT] # the old eras keep their bytes
|
||||
assert stored[c.ERA_COMPRESS_COUNT]["type"] == "era" and stored[c.ERA_COMPRESS_COUNT] not in old
|
||||
assert stored[c.ERA_COMPRESS_COUNT + 1:-1] == old[2 * c.ERA_COMPRESS_COUNT:]
|
||||
assert len(stored) == len(old) + 1 - c.ERA_COMPRESS_COUNT + 1
|
||||
|
||||
|
||||
def test_a_history_of_eras_alone_makes_no_call_and_keeps_every_block(tmp_path, fit, monkeypatch):
|
||||
old = [_era_block(str(i)) for i in range(c.MAX_SUMMARY_BLOCKS)]
|
||||
chat, blocks_path, meta_path = _seed_run(tmp_path, old)
|
||||
seen = _fake_era(monkeypatch, shorter=True)
|
||||
assert c._run_block_consolidation(chat, blocks_path, meta_path, _LLM(), "", force_tail=True)
|
||||
assert seen == [] # the newest summary block is never its own era; nothing else is compressible
|
||||
stored = json.loads(blocks_path.read_text(encoding="utf-8"))
|
||||
assert stored[:len(old)] == old and len(stored) == len(old) + 1
|
||||
|
||||
|
||||
def test_the_chronicle_pass_treats_eras_and_gaps_as_boundaries(tmp_path, fit, monkeypatch):
|
||||
blocks_path = tmp_path / "memory" / "dialogue_blocks.json"
|
||||
blocks_path.parent.mkdir(parents=True)
|
||||
gap = {"gap_id": "g1", "type": "gap", "content": "[MEMORY GAP]"}
|
||||
summaries = _summary_blocks(3)
|
||||
blocks = [_era_block(), summaries[0], summaries[1], gap, summaries[2]]
|
||||
blocks_path.write_text(json.dumps(blocks), encoding="utf-8")
|
||||
seen = _fake_era(monkeypatch, shorter=True)
|
||||
c._compact_chronicle(blocks_path, _LLM(), "", None)
|
||||
assert seen == [summaries[:2], summaries[2:]]
|
||||
stored = json.loads(blocks_path.read_text(encoding="utf-8"))
|
||||
assert stored[0] == blocks[0] and stored[2] == gap
|
||||
assert [b["type"] for b in stored] == ["era", "era", "gap", "era"]
|
||||
|
||||
|
||||
def test_the_chronicle_pass_consults_and_records_the_same_era_retry(tmp_path, fit, monkeypatch):
|
||||
# Review F2: a throwaway meta let every pressure pass pay again for a run that was not
|
||||
# shorter, and a recorded refusal (no call, no usage) would have crashed the pass.
|
||||
blocks_path = tmp_path / "memory" / "dialogue_blocks.json"
|
||||
meta_path = tmp_path / "memory" / "dialogue_meta.json"
|
||||
blocks_path.parent.mkdir(parents=True)
|
||||
blocks = _summary_blocks(3)
|
||||
blocks_path.write_text(json.dumps(blocks), encoding="utf-8")
|
||||
c.atomic_write_json(meta_path, {"last_consolidated_offset": 7})
|
||||
seen = _fake_era(monkeypatch, shorter=False)
|
||||
|
||||
first = c._compact_chronicle(blocks_path, _LLM(), "", None, meta_path=meta_path)
|
||||
assert len(seen) == 1 and first["cost"] == 0.01
|
||||
meta = json.loads(meta_path.read_text(encoding="utf-8"))
|
||||
assert meta["last_consolidated_offset"] == 7 # the rest of meta survives the record
|
||||
(record,) = meta["era_retry"].values()
|
||||
assert record["route"] == {"model": "test/model", "use_local": False}
|
||||
assert json.loads(blocks_path.read_text(encoding="utf-8")) == blocks
|
||||
|
||||
second = c._compact_chronicle(blocks_path, _LLM(), "", None, meta_path=meta_path)
|
||||
assert len(seen) == 1 # the recorded refusal made no second paid call
|
||||
assert second["cost"] == 0 and second["_consolidation_errors"] == [] # no call was made
|
||||
assert [e["attempted"] for e in _events(tmp_path, "era_not_shorter")] == [True, False]
|
||||
assert json.loads(blocks_path.read_text(encoding="utf-8")) == blocks
|
||||
|
||||
# A pass without a meta path still pays the attempt; it reads and records no era_retry.
|
||||
assert c._compact_chronicle(blocks_path, _LLM(), "", None)["cost"] == 0.01
|
||||
assert len(seen) == 2 and json.loads(meta_path.read_text(encoding="utf-8")) == meta
|
||||
|
||||
|
||||
def test_each_run_keeps_its_own_refusal_and_a_success_erases_only_its_own(tmp_path, fit, monkeypatch):
|
||||
# Review N1: one refusal record for the whole chronicle paid again for run A on every
|
||||
# pass once run B's refusal overwrote it, and run B's success erased run A's record.
|
||||
blocks_path = tmp_path / "memory" / "dialogue_blocks.json"
|
||||
meta_path = tmp_path / "memory" / "dialogue_meta.json"
|
||||
blocks_path.parent.mkdir(parents=True)
|
||||
gap = {"gap_id": "g1", "type": "gap", "content": "[MEMORY GAP]"}
|
||||
run_a, run_b = _summary_blocks(2), [{**b, "range": "b-" + b["range"]} for b in _summary_blocks(2)]
|
||||
blocks_path.write_text(json.dumps([*run_a, gap, *run_b]), encoding="utf-8")
|
||||
c.atomic_write_json(meta_path, {})
|
||||
seen = _fake_era(monkeypatch, shorter=False)
|
||||
c._compact_chronicle(blocks_path, _LLM(), "", None, meta_path=meta_path)
|
||||
assert seen == [run_a, run_b]
|
||||
runs = c._era_retry_runs(json.loads(meta_path.read_text(encoding="utf-8")))
|
||||
assert len(runs) == 2 # both refusals remembered, keyed by source
|
||||
c._compact_chronicle(blocks_path, _LLM(), "", None, meta_path=meta_path)
|
||||
assert seen == [run_a, run_b] # neither run is paid for again
|
||||
|
||||
seen.clear()
|
||||
shorter = _fake_era(monkeypatch, shorter=True)
|
||||
monkeypatch.setattr(c, "_consolidation_route", lambda: ("other/model", False)) # a new route re-attempts
|
||||
c._compact_chronicle(blocks_path, _LLM(), "", None, meta_path=meta_path)
|
||||
assert shorter == [run_a, run_b]
|
||||
assert "era_retry" not in json.loads(meta_path.read_text(encoding="utf-8")) # each success cleared its own
|
||||
assert [b["type"] for b in json.loads(blocks_path.read_text(encoding="utf-8"))] == ["era", "gap", "era"]
|
||||
|
||||
|
||||
def test_era_retry_is_keyed_on_the_effective_binding_dispatch_uses(tmp_path, fit, monkeypatch):
|
||||
# Round 2 (critical 2): dispatch applies the Light account pin and the live
|
||||
# model-wait override; a refusal recorded under one effective binding must not
|
||||
# suppress the paid retry under another, while an unchanged binding still does.
|
||||
from contextlib import nullcontext
|
||||
|
||||
from ouroboros import model_slots, model_wait
|
||||
|
||||
blocks_path = tmp_path / "memory" / "dialogue_blocks.json"
|
||||
meta_path = tmp_path / "memory" / "dialogue_meta.json"
|
||||
blocks_path.parent.mkdir(parents=True)
|
||||
blocks_path.write_text(json.dumps(_summary_blocks(3)), encoding="utf-8")
|
||||
c.atomic_write_json(meta_path, {})
|
||||
seen = _fake_era(monkeypatch, shorter=False)
|
||||
|
||||
def record():
|
||||
(row,) = c._era_retry_runs(json.loads(meta_path.read_text(encoding="utf-8"))).values()
|
||||
return row["route"]
|
||||
|
||||
c._compact_chronicle(blocks_path, _LLM(), "", None, meta_path=meta_path)
|
||||
c._compact_chronicle(blocks_path, _LLM(), "", None, meta_path=meta_path)
|
||||
assert len(seen) == 1 and record() == {"model": "test/model", "use_local": False} # unchanged: suppressed
|
||||
|
||||
# The Light account pin changes; the configured lane does not.
|
||||
monkeypatch.setattr(model_slots, "model_role_option", lambda key, role, **_kw: "acct-B" if role == "light" else "")
|
||||
c._compact_chronicle(blocks_path, _LLM(), "", None, meta_path=meta_path)
|
||||
assert len(seen) == 2 and record() == {"model": "test/model", "use_local": False, "model_account_override": "acct-B"}
|
||||
c._compact_chronicle(blocks_path, _LLM(), "", None, meta_path=meta_path)
|
||||
assert len(seen) == 2 # the same pin: suppressed again
|
||||
|
||||
# A model-wait override rebinds the role; lane and pin are unchanged.
|
||||
override = {"model": "override/model", "use_local": False, "model_account_override": "acct-B"}
|
||||
waiter = SimpleNamespace(overrides={"light": override}, register_reprepare=lambda role, callback: nullcontext())
|
||||
monkeypatch.setattr(model_wait, "current_model_wait", lambda: waiter)
|
||||
c._compact_chronicle(blocks_path, _LLM(), "", None, meta_path=meta_path)
|
||||
assert len(seen) == 3 and record() == override
|
||||
c._compact_chronicle(blocks_path, _LLM(), "", None, meta_path=meta_path)
|
||||
assert len(seen) == 3
|
||||
assert [e["attempted"] for e in _events(tmp_path, "era_not_shorter")] == [True, False, True, False, True, False]
|
||||
|
||||
# One helper feeds both sites: the key IS the binding a real call dispatches on.
|
||||
assert c._light_dispatch_binding() == override
|
||||
llm = _LLM()
|
||||
assert c._call_consolidation_llm(llm, "prompt", "probe")[0] == "summary-1"
|
||||
sent = llm.calls[0]
|
||||
assert {key: sent[key] for key in override} == override
|
||||
|
||||
|
||||
def test_era_retry_is_keyed_to_the_binding_the_era_call_executed_on(tmp_path, fit, monkeypatch):
|
||||
# Round 3 (critical F2): an owner ``switch`` during a model wait INSIDE the era call
|
||||
# rebinds the role's override before the paid send, so a key captured before the
|
||||
# call named the binding that never answered: the one that did paid again on the
|
||||
# next pass, and a fresh attempt on the captured one was suppressed by its outcome.
|
||||
from contextlib import nullcontext
|
||||
|
||||
from ouroboros import model_wait
|
||||
|
||||
blocks_path = tmp_path / "memory" / "dialogue_blocks.json"
|
||||
meta_path = tmp_path / "memory" / "dialogue_meta.json"
|
||||
blocks_path.parent.mkdir(parents=True)
|
||||
blocks = _summary_blocks(3)
|
||||
blocks_path.write_text(json.dumps(blocks), encoding="utf-8")
|
||||
c.atomic_write_json(meta_path, {})
|
||||
route_a = {"model": "test/model", "use_local": False}
|
||||
route_b = {"model": "switched/model", "use_local": False, "model_account_override": "acct-B"}
|
||||
waiter = SimpleNamespace(overrides={}, register_reprepare=lambda role, callback: nullcontext())
|
||||
monkeypatch.setattr(model_wait, "current_model_wait", lambda: waiter)
|
||||
not_shorter = "e" * (sum(len(b["content"]) for b in blocks) + 10)
|
||||
|
||||
def switch_inside_the_call(llm, prompt):
|
||||
# The real era path: the first send dispatched on A; the owner switches the
|
||||
# role while it is in flight, and the rest of the unit dispatches on B.
|
||||
if len(llm.calls) == 1:
|
||||
assert c._light_route() == route_a
|
||||
waiter.overrides["light"] = dict(route_b)
|
||||
return {"content": not_shorter}, dict(llm.usage)
|
||||
|
||||
def record():
|
||||
(row,) = c._era_retry_runs(json.loads(meta_path.read_text(encoding="utf-8"))).values()
|
||||
return row["route"]
|
||||
|
||||
llm = _LLM(effect=switch_inside_the_call)
|
||||
c._compact_chronicle(blocks_path, llm, "", None, meta_path=meta_path)
|
||||
assert llm.calls[0]["model"] == "test/model" and llm.calls[-1]["model"] == "switched/model"
|
||||
assert record() == route_b # keyed to the executed binding, not the one captured before the call
|
||||
(event,) = _events(tmp_path, "era_not_shorter")
|
||||
assert event["attempted"] is True and event["route"] == route_b
|
||||
assert json.loads(blocks_path.read_text(encoding="utf-8")) == blocks
|
||||
|
||||
# The override still binds B: B's own refusal suppresses the paid retry on B.
|
||||
still_b = _LLM(effect=lambda llm, prompt: ({"content": not_shorter}, dict(llm.usage)))
|
||||
c._compact_chronicle(blocks_path, still_b, "", None, meta_path=meta_path)
|
||||
assert still_b.calls == [] and record() == route_b
|
||||
|
||||
# Back on A, which never answered for this run: A is paid for once, then keyed honestly.
|
||||
waiter.overrides.clear()
|
||||
on_a = _LLM(effect=lambda llm, prompt: ({"content": not_shorter}, dict(llm.usage)))
|
||||
c._compact_chronicle(blocks_path, on_a, "", None, meta_path=meta_path)
|
||||
assert on_a.calls and all(call["model"] == "test/model" for call in on_a.calls)
|
||||
assert record() == route_a
|
||||
assert [(e["attempted"], e["route"]) for e in _events(tmp_path, "era_not_shorter")] == [
|
||||
(True, route_b), (False, route_b), (True, route_a)]
|
||||
|
||||
|
||||
def test_a_legacy_single_era_retry_record_is_read_as_one_run(tmp_path):
|
||||
legacy = {"era_retry": {"source_sha256": "abc", "route": {"model": "m", "use_local": False}}}
|
||||
assert c._era_retry_runs(legacy) == {"abc": {"route": {"model": "m", "use_local": False}}}
|
||||
assert c._era_retry_runs({"era_retry": "garbage"}) == {} and c._era_retry_runs({}) == {}
|
||||
|
||||
|
||||
# --- a not-shorter era is recorded, visible, and not paid for twice -----------------
|
||||
|
||||
|
||||
def test_a_not_shorter_era_records_era_retry_and_the_event(tmp_path, fit, monkeypatch):
|
||||
old = _summary_blocks(c.MAX_SUMMARY_BLOCKS)
|
||||
chat, blocks_path, meta_path = _seed_run(tmp_path, old)
|
||||
seen = _fake_era(monkeypatch, shorter=False)
|
||||
assert c._run_block_consolidation(chat, blocks_path, meta_path, _LLM(), "", force_tail=True)
|
||||
assert len(seen) == 1
|
||||
stored = json.loads(blocks_path.read_text(encoding="utf-8"))
|
||||
assert stored[:c.MAX_SUMMARY_BLOCKS] == old and not any(b["type"] == "era" for b in stored)
|
||||
retry = json.loads(meta_path.read_text(encoding="utf-8"))["era_retry"]
|
||||
(source_sha256, record), = retry.items()
|
||||
assert record == {"route": {"model": "test/model", "use_local": False},
|
||||
"observed_route": store.UNKNOWN_STAMP} # the fake era usage names no physical route
|
||||
events = _events(tmp_path, "era_not_shorter")
|
||||
assert len(events) == 1 and events[0]["attempted"] is True and events[0]["observed_route"] == store.UNKNOWN_STAMP
|
||||
assert events[0]["source_sha256"] == source_sha256 and events[0]["blocks"] == c.ERA_COMPRESS_COUNT
|
||||
assert events[0]["era_chars"] > events[0]["source_chars"]
|
||||
|
||||
|
||||
def test_the_same_run_on_the_same_route_is_not_paid_again_until_either_changes(tmp_path, fit, monkeypatch):
|
||||
old = _summary_blocks(c.MAX_SUMMARY_BLOCKS)
|
||||
chat, blocks_path, meta_path = _seed_run(tmp_path, old)
|
||||
seen = _fake_era(monkeypatch, shorter=False)
|
||||
assert c._run_block_consolidation(chat, blocks_path, meta_path, _LLM(), "", force_tail=True)
|
||||
_write_chat(chat, count=2 * c.BLOCK_SIZE, text_size=2)
|
||||
assert c._run_block_consolidation(chat, blocks_path, meta_path, _LLM(), "", force_tail=True)
|
||||
assert len(seen) == 1 # the second run made no era call
|
||||
events = _events(tmp_path, "era_not_shorter")
|
||||
assert [event["attempted"] for event in events] == [True, False]
|
||||
assert events[1]["source_sha256"] == events[0]["source_sha256"]
|
||||
|
||||
monkeypatch.setattr(c, "_consolidation_route", lambda: ("other/model", False))
|
||||
_write_chat(chat, count=3 * c.BLOCK_SIZE, text_size=2)
|
||||
assert c._run_block_consolidation(chat, blocks_path, meta_path, _LLM(), "", force_tail=True)
|
||||
assert len(seen) == 2 # a new route earns a new attempt
|
||||
(record,) = c._era_retry_runs(json.loads(meta_path.read_text(encoding="utf-8"))).values()
|
||||
assert record["route"] == {"model": "other/model", "use_local": False}
|
||||
|
||||
|
||||
def test_a_shorter_era_replaces_the_run_and_clears_era_retry(tmp_path, fit, monkeypatch):
|
||||
old = _summary_blocks(c.MAX_SUMMARY_BLOCKS)
|
||||
chat, blocks_path, meta_path = _seed_run(tmp_path, old)
|
||||
meta_path.write_text(json.dumps({"era_retry": {"source_sha256": "stale", "route": "unknown"}}), encoding="utf-8")
|
||||
seen = _fake_era(monkeypatch, shorter=True)
|
||||
assert c._run_block_consolidation(chat, blocks_path, meta_path, _LLM(), "", force_tail=True)
|
||||
assert len(seen) == 1
|
||||
stored = json.loads(blocks_path.read_text(encoding="utf-8"))
|
||||
assert stored[0]["type"] == "era" and stored[1:c.MAX_SUMMARY_BLOCKS - c.ERA_COMPRESS_COUNT + 1] == old[c.ERA_COMPRESS_COUNT:]
|
||||
# Another run's (legacy-shaped) refusal survives this run's success; this run never had one.
|
||||
assert c._era_retry_runs(json.loads(meta_path.read_text(encoding="utf-8"))) == {"stale": {"route": "unknown"}}
|
||||
assert _events(tmp_path, "era_not_shorter") == []
|
||||
|
||||
|
||||
def test_health_names_a_withheld_era_without_a_timestamp(tmp_path):
|
||||
(tmp_path / "memory").mkdir(parents=True, exist_ok=True)
|
||||
env = SimpleNamespace(drive_root=tmp_path, repo_dir=tmp_path,
|
||||
repo_path=lambda p: tmp_path / p, drive_path=lambda p: tmp_path / p)
|
||||
c.atomic_write_json(tmp_path / "memory" / "dialogue_meta.json", {"last_consolidated_offset": 100})
|
||||
assert not any("ERA COMPRESSION" in line for line in context_health._memory_health_lines(env))
|
||||
c.atomic_write_json(tmp_path / "memory" / "dialogue_meta.json", {
|
||||
"era_retry": {"abcdef0123456789": {"route": {"model": "light/model", "use_local": False}},
|
||||
"0123456789abcdef": {"route": "unknown"}}})
|
||||
rows = [line for line in context_health._memory_health_lines(env) if "ERA COMPRESSION WITHHELD" in line]
|
||||
assert len(rows) == 2 and "abcdef012345" in rows[0] and "light/model" in rows[0] and "2026-" not in rows[0]
|
||||
assert "0123456789ab" in rows[1]
|
||||
|
||||
|
||||
# --- every scratchpad pass names its outcome ------------------------------------------
|
||||
|
||||
|
||||
class _Scratch:
|
||||
def __init__(self, content):
|
||||
self.content = content
|
||||
|
||||
def chat(self, **_kwargs):
|
||||
return {"content": self.content}, {"prompt_tokens": 10, "completion_tokens": 5, "cost": 0.02,
|
||||
"provider": "openrouter", "resolved_model": "light/served"}
|
||||
|
||||
|
||||
def _scratchpad(tmp_path, count=4):
|
||||
memory = Memory(tmp_path)
|
||||
for index in range(count):
|
||||
memory.append_scratchpad_block(f"block-{index}-" + (chr(97 + index) * 8_000), source=f"source-{index}")
|
||||
return memory
|
||||
|
||||
|
||||
def test_a_scratchpad_replacement_reports_its_counts_and_source(tmp_path):
|
||||
memory = _scratchpad(tmp_path)
|
||||
c.consolidate_scratchpad(memory, tmp_path / "memory" / "knowledge",
|
||||
_Scratch(json.dumps({"knowledge_entries": [], "compressed_block": "compressed"})))
|
||||
events = _events(tmp_path, "scratchpad_consolidation")
|
||||
assert len(events) == 1
|
||||
event = events[0]
|
||||
assert event["outcome"] == "replaced" and event["pressure"] is False
|
||||
assert (event["blocks_before"], event["compressed_blocks"], event["blocks_after"]) == (4, 2, 3)
|
||||
assert event["chars_before"] > event["chars_after"] > 0
|
||||
assert event["source_entry_id"] == memory.load_scratchpad_blocks()[0]["metadata"]["source_ref"]["entry_id"]
|
||||
assert event["knowledge_writes"] == {"ok": 0, "failed": 0}
|
||||
assert event["last_error_kind"] is None and event["accounted_upper_bound_usd"] == 0.02
|
||||
|
||||
|
||||
@pytest.mark.parametrize("content, outcome", [
|
||||
("not the requested JSON", "failed"),
|
||||
(json.dumps({"knowledge_entries": [], "compressed_block": " "}), "empty_block"),
|
||||
])
|
||||
def test_a_refused_scratchpad_pass_names_why_and_keeps_the_blocks(tmp_path, content, outcome):
|
||||
memory = _scratchpad(tmp_path)
|
||||
before = memory.load_scratchpad_blocks()
|
||||
c.consolidate_scratchpad(memory, tmp_path / "memory" / "knowledge", _Scratch(content))
|
||||
assert memory.load_scratchpad_blocks() == before
|
||||
events = _events(tmp_path, "scratchpad_consolidation")
|
||||
assert len(events) == 1 and events[0]["outcome"] == outcome
|
||||
assert events[0]["blocks_after"] == 4 and events[0]["source_entry_id"] == ""
|
||||
assert events[0]["last_error_kind"] == ("scratchpad_consolidation_failed" if outcome == "failed" else None)
|
||||
|
||||
|
||||
def test_no_scratchpad_pass_means_no_event(tmp_path):
|
||||
memory = Memory(tmp_path)
|
||||
memory.append_scratchpad_block("small", source="task")
|
||||
assert c.consolidate_scratchpad(memory, tmp_path / "memory" / "knowledge", _Scratch("unused")) is None
|
||||
assert _events(tmp_path, "scratchpad_consolidation") == []
|
||||
|
||||
|
||||
# --- a declined reflection lesson is a fact --------------------------------------------
|
||||
|
||||
|
||||
def test_project_scoped_reflection_skips_are_typed_events(tmp_path):
|
||||
env = SimpleNamespace(drive_root=tmp_path, repo_dir=tmp_path)
|
||||
applied = reflection.apply_memory_actions(env, [
|
||||
{"type": "scratchpad_append", "content": "a lesson", "task_id": "t1"},
|
||||
{"type": "identity_update_candidate", "content": "a trait", "task_id": "t1"},
|
||||
{"type": "knowledge_write", "content": "topic-less", "task_id": "t1"},
|
||||
], project_id="proj_x")
|
||||
assert applied == 0
|
||||
events = _events(tmp_path, "reflection_memory_action_skipped")
|
||||
assert [(e["action_type"], e["reason"]) for e in events] == [
|
||||
("scratchpad_append", "project_scoped_task"), ("identity_update_candidate", "project_scoped_task"),
|
||||
("knowledge_write", "missing_topic")]
|
||||
assert all(e["project_id"] == "proj_x" and e["task_id"] == "t1" and e["content_chars"] > 0 for e in events)
|
||||
assert events[0]["input_ref"] == {"status": "source_unavailable", "project_id": "proj_x"}
|
||||
assert not (tmp_path / "memory" / "scratchpad_blocks.json").exists()
|
||||
|
||||
|
||||
def _reflect(tmp_path, llm, task_id="t-skip", project_id="proj_x"):
|
||||
return reflection.generate_reflection(
|
||||
{"id": task_id, "text": "the episode", "drive_root": str(tmp_path), "project_id": project_id},
|
||||
{}, "trace", llm, {"rounds": 1, "cost": 0.1})
|
||||
|
||||
|
||||
_REJECTED_RAW = [
|
||||
{"type": "knowledge_write", "topic": "people/alex", "content": " "},
|
||||
{"type": "knowledge_write", "content": "topic-less"},
|
||||
{"type": "delete_everything", "content": "nope"},
|
||||
{"type": "scratchpad_append", "content": "a kept lesson"},
|
||||
]
|
||||
|
||||
|
||||
def test_rejected_raw_reflection_actions_are_typed_events_on_the_real_path(tmp_path, fit):
|
||||
# Round 2 (critical 3): production drops empty and topic-less actions in the
|
||||
# validator, before apply_memory_actions ever runs; the event must fire there.
|
||||
from tests.test_knowledge_consolidation import MemoryLLM
|
||||
|
||||
entry = _reflect(tmp_path, MemoryLLM("Reflection.\nMEMORY_ACTIONS_JSON: " + json.dumps(_REJECTED_RAW)))
|
||||
assert entry["reflection"] == "Reflection." and [a["type"] for a in entry["memory_actions"]] == ["scratchpad_append"]
|
||||
events = _events(tmp_path, "reflection_memory_action_skipped")
|
||||
assert [(e["action_type"], e["reason"], e["content_chars"]) for e in events] == [
|
||||
("knowledge_write", "empty_content", 3), ("knowledge_write", "missing_topic", len("topic-less")),
|
||||
("delete_everything", "unsupported_type", len("nope"))]
|
||||
assert all((e["task_id"], e["project_id"]) == ("t-skip", "proj_x") for e in events)
|
||||
# The retained exact task input the model answered is what the validator seam
|
||||
# kept; the rejected reflection text itself is not retained.
|
||||
assert entry["source_ref"]["kind"] == "task_source"
|
||||
assert all(e["input_ref"] == entry["source_ref"] for e in events)
|
||||
|
||||
|
||||
def test_a_failed_validator_skip_event_cannot_abort_the_reflection(tmp_path, fit, monkeypatch, caplog):
|
||||
from tests.test_knowledge_consolidation import MemoryLLM
|
||||
|
||||
original = reflection.append_jsonl
|
||||
|
||||
def broken_event(path, row, **kwargs):
|
||||
if row.get("type") == "reflection_memory_action_skipped":
|
||||
raise OSError("event store unavailable")
|
||||
return original(path, row, **kwargs)
|
||||
|
||||
monkeypatch.setattr(reflection, "append_jsonl", broken_event)
|
||||
entry = _reflect(tmp_path, MemoryLLM("Reflection.\nMEMORY_ACTIONS_JSON: " + json.dumps(_REJECTED_RAW)))
|
||||
assert entry["reflection"] == "Reflection." # not "(reflection generation failed ...)"
|
||||
assert [a["content"] for a in entry["memory_actions"]] == ["a kept lesson"]
|
||||
assert caplog.text.count("Reflection memory skip event could not be written") == 3
|
||||
assert _events(tmp_path, "reflection_memory_action_skipped") == []
|
||||
|
||||
|
||||
def test_failed_skip_event_cannot_discard_later_reflection_lessons(tmp_path, monkeypatch, caplog):
|
||||
env = SimpleNamespace(drive_root=tmp_path, repo_dir=tmp_path)
|
||||
original = reflection.append_jsonl
|
||||
|
||||
def broken_event(path, row, **kwargs):
|
||||
if row.get("type") == "reflection_memory_action_skipped":
|
||||
raise OSError("event store unavailable")
|
||||
return original(path, row, **kwargs)
|
||||
|
||||
monkeypatch.setattr(reflection, "append_jsonl", broken_event)
|
||||
assert reflection.apply_memory_actions(env, [
|
||||
{"type": "knowledge_write", "content": "topic-less", "task_id": "t1"},
|
||||
{"type": "scratchpad_append", "content": "a real lesson", "task_id": "t1"},
|
||||
], project_id="") == 1
|
||||
assert "Reflection memory skip event could not be written" in caplog.text
|
||||
assert "a real lesson" in Memory(tmp_path).load_scratchpad()
|
||||
|
||||
|
||||
def test_applied_reflection_actions_emit_no_skip(tmp_path):
|
||||
env = SimpleNamespace(drive_root=tmp_path, repo_dir=tmp_path)
|
||||
assert reflection.apply_memory_actions(env, [
|
||||
{"type": "scratchpad_append", "content": "a lesson", "task_id": "t1"},
|
||||
{"type": "identity_update_candidate", "content": "a trait", "task_id": "t1"},
|
||||
{"type": "scratchpad_append", "content": " ", "task_id": "t1"},
|
||||
]) == 2
|
||||
events = _events(tmp_path, "reflection_memory_action_skipped")
|
||||
assert [(e["action_type"], e["reason"], e["project_id"]) for e in events] == [
|
||||
("scratchpad_append", "empty_content", "")]
|
||||
|
||||
|
||||
# --- the host stamp on every source_capture row --------------------------------------
|
||||
|
||||
|
||||
def _address(root, topic="people/alex"):
|
||||
return store.resolve_knowledge_address(root, topic, "global")
|
||||
|
||||
|
||||
def test_an_unnamed_writer_and_a_named_one_both_stamp_the_capture_row(tmp_path):
|
||||
target = _address(tmp_path)
|
||||
assert store.write_knowledge_note(target, "# Alex\n\nFirst.").ok
|
||||
legacy_shaped = _history(tmp_path)[-1]
|
||||
assert legacy_shaped["publication"] == "source_capture"
|
||||
assert (legacy_shaped["writer"], legacy_shaped["route"], legacy_shaped["writer_input_ref"]) == (
|
||||
store.UNKNOWN_STAMP, store.UNKNOWN_STAMP, store.UNKNOWN_STAMP)
|
||||
assert (legacy_shaped["old_chars"], legacy_shaped["new_chars"]) == (0, len(legacy_shaped["new_content"]))
|
||||
|
||||
current = store.read_knowledge_note(target)
|
||||
result = store.write_knowledge_note(target, "# Alex\n\nFirst. Second.", expected_revision=current.revision,
|
||||
writer="turn", route={"model": "m"}, writer_input_ref={"chat": 1})
|
||||
assert result.ok
|
||||
row = _history(tmp_path)[-1]
|
||||
assert (row["writer"], row["route"], row["writer_input_ref"]) == ("turn", {"model": "m"}, {"chat": 1})
|
||||
assert row["old_chars"] == len(current.text) and row["new_chars"] == len(result.current.text)
|
||||
assert row["delta"]["old_chars"] == row["old_chars"] and row["source_ref"] == result.current.source_ref()
|
||||
assert "writer" not in result.current.text # the stamp lives on the history row, never in the note
|
||||
|
||||
|
||||
def test_the_route_stamp_is_what_answered_never_the_configuration():
|
||||
# Review F3: a model-wait override or account rotation changes what ANSWERED;
|
||||
# the configured Light route cannot say which model wrote the note.
|
||||
assert store.observed_route_stamp({"cost": 0.01}) == store.UNKNOWN_STAMP # no physical fact at all
|
||||
assert store.observed_route_stamp(None) == store.UNKNOWN_STAMP
|
||||
assert store.observed_route_stamp({"provider": "openrouter", "resolved_model": "openai/gpt-x"}) == {
|
||||
"provider": "openrouter", "model": "openai/gpt-x"}
|
||||
# Round 2 (advisory 5): a partial physical fact leaves the missing field unknown,
|
||||
# never the caller's configured model; the local lane stamps its own resolved_model.
|
||||
assert store.observed_route_stamp({"provider": "openrouter"}) == {"provider": "openrouter", "model": "unknown"}
|
||||
assert store.observed_route_stamp({"resolved_model": "local-model"}) == {"provider": "unknown", "model": "local-model"}
|
||||
assert store.observed_route_stamp({"provider": "local", "resolved_model": "local-model"}) == {
|
||||
"provider": "local", "model": "local-model"}
|
||||
# The production-shaped served route names the account as credentialProfileId.
|
||||
served = {"provider": "claudexor", "resolved_model": "claude-fable", "claudexor": {"route": {
|
||||
"source": "claude", "model": "claude-fable", "credentialProfileId": "acct-B", "accountFingerprint": "fp-B"}}}
|
||||
assert store.observed_route_stamp(served) == {
|
||||
"provider": "claudexor", "model": "claude-fable", "source": "claude", "account": "acct-B",
|
||||
"account_fingerprint": "fp-B"}
|
||||
rotated = {**served, "claudexor": {"route": {**served["claudexor"]["route"], "credentialProfileId": "acct-C"}}}
|
||||
assert store.observed_route_stamp(rotated)["account"] == "acct-C" # rotation is visible in the stamp
|
||||
# A merged consolidation usage forwards the LAST physical route of the unit.
|
||||
merged = c._merge_consolidation_usage({"cost": 0.01, "provider": "openrouter", "resolved_model": "a"},
|
||||
{"cost": 0.01, "provider": "openrouter", "resolved_model": "b"})
|
||||
assert merged["_observed_route"] == {"provider": "openrouter", "model": "b"}
|
||||
assert store.observed_route_stamp(merged) == {"provider": "openrouter", "model": "b"}
|
||||
assert "_observed_route" not in c._merge_consolidation_usage({"cost": 0.01})
|
||||
|
||||
|
||||
def test_the_route_stamp_never_reads_a_configured_route_from_an_empty_usage():
|
||||
# Round 2 (advisory 5): the former ``model``/``use_local`` fallback turned an
|
||||
# empty usage into the configured route. Only physical facts are read now.
|
||||
import inspect
|
||||
|
||||
assert list(inspect.signature(store.observed_route_stamp).parameters) == ["usage"]
|
||||
assert store.observed_route_stamp({}) == store.UNKNOWN_STAMP
|
||||
assert store.observed_route_stamp({"cost": 0.0, "prompt_tokens": 0}) == store.UNKNOWN_STAMP
|
||||
assert store.observed_route_stamp({"_observed_route": "unknown"}) == store.UNKNOWN_STAMP # not a dict stamp
|
||||
|
||||
|
||||
def test_a_merge_never_lets_an_earlier_stamp_masquerade_as_the_final_call():
|
||||
# Round 2 (advisory 5): the final call of a unit answered without a physical
|
||||
# fact (a released send, an exception's usage); an earlier known stamp must not
|
||||
# be forwarded as if that call had produced it.
|
||||
stamped = {"cost": 0.01, "provider": "openrouter", "resolved_model": "a"}
|
||||
merged = c._merge_consolidation_usage(stamped, {"cost": 0.01})
|
||||
assert "_observed_route" not in merged and store.observed_route_stamp(merged) == store.UNKNOWN_STAMP
|
||||
# A forwarded merged stamp counts as the (known) last call when it IS the last usage ...
|
||||
again = c._merge_consolidation_usage({"cost": 0.01}, c._merge_consolidation_usage(stamped))
|
||||
assert again["_observed_route"] == {"provider": "openrouter", "model": "a"}
|
||||
# ... and an unknown-last merged unit stays unknown through a further merge.
|
||||
nested = c._merge_consolidation_usage(stamped, c._merge_consolidation_usage(stamped, {"cost": 0.01}))
|
||||
assert store.observed_route_stamp(nested) == store.UNKNOWN_STAMP
|
||||
|
||||
|
||||
def test_a_direct_turn_stamps_itself_and_its_observed_route_when_the_loop_recorded_one(tmp_path):
|
||||
ctx = ToolContext(repo_dir=tmp_path, drive_root=tmp_path, task_id="turn-1")
|
||||
assert "✅" in knowledge_tools._knowledge_write(ctx, "notes/a", "Plain observation.")
|
||||
first = _history(tmp_path)[-1]
|
||||
assert (first["writer"], first["route"], first["task_id"]) == ("turn", store.UNKNOWN_STAMP, "turn-1")
|
||||
|
||||
# The loop records what answered its last round on every lane, not only Claudexor.
|
||||
ctx._accumulated_usage = {"_observed_route": {"provider": "openrouter", "model": "openai/gpt-x"}}
|
||||
assert "✅" in knowledge_tools._knowledge_write(ctx, "notes/b", "Another observation.")
|
||||
second = _history(tmp_path)[-1]
|
||||
assert second["writer"] == "turn" and second["route"] == {"provider": "openrouter", "model": "openai/gpt-x"}
|
||||
|
||||
|
||||
class _ThreeRoomNominating:
|
||||
"""One block of three rooms; each room's correction answers on its own route.
|
||||
|
||||
Review N2 / round 2 (advisory 6): provenance is per nomination. Room A's
|
||||
correction has no physical route fact (unknown), rooms B and C answer on two
|
||||
different accounts, and every correction also tries to forge the host stamp.
|
||||
"""
|
||||
topics = {"A": "people/alex", "B": "people/bob", "C": "people/cara"}
|
||||
accounts = {"A": None, "B": "acct-1", "C": "acct-2"}
|
||||
|
||||
def __init__(self):
|
||||
self.corrections = []
|
||||
|
||||
@staticmethod
|
||||
def _served(account):
|
||||
return {"provider": "claudexor", "resolved_model": "light/served",
|
||||
"claudexor": {"route": {"source": "codex", "credentialProfileId": account}}}
|
||||
|
||||
def chat(self, **kwargs):
|
||||
prompt = kwargs["messages"][0]["content"]
|
||||
if prompt.startswith("Compare this draft memory"):
|
||||
room = next(name for name in "ABC" if f"## Draft memory\nEpisode {name}." in prompt)
|
||||
self.corrections.append(room)
|
||||
account = self.accounts[room]
|
||||
usage = {"cost": 0.01, **({} if account is None else self._served(account))}
|
||||
return {"content": f"Episode {room}, checked.\nKNOWLEDGE_ENTRIES_JSON: " + json.dumps([
|
||||
{"topic": self.topics[room], "content": f"Understanding {room}.", "_nomination_route": _FORGED}])}, usage
|
||||
room = "A" if "entry-0 " in prompt else "B" if "entry-34 " in prompt else "C"
|
||||
return {"content": f"Episode {room}.\nKNOWLEDGE_ENTRIES_JSON: " + json.dumps([
|
||||
{"topic": self.topics[room], "content": f"Understanding {room}."}])}, {"cost": 0.01, **self._served("acct-draft")}
|
||||
|
||||
|
||||
def test_dialogue_consolidation_stamps_each_nomination_with_its_own_correction_route(tmp_path, fit):
|
||||
chat, blocks, meta = _paths(tmp_path)
|
||||
chat.parent.mkdir(parents=True, exist_ok=True)
|
||||
rows = [{"ts": f"2026-01-01T{index // 60:02d}:{index % 60:02d}:00Z", "direction": "in",
|
||||
"text": f"entry-{index} ", "chat_id": 1 if index < 34 else 2 if index < 67 else 3} for index in range(100)]
|
||||
chat.write_text("\n".join(json.dumps(row) for row in rows) + "\n", encoding="utf-8")
|
||||
ctx = ToolContext(repo_dir=tmp_path, drive_root=tmp_path, task_id="consolidate")
|
||||
llm = _ThreeRoomNominating()
|
||||
usage = c.consolidate(chat, blocks, meta, llm, knowledge_context=ctx)
|
||||
assert llm.corrections == ["A", "B", "C"] and usage["_blocks_written"] == 1
|
||||
# The block-wide stamp is the LAST call's known route (room C's correction) ...
|
||||
assert c._route_stamp(usage)["account"] == "acct-2"
|
||||
captures = {row["topic"]: row for row in _history(tmp_path) if row.get("publication") == "source_capture"}
|
||||
assert set(captures) == {"people/alex", "people/bob", "people/cara"}
|
||||
assert all(row["writer"] == "consolidation" for row in captures.values())
|
||||
# ... yet each nomination carries the route of the correction that released IT:
|
||||
# an explicit unknown outranks the known block stamp, two known routes stay
|
||||
# distinct within one block, and the forged model stamp reached none of them.
|
||||
assert captures["people/alex"]["route"] == store.UNKNOWN_STAMP
|
||||
assert captures["people/bob"]["route"] == {
|
||||
"provider": "claudexor", "model": "light/served", "source": "codex", "account": "acct-1"}
|
||||
assert captures["people/cara"]["route"] == {
|
||||
"provider": "claudexor", "model": "light/served", "source": "codex", "account": "acct-2"}
|
||||
block = json.loads(blocks.read_text(encoding="utf-8"))[0]
|
||||
assert all(row["writer_input_ref"] == block["knowledge_source_ref"] for row in captures.values())
|
||||
assert block["knowledge_source_ref"]["entry_id"] and len(block["knowledge_writes"]) == 3
|
||||
assert "_nomination_route" not in block # a history stamp, never a persisted block field
|
||||
# The retained nominations row keeps the HOST stamp per entry, never the model's.
|
||||
nominations = _history(tmp_path)[0]
|
||||
assert nominations["type"] == "dialogue_knowledge_nominations"
|
||||
assert [(entry["topic"], entry["_nomination_route"]) for entry in nominations["nominations"][0]["entries"]] == [
|
||||
("people/alex", store.UNKNOWN_STAMP), ("people/bob", captures["people/bob"]["route"]),
|
||||
("people/cara", captures["people/cara"]["route"])]
|
||||
|
||||
|
||||
def test_scratchpad_consolidation_stamps_its_journal_source(tmp_path):
|
||||
memory = _scratchpad(tmp_path)
|
||||
c.consolidate_scratchpad(memory, tmp_path / "memory" / "knowledge", _Scratch(json.dumps({
|
||||
"knowledge_entries": [{"topic": "lessons/one", "content": "A durable lesson."}],
|
||||
"compressed_block": "compressed"})))
|
||||
capture = next(row for row in _history(tmp_path) if row.get("publication") == "source_capture")
|
||||
assert capture["writer"] == "scratchpad_consolidation"
|
||||
assert capture["route"] == {"provider": "openrouter", "model": "light/served"}
|
||||
assert capture["writer_input_ref"] == memory.load_scratchpad_blocks()[0]["metadata"]["source_ref"]
|
||||
assert _events(tmp_path, "scratchpad_consolidation")[0]["knowledge_writes"] == {"ok": 1, "failed": 0}
|
||||
|
||||
|
||||
_FORGED = {"provider": "forged", "model": "forged/model", "account": "forged-acct"}
|
||||
|
||||
|
||||
def test_a_model_supplied_nomination_route_never_reaches_history_from_scratchpad(tmp_path):
|
||||
# Round 2 (critical 1): the scratchpad producer binds every key the model wrote
|
||||
# before the writer ran; a forged host stamp must be dropped at that binding.
|
||||
memory = _scratchpad(tmp_path)
|
||||
c.consolidate_scratchpad(memory, tmp_path / "memory" / "knowledge", _Scratch(json.dumps({
|
||||
"knowledge_entries": [{"topic": "lessons/one", "content": "A durable lesson.", "_nomination_route": _FORGED,
|
||||
"_expected_revision": "forged", "_task_id": "forged"}],
|
||||
"compressed_block": "compressed"})))
|
||||
capture = next(row for row in _history(tmp_path) if row.get("publication") == "source_capture")
|
||||
assert capture["topic"] == "lessons/one" and capture["writer"] == "scratchpad_consolidation"
|
||||
assert capture["route"] == {"provider": "openrouter", "model": "light/served"} # the host's observed stamp
|
||||
assert _events(tmp_path, "scratchpad_consolidation")[0]["knowledge_writes"] == {"ok": 1, "failed": 0}
|
||||
journal = [json.loads(line) for line in memory.journal_path().read_text(encoding="utf-8").splitlines() if line.strip()]
|
||||
(bound,) = next(row for row in journal if row.get("type") == "blocks_consolidated")["knowledge_entries"]
|
||||
assert not [key for key in bound if key.startswith("_")] # nothing model-written survives as a host key
|
||||
|
||||
|
||||
def test_a_model_supplied_nomination_route_never_reaches_history_from_knowledge_maintenance(tmp_path, fit):
|
||||
from tests.test_memory_pressure_maintenance import setup_memory
|
||||
|
||||
memory, ctx = setup_memory(tmp_path)
|
||||
assert store.write_knowledge_note(store.resolve_knowledge_address(tmp_path, "overview", "global"),
|
||||
"---\nsummary: Orientation.\n---\n" + "Detailed understanding. " * 200).ok
|
||||
|
||||
class Forging:
|
||||
def chat(self, **_kwargs):
|
||||
return {"content": json.dumps({"knowledge_entries": [
|
||||
{"topic": "lessons/forged", "scope": "global", "content": "A shorter detail note.",
|
||||
"_nomination_route": _FORGED}]})}, {
|
||||
"cost": 0.02, "provider": "claudexor", "resolved_model": "light/served",
|
||||
"claudexor": {"route": {"source": "codex", "credentialProfileId": "acct-real"}}}
|
||||
|
||||
result = c.maintain_memory_pressure(memory, Forging(), ctx, fits=lambda: False)
|
||||
action = next(row for row in result["actions"] if row["owner"] == "knowledge_maintenance")
|
||||
assert [(row["topic"], row["scope"], row["ok"]) for row in action["writes"]] == [("lessons/forged", "global", True)]
|
||||
capture = next(row for row in _history(tmp_path) if row.get("publication") == "source_capture"
|
||||
and row["topic"] == "lessons/forged")
|
||||
assert capture["writer"] == "knowledge_maintenance"
|
||||
assert capture["route"] == {"provider": "claudexor", "model": "light/served", "source": "codex", "account": "acct-real"}
|
||||
|
||||
|
||||
def test_bind_entries_keeps_model_fields_and_drops_every_host_key(tmp_path):
|
||||
reads = c.KnowledgeReadContext(ToolContext(repo_dir=tmp_path, drive_root=tmp_path, task_id="op-1"))
|
||||
(bound,) = reads.bind_entries([{"topic": "people/alex", "content": "Observed.", "scope": "global",
|
||||
"_nomination_route": _FORGED, "_anything": 1, "extra": "kept"}])
|
||||
assert bound == {"topic": "people/alex", "content": "Observed.", "scope": "global", "extra": "kept",
|
||||
"expected_revision": None, "canonical_root": str(tmp_path), "task_id": "op-1"}
|
||||
|
||||
|
||||
def test_project_reflection_action_uses_actor_readable_exact_source(tmp_path):
|
||||
from ouroboros.artifacts import read_actor_source_bytes
|
||||
env = SimpleNamespace(drive_root=tmp_path, budget_drive_root=tmp_path, repo_dir=tmp_path)
|
||||
entry = {"task_id": "t3", "ts": "2026-01-01T00:00:00Z", "route": {"provider": "claudexor", "model": "claude-fable"},
|
||||
"memory_actions": [
|
||||
{"type": "knowledge_write", "topic": "lessons/project", "content": "Grounded.", "task_id": "t3"}]}
|
||||
reflection.append_reflection_routed(env, {"id": "t3", "project_id": "proj_x",
|
||||
"budget_drive_root": str(tmp_path)}, entry)
|
||||
action = entry["memory_actions"][0]
|
||||
source = action["_reflection_source_ref"]
|
||||
assert source["kind"] == "task_source"
|
||||
assert json.loads(read_actor_source_bytes(tmp_path, "t3", source))["memory_actions"][0]["content"] == "Grounded."
|
||||
assert reflection.apply_memory_actions(env, entry["memory_actions"], project_id="proj_x") == 1
|
||||
history = tmp_path / "projects" / "proj_x" / "knowledge_history.jsonl"
|
||||
row = json.loads(history.read_text(encoding="utf-8").splitlines()[-1])
|
||||
assert row["writer_input_ref"]["sha256"] == source["sha256"]
|
||||
assert row["writer_input_ref"]["task_id"] == "t3"
|
||||
assert row["route"] == {"provider": "claudexor", "model": "claude-fable"} # the reflection's answering route
|
||||
|
||||
|
||||
def test_reflection_stamps_the_reflection_row_it_came_from(tmp_path):
|
||||
env = SimpleNamespace(drive_root=tmp_path, repo_dir=tmp_path)
|
||||
assert reflection.apply_memory_actions(env, [
|
||||
{"type": "knowledge_write", "topic": "lessons/two", "content": "Reusable fact.", "task_id": "t9"}]) == 1
|
||||
capture = next(row for row in _history(tmp_path) if row.get("publication") == "source_capture")
|
||||
assert capture["writer"] == "reflection" and capture["route"] == store.UNKNOWN_STAMP
|
||||
assert capture["writer_input_ref"]["task_id"] == "t9"
|
||||
assert capture["writer_input_ref"]["read"]["arguments"]["path"] == "logs/task_reflections.jsonl"
|
||||
|
||||
|
||||
def test_the_reader_less_knowledge_journal_is_no_longer_written(tmp_path):
|
||||
assert store.write_knowledge_note(_address(tmp_path), "# Alex\n\nFirst.").ok
|
||||
assert not (tmp_path / "memory" / "knowledge_journal.jsonl").exists()
|
||||
assert (tmp_path / "memory" / "knowledge_history.jsonl").exists()
|
||||
assert "knowledge_journal" not in inspect.getsource(store)
|
||||
|
|
@ -524,7 +524,6 @@ def scan_data_paths(root: pathlib.Path = REPO) -> frozenset[str]:
|
|||
paths.update({
|
||||
"projects/*/knowledge/*.md",
|
||||
"projects/*/knowledge_history.jsonl",
|
||||
"projects/*/knowledge_journal.jsonl",
|
||||
})
|
||||
return frozenset(paths)
|
||||
|
||||
|
|
@ -579,7 +578,10 @@ def scan_data_paths(root: pathlib.Path = REPO) -> frozenset[str]:
|
|||
# owner PEM rotates every cache); one section-2 row covers both.
|
||||
# 296 -> 297 (Presence resilience): Presence recovery inspects the retained quarantine
|
||||
# members (``task_results/quarantine/*``) before deciding whether an event ever started.
|
||||
EXPECTED_SCAN_PATHS = 297
|
||||
# 297 -> 296 (TZ-3 PR-1): the ``knowledge_journal.jsonl`` size-telemetry writer is
|
||||
# removed (its only reader was this inventory); ``knowledge_history.jsonl`` keeps the
|
||||
# complete captures, now host-stamped.
|
||||
EXPECTED_SCAN_PATHS = 296
|
||||
|
||||
# Scanned paths that must always be present — guards the scanner itself
|
||||
# against a silent regression that would shrink coverage while keeping counts
|
||||
|
|
|
|||
|
|
@ -167,7 +167,11 @@ CHAPTER_BYTE_BUDGETS: dict[str, int] = {
|
|||
# residuals; a mechanism the chapter lacked, so only its two stale clauses were replaced.
|
||||
# 313400 -> 314000 (PR #1300; measured 313858 on the merged tree): the transport paragraph gains
|
||||
# the trust-bundle seam every first-party client shares; no older text to displace.
|
||||
"docs/architecture/06-agent-core.md": 314000,
|
||||
# 314000 -> 314900 (TZ-3 PR-1, measured 314850 on the merged tree): the era run boundary with
|
||||
# its `era_retry` record keyed to the executed Light binding, the four typed memory-maintenance
|
||||
# events and the host stamp on `source_capture` history rows are mechanisms no older text
|
||||
# described; the sentences they extend were rewritten in place, not appended to.
|
||||
"docs/architecture/06-agent-core.md": 314900,
|
||||
# 36991 -> 37300: the facade paragraph names the three loop constants runtime_limits.py
|
||||
# gained (events batch bound, budget-projection retry interval); no older text to displace.
|
||||
# 37300 -> 38400 (PR #1207): the Z.ai (`zai::`) direct provider gets its own route
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue