wip(tz3-pr1): writer-neutral PR-1 candidate (reviewed packet 40dc2585, index tree 5a00e3cd)

This commit is contained in:
Ouroboros 2026-09-25 23:13:43 +03:00
parent fa741a7fe9
commit d7a6cf8b31
11 changed files with 565 additions and 135 deletions

View file

@ -143,10 +143,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 + Light route of a not-shorter era) | blocks reduced by era compression (over 10 blocks, the oldest run of up to 4 SUMMARY blocks; an era is never re-compressed until the calendar chronicle, 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
@ -167,7 +167,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 |

View file

@ -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 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; every skipped action (and an empty or topic-less one) is a `reflection_memory_action_skipped` event naming its reason and the correct project or canonical reflection source (the canonical Project row is only a pointer). Reflection may propose a future campaign or backlog item, but it cannot 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.
@ -559,11 +559,11 @@ Every IMPLICIT claim — the UI conversion, that admission, the reaper's retry a
`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.
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 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 over one run of summary blocks bounded by gaps and earlier eras (never an era of an era, until the calendar chronicle replaces eras); a failed or not-shorter era retains the blocks, the latter as `era_retry` (source hash + Light route: no paid repeat while both hold, shown in Health) and an `era_not_shorter` event; a lock skip is a `consolidation_skipped_locked` event; each scratchpad pass ends in a `scratchpad_consolidation` event naming its outcome; 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`.
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.
`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, 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, whose `source_capture` rows carry the host stamp `writer`/`route`/`writer_input_ref` and `old_chars`/`new_chars` (rows older than the stamp read `unknown`); the reader-less `knowledge_journal.jsonl` size telemetry is no longer written. 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.
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.

View file

@ -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,22 @@ def _consolidation_route() -> Tuple[str, bool]:
return resolve_credentialed_model(lane.model), False
def _light_route() -> Any:
"""The configured Light route as a history/meta stamp; unknown when it cannot be resolved."""
try:
return dict(zip(("model", "use_local"), _consolidation_route()))
except Exception:
return "unknown"
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 +131,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,18 +145,14 @@ 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)
@ -281,9 +281,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 +340,14 @@ 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"]} if usage.get("_knowledge_entries") else {})})
processed += len(chunk)
# Set after the last merge of this stretch: _merge_consolidation_usage forwards
@ -372,8 +368,7 @@ def _run_block_consolidation(
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"}},
@ -398,30 +393,22 @@ def _run_block_consolidation(
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.
run_start = next((i for i, b in enumerate(old_blocks) 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(old_blocks) and not _is_run_boundary(old_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(old_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 = [*old_blocks[:run_start], era, *old_blocks[run_end:], *all_blocks[compress_count:]]
_write_locked_json(blocks_path, all_blocks)
@ -430,7 +417,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": _light_route(),
"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", {
@ -949,26 +938,55 @@ 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")
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
(source hash + Light route, after ``consolidation_retry``): the same source on the
same route is not paid for again, 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)}
retry = meta.get("era_retry") or {}
if retry.get("source_sha256") == fact["source_sha256"] and retry.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"]:
meta["era_retry"] = {"source_sha256": fact["source_sha256"], "route": fact["route"]}
_emit_event(logs_dir, "era_not_shorter", attempted=True, era_chars=len(era["content"]), **fact)
return None, usage
if era is not None:
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."""
"""Reduce every contiguous run of summary blocks, preserving gaps, earlier eras and exact sources."""
blocks = _load_blocks(blocks_path)
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)
# A measured-pressure pass always attempts: its throwaway meta consults no retry record.
era, usage = _era_for_run(run, {}, blocks_path.parent.parent / "logs", 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)
reduced.extend([era] if era is not None else run)
if usage.get("_consolidation_errors"):
reduced.extend(blocks[end:])
break
@ -1054,7 +1072,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": _light_route(), "writer_input_ref": source_ref})
except (ValueError, TypeError, AttributeError) as exc:
action["error"] = str(exc)
actions.append(action)
@ -1121,17 +1140,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 +1309,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 +1321,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 +1364,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 +1376,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 +1386,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": _light_route(), "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 +1425,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 +1436,29 @@ 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``."""
from ouroboros.knowledge import KnowledgeAddress, sanitize_topic, write_knowledge_note
from ouroboros.tools.knowledge import _address, _record_backlog_history
@ -1476,7 +1485,7 @@ def _write_knowledge_entries(
"reason": "backlog_merge" if merged >= 0 else "unparseable_backlog"})
continue
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 ""), **(stamp or {}))
outcomes.append({"topic": topic, "scope": address.scope, "ok": result.ok,
"reason": result.reason,
"source_ref": result.current.source_ref() if result.current else None})

View file

@ -273,6 +273,15 @@ 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')}"
)
retry = meta.get("era_retry")
if isinstance(retry, dict):
route = retry.get("route")
lines.append(
f"WARNING: DIALOGUE ERA COMPRESSION WITHHELD — the era for source run "
f"{str(retry.get('source_sha256') or '')[: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

View file

@ -23,6 +23,7 @@ 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"
@ -328,8 +329,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 +390,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 +410,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)

View file

@ -665,12 +665,34 @@ 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"
# The reflection row that nominated an action is what its writer saw.
from ouroboros.project_facts import sanitize_project_id
# Project reflections live in their own durable project log; the canonical
# log has only a bounded pointer and cannot witness the nominated action.
reflection_path = (f"projects/{sanitize_project_id(pid)}/logs/{REFLECTIONS_FILENAME}"
if pid else f"logs/{REFLECTIONS_FILENAME}")
reflection_ref = {"read": {"tool": "read_file", "arguments": {
"root": "runtime_data", "path": reflection_path}}}
def skipped(action: Dict[str, Any], reason: str) -> None:
# A lesson the host declines is a fact, not silence (I4): name it where
# Health and the owner can count it, with the reflection row to reread.
append_jsonl(events, {"ts": utc_now_iso(), "type": "reflection_memory_action_skipped",
"task_id": str(action.get("task_id") or ""), "project_id": pid,
"action_type": str(action.get("type") or ""), "reason": reason,
"content_chars": len(str(action.get("content") or "")),
"reflection_ref": reflection_ref})
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 +707,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 +717,9 @@ 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", "writer_input_ref": {**reflection_ref, "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)

View file

@ -164,9 +164,12 @@ 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; its observed model route is the loop's fact when
# the lane reported one, otherwise the stamp stays 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("_model_route") or None,
)
except ValueError as exc:
return _publish_tool_result(ctx, ToolResult(

View file

@ -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

View file

@ -0,0 +1,373 @@
"""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_consolidation_honesty import _Nominating
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)
assert seen == [old[1:c.ERA_COMPRESS_COUNT]] # the window's summary blocks, the era 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) - c.ERA_COMPRESS_COUNT] == old[c.ERA_COMPRESS_COUNT:] # the rest untouched
def test_a_window_of_eras_only_makes_no_call_and_keeps_every_block(tmp_path, fit, monkeypatch):
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)
assert seen == []
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"]
# --- 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"]
assert retry == {"source_sha256": retry["source_sha256"], "route": {"model": "test/model", "use_local": False}}
events = _events(tmp_path, "era_not_shorter")
assert len(events) == 1 and events[0]["attempted"] is True
assert events[0]["source_sha256"] == retry["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
assert json.loads(meta_path.read_text(encoding="utf-8"))["era_retry"]["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:]
assert "era_retry" not in json.loads(meta_path.read_text(encoding="utf-8"))
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": {"source_sha256": "abcdef0123456789", "route": {"model": "light/model", "use_local": False}}})
row = next(line for line in context_health._memory_health_lines(env) if "ERA COMPRESSION WITHHELD" in line)
assert "abcdef012345" in row and "light/model" in row and "2026-" not in row
# --- 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}
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]["reflection_ref"]["read"]["arguments"]["path"] == (
"projects/proj_x/logs/task_reflections.jsonl")
assert not (tmp_path / "memory" / "scratchpad_blocks.json").exists()
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_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")
ctx._accumulated_usage = {"_model_route": {"source": "claude", "model": "fable"}}
assert "✅" in knowledge_tools._knowledge_write(ctx, "notes/b", "Another observation.")
second = _history(tmp_path)[-1]
assert second["writer"] == "turn" and second["route"] == {"source": "claude", "model": "fable"}
def test_dialogue_consolidation_stamps_its_seam_route_and_source(tmp_path, fit):
chat, blocks, meta = _paths(tmp_path)
_write_chat(chat, count=100, text_size=0)
ctx = ToolContext(repo_dir=tmp_path, drive_root=tmp_path, task_id="consolidate")
c.consolidate(chat, blocks, meta, _Nominating(), knowledge_context=ctx)
capture = next(row for row in _history(tmp_path) if row.get("publication") == "source_capture")
assert capture["writer"] == "consolidation"
assert capture["route"] == {"model": "test/model", "use_local": False}
block = json.loads(blocks.read_text(encoding="utf-8"))[0]
assert capture["writer_input_ref"] == block["knowledge_source_ref"]
assert capture["writer_input_ref"]["entry_id"]
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["writer_input_ref"] == memory.load_scratchpad_blocks()[0]["metadata"]["source_ref"]
assert _events(tmp_path, "scratchpad_consolidation")[0]["knowledge_writes"] == {"ok": 1, "failed": 0}
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)

View file

@ -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)
@ -574,7 +573,10 @@ def scan_data_paths(root: pathlib.Path = REPO) -> frozenset[str]:
# 295 -> 294: TZ-3 removed the destructive memory journal rewrite and its
# ``.compact.tmp`` sibling path; PERSISTENCE.md keeps the journals, now
# read-only observed and never age-digested.
EXPECTED_SCAN_PATHS = 294
# 294 -> 293: TZ-3 PR-1 removed the ``knowledge_journal.jsonl`` size-telemetry
# writer (its only reader was this inventory); ``knowledge_history.jsonl`` keeps
# the complete captures, now host-stamped.
EXPECTED_SCAN_PATHS = 293
# Scanned paths that must always be present — guards the scanner itself
# against a silent regression that would shrink coverage while keeping counts

View file

@ -159,7 +159,11 @@ CHAPTER_BYTE_BUDGETS: dict[str, int] = {
# 309800 -> 310100 (TZ2 + #1262 merge, measured 309985): the Presence task-message
# own-binding boundary and forced declaration remain beside #1262's name-miss
# contract; both are independent rules in the same chapter, not duplicate prose.
"docs/architecture/06-agent-core.md": 310100,
# 310100 -> 310800 (TZ-3 PR-1, measured 310721): the era run boundary with its
# `era_retry` record, 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": 310800,
# 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.
"docs/architecture/07-configuration.md": 37300,