mirror of
https://github.com/razzant/ouroboros.git
synced 2026-10-02 19:58:46 +00:00
feat: preserve memory history and expose pending nominations
This commit is contained in:
parent
86d4e62986
commit
b2b5b4f5a3
16 changed files with 374 additions and 443 deletions
|
|
@ -22,13 +22,13 @@ The manifest is the SSOT of the module→domain assignment (1:1, complete over t
|
|||
| D12 | Settings & configuration | 15 | 0 |
|
||||
| D13 | Safety, guards & runtime mode | 9 | 0 |
|
||||
| D14 | Skills & extensions | 56 | 0 |
|
||||
| D15 | Memory, knowledge, consciousness & self-evolution | 22 | 0 |
|
||||
| D15 | Memory, knowledge, consciousness & self-evolution | 23 | 0 |
|
||||
| D16 | Observability, usage accounting & cost | 11 | 0 |
|
||||
| D17 | Projects, workspaces & task results | 23 | 0 |
|
||||
| D18 | Launcher, packaging, platform & shared substrate | 15 | 0 |
|
||||
| D19 | Frozen contracts (ABI) | 10 | 0 |
|
||||
| D20 | Presence | 10 | 0 |
|
||||
| **total** | | **574** | **0** |
|
||||
| **total** | | **575** | **0** |
|
||||
|
||||
## Dependency direction matrix (strict, pinned)
|
||||
|
||||
|
|
@ -720,6 +720,7 @@ No function body (≥ 10 normalized lines) is shared verbatim across domains. Ne
|
|||
- `ouroboros/knowledge.py`
|
||||
- `ouroboros/memory.py`
|
||||
- `ouroboros/memory_journal_compaction.py`
|
||||
- `ouroboros/memory_nomination_receipts.py`
|
||||
- `ouroboros/post_task_evolution.py`
|
||||
- `ouroboros/project_facts.py`
|
||||
- `ouroboros/reflection.py`
|
||||
|
|
|
|||
|
|
@ -22,7 +22,7 @@ scanned data-relative path to be covered by a row here (count-anchored both ways
|
|||
keys migrate). Governs subagent worktrees, headless/task drives, task trees,
|
||||
service logs, consumed schedule receipts,
|
||||
confirmed capability probes, delegate recovery/supervision sweeps, code_intel
|
||||
and reconcile-failed prunes, memory-journal digesting and agent media.
|
||||
reconcile-failed prunes, and agent media. Memory journals retain full new rows independently of this knob.
|
||||
- **Rotation** — `supervisor/state.py::rotate_jsonl_log_if_needed`: >800 KB →
|
||||
atomic rename to `archive/<prefix>_<ts>.jsonl` under the append lock.
|
||||
Applied on the supervisor tick to `chat.jsonl`, `progress.jsonl`,
|
||||
|
|
@ -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` (locked atomic) | none | bounded by era compression (10 blocks, oldest 4 compressed) | blocks: compressed biography irreproducible; meta: full re-consolidation (cost, not loss) |
|
||||
| `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_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; digested rows carry `content_digested: true` | full old+new text only inside GC retention: older identity/knowledge/patterns rows go digest-only (sha256+len) at startup (`memory_journal_compaction.py`, under the append lock, unreadable lines byte-preserved); scratchpad journal keeps its own eviction contract | undo/provenance record lost (live .md survives); eviction/rewrite paths fail closed when journal append fails; digested history is irreversible by design |
|
||||
| `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/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
|
||||
|
|
@ -180,13 +180,15 @@ scanned data-relative path to be covered by a row here (count-anchored both ways
|
|||
Always safe (pure caches, recreated): `state/pycache`, `state/code_intel`,
|
||||
`state/evolution_metrics_cache.json`, `playwright-browsers/`, `state/cx`,
|
||||
`state/betterleaks`, lock files, `state/server_port`.
|
||||
Safe with bounded cost: `WORLD.md` (regenerates), `dialogue_meta.json`
|
||||
(re-consolidation), `state/usage_import_watermark.json` (safe re-import),
|
||||
Safe with bounded cost: `WORLD.md` (regenerates),
|
||||
`state/usage_import_watermark.json` (safe re-import),
|
||||
`ui_preferences.json`, `auth_secret.key` (one re-login).
|
||||
Fail-closed losses (system stays correct, work/authority is forgone):
|
||||
skill state dirs, `advisory_review.json`, `capability_evidence.json`,
|
||||
`pending_restart_verify.json`.
|
||||
Dangerous (authority/history destruction): `settings.json`,
|
||||
`state/usage_attempts.jsonl`, `task_results/**`, `logs/events.jsonl`,
|
||||
`memory/**`, `archive/**`, `observability/**`, `state/subagent_worktrees.json`
|
||||
`memory/**` (including `dialogue_meta.json`: deletion erases cursor and pending
|
||||
nomination obligations; re-consolidation cannot reconstruct the old IDs),
|
||||
`archive/**`, `observability/**`, `state/subagent_worktrees.json`
|
||||
(leak), `claudexor/**`, `state/python-userbase` (real deps).
|
||||
|
|
|
|||
|
|
@ -186,11 +186,12 @@ server.py (Starlette+uvicorn) ← HTTP + WebSocket on configurable host:port (de
|
|||
├── consciousness_wake.py ← The wake-up MESSAGE (`prompts/CONSCIOUSNESS.md` rendered as the turn's USER message, cuts disclosed as `(+N more)`) and the origin/authority envelope `wake_task_metadata`
|
||||
├── consciousness_authority.py ← The three autonomy levels of a wake (observe/act/full) and their consequences — `disabled_tools`, bound at dispatch only so the prompt prefix matches an owner turn's, `runtime_mode_cap=light` below Full, and Observe's argument-level narrowing of the mutating names it keeps (§6 Background consciousness and Evolution)
|
||||
├── consciousness_allowance.py ← Rolling-24h consciousness spend read off the usage ledger; typed `allowance_unknown` on a read failure; read by the alarm and the single admission door in `supervisor/queue.py`
|
||||
├── room_consolidation.py ← Per-room memory: one Light draft + one source-grounded correction per room, deterministic assembly of typed room sections into one block/era; no cross-room LLM recombine, legacy blocks keep unknown provenance (§6)
|
||||
├── consolidator.py ← Dialogue consolidation with a generation-aware cursor; an unfindable generation appends a loud `[MEMORY GAP]` block, never a silent offset reset; `last_consolidation_error` / `last_unpublished_nominations` in `dialogue_meta.json` (§6 Durable memory and project focus)
|
||||
├── room_consolidation.py ← Per-room Light draft/correction and deterministic assembly; no cross-room LLM recombine (§6)
|
||||
├── consolidator.py ← Generation cursor, explicit `[MEMORY GAP]`, and knowledge nomination outcomes in `dialogue_meta.json` (§6)
|
||||
├── memory_nomination_receipts.py ← Source-addressed pending nominations; no cross-batch retirement (§6)
|
||||
├── memory.py ← Scratchpad, identity, chat history
|
||||
├── knowledge.py ← `ouroboros/knowledge.py`: linked-Markdown note addressing, exact source reads, generated shelf indexes for global and project knowledge, and revision-checked writes, so concurrent cognition cannot silently overwrite a newer note (§6 Durable memory and project focus)
|
||||
├── memory_journal_compaction.py ← Digest-only compaction of old memory-journal snapshots: the digest replaces the snapshots it summarizes, never a silent drop
|
||||
├── memory_journal_compaction.py ← Startup read-only size facts (`memory_journal_observation`); new history stays complete, old digests unrecoverable
|
||||
├── project_facts.py ← project_id resolution (explicit `--project-id` or workspace-path hash); per-project knowledge dir `projects/<id>/knowledge` isolated from `memory/knowledge`; journal/workpad helpers
|
||||
├── task_tree_ledger.py ← Append-only `data/task_trees/<root>/blackboard.jsonl`: EPHEMERAL typed swarm coordination (`tree_note`/`tree_read`), mirrored into the durable project journal at root completion; pruned on root terminal
|
||||
├── projects_registry.py ← Durable `data/state/projects.json`: 80-char names, `active|deleting|tombstoned`; deletion preserves bindings/history/folder/memory; a tombstone blocks resurrection; reconcile NEVER prunes (§6 Project registry and lease)
|
||||
|
|
|
|||
|
|
@ -557,9 +557,9 @@ Every IMPLICIT claim — the UI conversion, that admission, the reaper's retry a
|
|||
|
||||
`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.
|
||||
|
||||
`consolidator.py` publishes a dialogue block only after every part succeeds, retaining raw generations and their cursor; an unfindable generation appends `[MEMORY GAP]`, never silently resets the offset. `context_fit` measures full Light requests against fresh route/account capacity and calibrated density (`llm_local.local_context_limits` owns local output reservation; missing/stale evidence remains unknown). `room_consolidation.py` drafts and corrects each room separately, then deterministically assembles the sections. Episodic text is grounded in that room's source; cumulative knowledge replacements require the complete current note plus the episode in BOTH stages. Corrected entries bind to the corrector's complete delivered read, never the draft's revision credit. A narrow or older episode cannot negate prior facts or later receipts; supported corrections and removals remain model judgment. Range reads, authored views, CAS and old/new history remain the publication path; no new stage or store. Failed, empty or truncated correction withholds the chunk, cursor and nominations.
|
||||
`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 source splits without clipping, including inside an entry; continuation context stays outside source bytes. A real refusal records its source hash and strictly smaller same-route byte bound in `dialogue_meta.json` (`consolidation_retry`); changed source, route, capacity or output reserve invalidates it. Era compression regroups each recorded room across blocks and reassembles deterministically; legacy untyped blocks remain explicitly unknown provenance. Failed/overflowed eras preserve old blocks. Failures retain `last_consolidation_error`, cleared by an advance without a new failure; incomplete knowledge publication retains `last_unpublished_nominations`. Unknown spend remains nullable and control/resource/unknown model errors retain `propagate_model_error` semantics.
|
||||
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`. An unreadable existing meta refuses consolidation and warns in Health, never replacing obligations with `{}`. Old digests and debts do not auto-resolve by same-topic writes. 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.
|
||||
|
||||
|
|
|
|||
|
|
@ -2,7 +2,7 @@
|
|||
|
||||
Machine extraction of the `docs/ARCHITECTURE.md` "Data layout (`~/Ouroboros/`)" tree — the durable-file orientation carrier (this tree's counterpart of the reference PERSISTENCE_OWNERS derivation checklist) — regenerated by `python scripts/regenerate_inventories.py`. Do not edit. Every entry is probed against reality: repo entries must exist as tracked paths; data-plane entries must appear as a literal in the runtime sources that construct them. A durable file renamed or removed in code while its tree row survives = red (`tests/test_generated_inventories.py`).
|
||||
|
||||
Source: `docs/architecture/01-high-level-architecture.md`, physical LF lines 602-691; UTF-8 SHA-256 `1839b923a939496f18c4d428807d8c876ca14422dd09c33001bf39a2b4a0e0ba`.
|
||||
Source: `docs/architecture/01-high-level-architecture.md`, physical LF lines 603-692; UTF-8 SHA-256 `b0e84be83d74aea31c3b247c57386cce150e7b5ceeb066c1a182e23d026d0473`.
|
||||
|
||||
- entries: **79** (code-ref: 72, repo-dir: 6, repo-path: 1)
|
||||
|
||||
|
|
|
|||
|
|
@ -10,7 +10,6 @@ from ouroboros import room_consolidation
|
|||
from ouroboros.utils import (
|
||||
append_jsonl,
|
||||
atomic_write_json,
|
||||
read_json_dict,
|
||||
replace_atomic,
|
||||
utc_now_iso,
|
||||
read_text,
|
||||
|
|
@ -388,6 +387,10 @@ def _run_block_consolidation(
|
|||
return total_usage
|
||||
for nominated_block, _entries in pending_knowledge:
|
||||
nominated_block["knowledge_source_ref"] = ref
|
||||
if knowledge_context is not None:
|
||||
from ouroboros.memory_nomination_receipts import prepare
|
||||
pending_ids = prepare(meta, source_id, pending_knowledge)
|
||||
atomic_write_json(meta_path, meta) # Debt precedes block and note publication.
|
||||
|
||||
existing_blocks = _load_blocks(blocks_path)
|
||||
all_blocks = existing_blocks + new_blocks
|
||||
|
|
@ -435,21 +438,11 @@ def _run_block_consolidation(
|
|||
"source_ref": block["knowledge_source_ref"], "outcomes": block["knowledge_writes"],
|
||||
})
|
||||
if pending_knowledge:
|
||||
# Nominations were durable before mutation. Outcome facts belong to
|
||||
# the same blocks, so a failed write is available to later learning.
|
||||
_write_locked_json(blocks_path, all_blocks)
|
||||
# Era compression later replaces these blocks with one object carrying no
|
||||
# knowledge_writes, so this batch receipt lives in meta, not in a scan of
|
||||
# dialogue_blocks.json. A fully published batch clears it; a run with no
|
||||
# nominations at all leaves the older receipt standing.
|
||||
# Count what was NOMINATED, not only what produced an outcome: an entry
|
||||
# the writer skipped as malformed was not published either.
|
||||
nominated = sum(len(entries) for _block, entries in pending_knowledge)
|
||||
failed = nominated - sum(1 for outcome in published if outcome["ok"])
|
||||
meta.pop("last_unpublished_nominations", None)
|
||||
if failed > 0:
|
||||
meta["last_unpublished_nominations"] = {"entry_id": ref["entry_id"],
|
||||
"failed": failed, "total": nominated}
|
||||
from ouroboros.memory_nomination_receipts import settle
|
||||
settle(meta, pending_ids, published)
|
||||
# Legacy batch-only receipts remain open: no positional evidence can
|
||||
# prove which old entry a later successful nomination resolved.
|
||||
|
||||
_advance_cursor(meta, segments, segment_sigs, segment_entries, last_offset + processed)
|
||||
if not run_failed: # An advance by a run that recorded no failure retires a stale error.
|
||||
|
|
@ -1237,7 +1230,9 @@ def _advance_cursor(
|
|||
|
||||
|
||||
def _load_meta(path: pathlib.Path) -> Dict[str, Any]:
|
||||
return read_json_dict(path) or {}
|
||||
from ouroboros.memory_nomination_receipts import load_meta
|
||||
|
||||
return load_meta(path)
|
||||
|
||||
|
||||
from ouroboros.utils import jsonl_generation_signature as _chat_log_signature
|
||||
|
|
@ -1452,9 +1447,11 @@ def _write_knowledge_entries(
|
|||
outcomes = []
|
||||
for entry in entries:
|
||||
if not isinstance(entry, dict):
|
||||
outcomes.append({"topic": "", "ok": False, "reason": "malformed_nomination"})
|
||||
continue
|
||||
topic, content = entry.get("topic"), entry.get("content")
|
||||
if not isinstance(content, str) or not content.strip():
|
||||
outcomes.append({"topic": topic, "ok": False, "reason": "empty_nomination"})
|
||||
continue
|
||||
try:
|
||||
topic = sanitize_topic(topic)
|
||||
|
|
|
|||
|
|
@ -233,14 +233,28 @@ def _memory_health_lines(env: Any) -> List[str]:
|
|||
pass
|
||||
|
||||
try:
|
||||
meta = read_json_dict(env.drive_path("memory/dialogue_meta.json")) or {}
|
||||
from ouroboros.memory_nomination_receipts import load_meta
|
||||
|
||||
meta = load_meta(env.drive_path("memory/dialogue_meta.json"))
|
||||
pending = meta.get("pending_knowledge_nominations")
|
||||
if pending:
|
||||
# load_meta already validates the whole list; malformed state raises.
|
||||
sample = ", ".join(row["id"].split(":")[0][:12] + ":" +
|
||||
":".join(row["id"].split(":")[-2:])
|
||||
for row in pending[:3])
|
||||
lines.append(
|
||||
f"WARNING: DIALOGUE KNOWLEDGE PUBLICATION OPEN — {len(pending)} source-addressed "
|
||||
f"nominations (first {min(3, len(pending))}: {sample}; omitted {max(0, len(pending)-3)}). "
|
||||
"Read memory/dialogue_meta.json and memory/knowledge_history.jsonl for full source. "
|
||||
"No automatic or tool-level discharge exists yet; later successes cannot retire older entries."
|
||||
)
|
||||
receipt = meta.get("last_unpublished_nominations")
|
||||
if isinstance(receipt, dict) and int(receipt.get("failed") or 0) > 0:
|
||||
# The recovery route is named because the reader may hold no read_file:
|
||||
# an external-channel turn has the cognitive memory tools and nothing else.
|
||||
lines.append(
|
||||
f"WARNING: LAST DIALOGUE KNOWLEDGE PUBLICATION INCOMPLETE — {receipt.get('failed')} of "
|
||||
f"{receipt.get('total')} nominations from the latest consolidation batch were not published "
|
||||
f"{receipt.get('total')} nominations in a legacy consolidation batch remain unresolved "
|
||||
f"(entry_id {receipt.get('entry_id')}); from the main chat, read_file(root='runtime_data', "
|
||||
"path='memory/knowledge_history.jsonl') and publish what still holds"
|
||||
)
|
||||
|
|
@ -251,7 +265,10 @@ def _memory_health_lines(env: Any) -> List[str]:
|
|||
f"at cursor {error.get('cursor_offset')}"
|
||||
)
|
||||
except Exception:
|
||||
pass
|
||||
# A broken existing meta file is NOT an empty nomination/cursor state.
|
||||
# Consolidation refuses to overwrite it; its reader must name the gap.
|
||||
lines.append("WARNING: DIALOGUE META UNREADABLE — memory/dialogue_meta.json; "
|
||||
"consolidation withheld to preserve existing bytes")
|
||||
return lines
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -273,6 +273,7 @@ D20 = "Presence"
|
|||
"ouroboros/mcp_client.py" = "D05"
|
||||
"ouroboros/memory.py" = "D15"
|
||||
"ouroboros/memory_journal_compaction.py" = "D15"
|
||||
"ouroboros/memory_nomination_receipts.py" = "D15"
|
||||
"ouroboros/model_concurrency.py" = "D02"
|
||||
"ouroboros/model_send_seal.py" = "D16"
|
||||
"ouroboros/mutation_attribution.py" = "D01"
|
||||
|
|
|
|||
|
|
@ -1,182 +1,20 @@
|
|||
"""Digest-only compaction of old memory-journal snapshots (CPL4-C16, owner 4A).
|
||||
"""Compatibility entry point for the retired destructive memory-journal sweep.
|
||||
|
||||
``memory/identity_journal.jsonl``, ``memory/knowledge_history.jsonl`` and
|
||||
``memory/knowledge/patterns_history.jsonl`` record the FULL old+new document
|
||||
text on every write — O(doc×edits) growth, the worst byte offenders in the
|
||||
memory plane. Owner decision 4A: entries younger than the unified GC
|
||||
retention keep their full text; older entries become digest-only — the
|
||||
content keys are replaced by their sha256 + length (existing hashes are
|
||||
never overwritten) and the row is marked ``content_digested``.
|
||||
|
||||
Strictly fail-closed per line: an unparseable line, a row without a
|
||||
readable ``ts``, a row with nothing to digest, or a row whose STORED digest
|
||||
disagrees with the text it claims to describe is carried through
|
||||
BYTE-IDENTICAL. The scratchpad journal (typed rows, its own eviction
|
||||
contract) is deliberately NOT in scope.
|
||||
|
||||
This is the only sweep that DESTROYS content rather than whole dead files,
|
||||
so its three guards are load-bearing (audit #15-11):
|
||||
|
||||
* **Digest truth before deletion.** The digest becomes the only surviving
|
||||
record of the text, so a stored ``*_sha256``/``*_len`` that does not match
|
||||
the text is never published over it: the row keeps its full content and the
|
||||
mismatch is reported as a typed fact (``digest_mismatch``, surfaced on the
|
||||
``memory_journal_compaction`` event).
|
||||
* **A lock nobody can steal.** The append lock is taken ``owner_aware_stale``
|
||||
so elapsed time alone can never hand a second writer the same journal.
|
||||
* **Publish only an unchanged source.** ``append_jsonl`` falls back to an
|
||||
UNLOCKED append after its own lock timeout, so a concurrent row can still
|
||||
land while this rewrite streams. The file is re-identified (size + inode)
|
||||
against the bytes actually consumed immediately before ``os.replace``; any
|
||||
delta aborts the publish and the journal stays as the appender left it.
|
||||
|
||||
The rewrite streams line by line into the temp sibling — the journals are the
|
||||
worst byte offenders in the memory plane and must never be loaded whole.
|
||||
The startup/maintenance caller still invokes this function. Retaining that call
|
||||
keeps old integrations working while new knowledge, identity and Pattern Register
|
||||
history stays complete. A previously digested row cannot be reconstructed;
|
||||
no new row loses its old/new text merely because of its age.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import hashlib
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import pathlib
|
||||
from typing import Any, Dict, Optional, Tuple
|
||||
|
||||
from ouroboros.deadline_utils import parse_deadline_ts
|
||||
from ouroboros.platform_layer import acquire_exclusive_file_lock, release_exclusive_file_lock
|
||||
from ouroboros.utils import jsonl_append_lock_path
|
||||
|
||||
log = logging.getLogger(__name__)
|
||||
|
||||
_JOURNAL_RELS = (
|
||||
pathlib.Path("memory") / "identity_journal.jsonl",
|
||||
pathlib.Path("memory") / "knowledge_history.jsonl",
|
||||
pathlib.Path("memory") / "knowledge" / "patterns_history.jsonl",
|
||||
)
|
||||
_CONTENT_KEYS = ("old_content", "new_content")
|
||||
from pathlib import Path
|
||||
from typing import Any, Dict, Optional
|
||||
import stat
|
||||
|
||||
|
||||
def _digest_row(row: Dict[str, Any]) -> str:
|
||||
"""Replace full-text keys with sha256+len.
|
||||
|
||||
Returns ``"digested"`` when text was dropped, ``"mismatch"`` when a STORED
|
||||
digest or length contradicts the text it describes, ``""`` when there was
|
||||
nothing to digest. On a mismatch the row is left EXACTLY as found: the
|
||||
digest is about to become the only surviving record of that content, and
|
||||
publishing a digest already known to be false while deleting the last
|
||||
correct copy is unrecoverable.
|
||||
"""
|
||||
dropped = []
|
||||
for key in _CONTENT_KEYS:
|
||||
value = row.get(key)
|
||||
if not isinstance(value, str):
|
||||
continue
|
||||
prefix = key[: -len("_content")]
|
||||
digest = hashlib.sha256(value.encode("utf-8")).hexdigest() if value else ""
|
||||
stored_digest = row.get(f"{prefix}_sha256")
|
||||
stored_len = row.get(f"{prefix}_len")
|
||||
if isinstance(stored_digest, str) and stored_digest != digest:
|
||||
return "mismatch"
|
||||
if isinstance(stored_len, int) and not isinstance(stored_len, bool) and stored_len != len(value):
|
||||
return "mismatch"
|
||||
dropped.append((key, prefix, digest, len(value)))
|
||||
if not dropped:
|
||||
return ""
|
||||
for key, prefix, digest, length in dropped:
|
||||
row[f"{prefix}_sha256"] = digest
|
||||
row[f"{prefix}_len"] = length
|
||||
del row[key]
|
||||
row["content_digested"] = True
|
||||
return "digested"
|
||||
|
||||
|
||||
def _digest_line(raw: bytes, cutoff: float) -> Tuple[bytes, str]:
|
||||
"""One journal line, transformed or carried through byte-identical."""
|
||||
stripped = raw.strip()
|
||||
if not stripped:
|
||||
return raw, ""
|
||||
try:
|
||||
row = json.loads(stripped.decode("utf-8"))
|
||||
except (UnicodeDecodeError, ValueError):
|
||||
return raw, "" # fail-closed: never rewrite what cannot be read
|
||||
if not isinstance(row, dict):
|
||||
return raw, ""
|
||||
parsed_ts = parse_deadline_ts(str(row.get("ts") or ""))
|
||||
if parsed_ts is None or parsed_ts.timestamp() >= cutoff:
|
||||
return raw, "" # fresh, or age unknowable: keep full text
|
||||
outcome = _digest_row(row)
|
||||
if outcome != "digested":
|
||||
return raw, outcome
|
||||
return json.dumps(row, ensure_ascii=False).encode("utf-8") + b"\n", "digested"
|
||||
|
||||
|
||||
def _publish_if_unchanged(
|
||||
path: pathlib.Path, tmp: pathlib.Path, expected: Tuple[int, int, int],
|
||||
) -> bool:
|
||||
"""Swap the rewritten journal in ONLY if the source is still what we read.
|
||||
|
||||
``append_jsonl`` appends WITHOUT the sidecar lock once its own acquisition
|
||||
times out, so a concurrent row can land while this rewrite streams. The
|
||||
identity is (bytes consumed, device, inode): a grown file means an append
|
||||
we did not carry over, a different inode means the journal was replaced
|
||||
outright. Either way the rewrite is dropped and the appender's file stands.
|
||||
"""
|
||||
try:
|
||||
stat = path.stat()
|
||||
except OSError:
|
||||
return False
|
||||
if (int(stat.st_size), int(stat.st_dev), int(stat.st_ino)) != expected:
|
||||
return False
|
||||
os.replace(tmp, path)
|
||||
return True
|
||||
|
||||
|
||||
def _compact_one(path: pathlib.Path, cutoff: float) -> Tuple[int, int, str]:
|
||||
"""Digest one journal in place.
|
||||
|
||||
Returns ``(digested, mismatched, error)``; a nonempty ``error``
|
||||
(``lock_unavailable`` / ``source_changed``) means nothing was published and
|
||||
the journal is byte-identical to what the appenders left.
|
||||
"""
|
||||
lock_path = jsonl_append_lock_path(path)
|
||||
lock_fd = acquire_exclusive_file_lock(
|
||||
lock_path, timeout_sec=2.0, stale_sec=10.0, owner_aware_stale=True,
|
||||
)
|
||||
if lock_fd is None:
|
||||
return 0, 0, "lock_unavailable"
|
||||
tmp = path.with_name(path.name + ".compact.tmp")
|
||||
published = False
|
||||
try:
|
||||
digested = 0
|
||||
mismatched = 0
|
||||
consumed = 0
|
||||
with path.open("rb") as source:
|
||||
start = os.fstat(source.fileno())
|
||||
with tmp.open("wb") as sink:
|
||||
for raw in source: # streaming: one line in flight, never the file
|
||||
consumed += len(raw)
|
||||
out, outcome = _digest_line(raw, cutoff)
|
||||
sink.write(out)
|
||||
if outcome == "digested":
|
||||
digested += 1
|
||||
elif outcome == "mismatch":
|
||||
mismatched += 1
|
||||
if not digested:
|
||||
return 0, mismatched, ""
|
||||
published = _publish_if_unchanged(
|
||||
path, tmp, (consumed, int(start.st_dev), int(start.st_ino)),
|
||||
)
|
||||
if not published:
|
||||
return 0, 0, "source_changed"
|
||||
return digested, mismatched, ""
|
||||
finally:
|
||||
if not published:
|
||||
try:
|
||||
tmp.unlink()
|
||||
except OSError:
|
||||
log.debug("Failed to drop the journal compaction temp file", exc_info=True)
|
||||
release_exclusive_file_lock(lock_path, lock_fd)
|
||||
_JOURNALS = ("memory/identity_journal.jsonl", "memory/knowledge_history.jsonl",
|
||||
"memory/knowledge/patterns_history.jsonl")
|
||||
|
||||
|
||||
def compact_memory_journal_snapshots(
|
||||
|
|
@ -185,33 +23,26 @@ def compact_memory_journal_snapshots(
|
|||
*,
|
||||
now: Optional[float] = None,
|
||||
) -> Dict[str, Any]:
|
||||
"""Digest old full-text snapshots in the three memory journals."""
|
||||
from ouroboros.retention import age_cutoff, get_gc_retention_days
|
||||
"""Preserve every journal byte, including malformed and historical rows.
|
||||
|
||||
if retention_days is None:
|
||||
retention_days = get_gc_retention_days()
|
||||
cutoff = age_cutoff(retention_days, now)
|
||||
report: Dict[str, Any] = {"digested": {}, "digest_mismatch": {}, "errors": []}
|
||||
root = pathlib.Path(drive_root)
|
||||
for rel in _JOURNAL_RELS:
|
||||
path = root / rel
|
||||
if not path.exists():
|
||||
continue
|
||||
The arguments keep the previous call contract. Size facts in the existing
|
||||
startup report measure growth; a missing journal is not a measured zero.
|
||||
"""
|
||||
sizes: Dict[str, Optional[int]] = {}
|
||||
errors: list[str] = []
|
||||
for relative in _JOURNALS:
|
||||
try:
|
||||
digested, mismatched, error = _compact_one(path, cutoff)
|
||||
except OSError:
|
||||
report["errors"].append({"journal": rel.as_posix(), "error": "io_error"})
|
||||
continue
|
||||
if error:
|
||||
report["errors"].append({"journal": rel.as_posix(), "error": error})
|
||||
if digested:
|
||||
report["digested"][rel.as_posix()] = digested
|
||||
if mismatched:
|
||||
# Typed fact, not a silent skip: a stored digest that contradicts
|
||||
# its own text means one of the two is already corrupt, and the
|
||||
# content stays in full until a human looks.
|
||||
report["digest_mismatch"][rel.as_posix()] = mismatched
|
||||
return report
|
||||
info = (Path(drive_root) / relative).lstat()
|
||||
sizes[relative] = info.st_size if stat.S_ISREG(info.st_mode) else None
|
||||
if sizes[relative] is None:
|
||||
errors.append(f"{relative}: not_regular")
|
||||
except FileNotFoundError:
|
||||
sizes[relative] = None
|
||||
except OSError as exc:
|
||||
sizes[relative] = None
|
||||
errors.append(f"{relative}: {type(exc).__name__}")
|
||||
return {"digested": {}, "digest_mismatch": {}, "errors": errors,
|
||||
"journal_bytes": sizes}
|
||||
|
||||
|
||||
__all__ = ["compact_memory_journal_snapshots"]
|
||||
|
|
|
|||
98
ouroboros/memory_nomination_receipts.py
Normal file
98
ouroboros/memory_nomination_receipts.py
Normal file
|
|
@ -0,0 +1,98 @@
|
|||
"""Source-addressed pending knowledge nominations in dialogue_meta.json.
|
||||
|
||||
The history log carries full nomination bytes. This compact index makes an
|
||||
unpublished nomination visible even when its summary block becomes an era.
|
||||
A later unrelated success cannot discharge an older source identity.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
|
||||
KEY = "pending_knowledge_nominations"
|
||||
|
||||
|
||||
def load_meta(path: Path) -> dict[str, Any]:
|
||||
"""An absent cursor is new; an unreadable existing cursor is not empty.
|
||||
|
||||
This meta file now owns durable pending obligations. A permissive JSON read
|
||||
would erase them on the next consolidation. Reject duplicate keys as well:
|
||||
the second copy of an obligation field cannot silently replace the first.
|
||||
"""
|
||||
def unique_pairs(pairs: list[tuple[str, Any]]) -> dict[str, Any]:
|
||||
result: dict[str, Any] = {}
|
||||
for key, value in pairs:
|
||||
if key in result:
|
||||
raise ValueError(f"Duplicate dialogue meta key: {key}")
|
||||
result[key] = value
|
||||
return result
|
||||
|
||||
try:
|
||||
with path.open("r", encoding="utf-8") as source:
|
||||
value = json.load(source, object_pairs_hook=unique_pairs)
|
||||
except FileNotFoundError:
|
||||
# A dangling link is an existing, unreadable source, not a new cursor.
|
||||
try:
|
||||
path.lstat()
|
||||
except FileNotFoundError:
|
||||
return {}
|
||||
raise ValueError("Dialogue meta exists but cannot be read") from None
|
||||
if not isinstance(value, dict):
|
||||
raise ValueError("Dialogue meta must be a JSON object")
|
||||
_pending(value) # Refuse corrupt obligations before the first paid correction call.
|
||||
return value
|
||||
|
||||
|
||||
def _pending(meta: dict[str, Any]) -> dict[str, dict[str, Any]]:
|
||||
rows = meta.get(KEY, [])
|
||||
if not isinstance(rows, list) or any(not isinstance(row, dict) or not isinstance(row.get("id"), str)
|
||||
for row in rows):
|
||||
raise ValueError("Unreadable nomination obligations; refusing to replace their bytes")
|
||||
if len({row["id"] for row in rows}) != len(rows):
|
||||
raise ValueError("Duplicate nomination obligation IDs")
|
||||
return {row["id"]: row for row in rows}
|
||||
|
||||
|
||||
def prepare(meta: dict[str, Any], source_id: str, batches: list[tuple[Any, list[Any]]]) -> list[str]:
|
||||
"""Record every proposed entry before publication; return positional IDs."""
|
||||
pending = _pending(meta)
|
||||
ids: list[str] = []
|
||||
for block_index, (_block, entries) in enumerate(batches):
|
||||
for entry_index, entry in enumerate(entries):
|
||||
identifier = f"{source_id}:{block_index}:{entry_index}"
|
||||
ids.append(identifier)
|
||||
if identifier not in pending:
|
||||
pending[identifier] = {
|
||||
"id": identifier,
|
||||
"scope": str(entry.get("scope") or "default") if isinstance(entry, dict) else "invalid",
|
||||
"topic": str(entry.get("topic") or "") if isinstance(entry, dict) else "",
|
||||
"reason": "publication_pending",
|
||||
}
|
||||
meta[KEY] = list(pending.values())
|
||||
return ids
|
||||
|
||||
|
||||
def settle(meta: dict[str, Any], ids: list[str], outcomes: list[dict[str, Any]]) -> None:
|
||||
"""Only this exact source's successful entries retire; missing outcomes stay owed.
|
||||
|
||||
A failed or unobserved entry never expires. A future source-grounded
|
||||
resolution must address its ID explicitly; same-topic later writes cannot.
|
||||
"""
|
||||
pending = _pending(meta)
|
||||
for index, identifier in enumerate(ids):
|
||||
outcome = outcomes[index] if index < len(outcomes) else {}
|
||||
if outcome.get("ok") is True:
|
||||
pending.pop(identifier, None)
|
||||
elif identifier in pending:
|
||||
row = pending[identifier]
|
||||
row["reason"] = str(outcome.get("reason") or "outcome_missing")
|
||||
if isinstance(outcome.get("scope"), str):
|
||||
row["scope"] = outcome["scope"]
|
||||
if isinstance(outcome.get("topic"), str):
|
||||
row["topic"] = outcome["topic"]
|
||||
if pending:
|
||||
meta[KEY] = list(pending.values())
|
||||
else:
|
||||
meta.pop(KEY, None)
|
||||
|
|
@ -452,9 +452,8 @@ def _startup_retired_settings_notice(settings: dict) -> None:
|
|||
|
||||
|
||||
def _prune_event(event_type: str, keys: tuple, **reports: dict) -> None:
|
||||
"""One ``events.jsonl`` row for a GC/sweep step that did or failed something:
|
||||
``keys`` are its own evidence of material work, read across every report it
|
||||
hands in, so a healthy no-op pass stays silent instead of rowing every boot."""
|
||||
"""Emit when a report has evidence under ``keys``. GC no-ops stay silent;
|
||||
an observation report with measured journal sizes intentionally rows at boot."""
|
||||
from supervisor.state import append_jsonl
|
||||
|
||||
if any(report.get(key) for report in reports.values() for key in keys):
|
||||
|
|
@ -565,11 +564,11 @@ def _startup_prune_sweeps(*, preserve_task_sources: bool = False) -> None:
|
|||
except Exception:
|
||||
log.debug("Stale cache prune failed", exc_info=True)
|
||||
try:
|
||||
# CPL4-C16 (owner 4A): memory-journal snapshots older than GC retention
|
||||
# become digest-only (sha256 + length); fresh entries keep full text.
|
||||
# TZ-3 11A/B5 supersedes old age-digestion: measure journal growth
|
||||
# without touching historical or new full-text snapshots.
|
||||
from ouroboros.memory_journal_compaction import compact_memory_journal_snapshots
|
||||
|
||||
_prune_event("memory_journal_compaction", ("digested", "digest_mismatch", "errors"),
|
||||
_prune_event("memory_journal_observation", ("journal_bytes", "errors"),
|
||||
report=compact_memory_journal_snapshots(DATA_DIR))
|
||||
except Exception:
|
||||
log.debug("Memory journal compaction failed", exc_info=True)
|
||||
|
|
|
|||
|
|
@ -198,18 +198,39 @@ def test_partial_publication_records_the_batch_receipt_in_meta(tmp_path, fit, mo
|
|||
lambda *_a, **_k: [{"topic": "people/alex", "ok": False, "reason": "revision_conflict"}])
|
||||
ctx = ToolContext(repo_dir=tmp_path, drive_root=tmp_path, task_id="partial")
|
||||
c.consolidate(chat, blocks, meta, _Nominating(), knowledge_context=ctx)
|
||||
receipt = json.loads(meta.read_text())["last_unpublished_nominations"]
|
||||
assert receipt["failed"] == 1 and receipt["total"] == 1 and receipt["entry_id"]
|
||||
pending = json.loads(meta.read_text())["pending_knowledge_nominations"]
|
||||
assert len(pending) == 1
|
||||
assert pending[0]["topic"] == "people/alex" and pending[0]["reason"] == "revision_conflict"
|
||||
assert pending[0]["id"].endswith(":0:0")
|
||||
|
||||
|
||||
def test_a_fully_published_batch_clears_the_receipt(tmp_path, fit):
|
||||
def test_new_success_does_not_erase_an_old_failed_entry_or_legacy_receipt(tmp_path, fit, monkeypatch):
|
||||
chat, blocks, meta = _paths(tmp_path)
|
||||
_write_chat(chat, count=100, text_size=0)
|
||||
meta.parent.mkdir(parents=True, exist_ok=True)
|
||||
c.atomic_write_json(meta, {"last_unpublished_nominations": {"entry_id": "old", "failed": 3, "total": 4}})
|
||||
legacy = {"entry_id": "old", "failed": 3, "total": 4}
|
||||
c.atomic_write_json(meta, {"last_unpublished_nominations": legacy})
|
||||
ctx = ToolContext(repo_dir=tmp_path, drive_root=tmp_path, task_id="clean")
|
||||
original = c._write_knowledge_entries
|
||||
calls = 0
|
||||
|
||||
def fail_once(*args, **kwargs):
|
||||
nonlocal calls
|
||||
calls += 1
|
||||
if calls == 1:
|
||||
return [{"topic": "people/alex", "scope": "global", "ok": False,
|
||||
"reason": "revision_conflict"}]
|
||||
return original(*args, **kwargs)
|
||||
|
||||
monkeypatch.setattr(c, "_write_knowledge_entries", fail_once)
|
||||
c.consolidate(chat, blocks, meta, _Nominating(), knowledge_context=ctx)
|
||||
assert "last_unpublished_nominations" not in json.loads(meta.read_text())
|
||||
older = json.loads(meta.read_text())["pending_knowledge_nominations"][0]
|
||||
_write_chat(chat, count=200, text_size=0)
|
||||
c.consolidate(chat, blocks, meta, _Nominating(), knowledge_context=ctx)
|
||||
saved = json.loads(meta.read_text())
|
||||
assert saved["last_unpublished_nominations"] == legacy
|
||||
assert saved["pending_knowledge_nominations"] == [older]
|
||||
assert calls == 2
|
||||
|
||||
|
||||
def test_a_run_without_nominations_leaves_the_receipt_alone(tmp_path, fit):
|
||||
|
|
@ -235,8 +256,9 @@ def test_the_receipt_survives_era_compression(tmp_path, fit, monkeypatch):
|
|||
saved_blocks = json.loads(blocks.read_text())
|
||||
assert saved_blocks[0]["type"] == "era"
|
||||
assert "knowledge_writes" not in saved_blocks[0]
|
||||
receipt = json.loads(meta.read_text())["last_unpublished_nominations"]
|
||||
assert receipt["failed"] == 11 and receipt["total"] == 11
|
||||
pending = json.loads(meta.read_text())["pending_knowledge_nominations"]
|
||||
assert len(pending) == 11
|
||||
assert len({row["id"] for row in pending}) == 11
|
||||
|
||||
|
||||
# --- the Health block is where stale memory becomes visible -----------------------
|
||||
|
|
@ -297,3 +319,90 @@ def test_unreadable_receipts_do_not_raise_or_shout(tmp_path, payload):
|
|||
env = _health_env(tmp_path)
|
||||
c.atomic_write_json(tmp_path / "memory" / "dialogue_meta.json", payload)
|
||||
assert not any("DIALOGUE" in line for line in context_health._memory_health_lines(env))
|
||||
|
||||
|
||||
def test_pending_receipt_precedes_the_note_writer_and_cannot_be_replaced_by_corrupt_meta(tmp_path, fit, monkeypatch):
|
||||
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="interrupted")
|
||||
|
||||
def interrupted(*_args, **_kwargs):
|
||||
saved = json.loads(meta.read_text())
|
||||
assert len(saved["pending_knowledge_nominations"]) == 1
|
||||
assert saved["pending_knowledge_nominations"][0]["reason"] == "publication_pending"
|
||||
raise RuntimeError("simulated stop after pending publication")
|
||||
|
||||
monkeypatch.setattr(c, "_write_knowledge_entries", interrupted)
|
||||
with pytest.raises(RuntimeError, match="simulated stop"):
|
||||
c.consolidate(chat, blocks, meta, _Nominating(), knowledge_context=ctx)
|
||||
saved = json.loads(meta.read_text())
|
||||
assert saved["pending_knowledge_nominations"][0]["reason"] == "publication_pending"
|
||||
assert saved.get("last_consolidated_offset", 0) == 0
|
||||
|
||||
|
||||
def test_health_projects_three_owed_addresses_and_omission_count(tmp_path):
|
||||
env = _health_env(tmp_path)
|
||||
rows = [{"id": f"source{i}:0:0", "scope": "global", "topic": f"people/{i}",
|
||||
"reason": "revision_conflict"} for i in range(5)]
|
||||
c.atomic_write_json(tmp_path / "memory" / "dialogue_meta.json",
|
||||
{"pending_knowledge_nominations": rows})
|
||||
lines = context_health._memory_health_lines(env)
|
||||
row = next(line for line in lines if "KNOWLEDGE PUBLICATION OPEN" in line)
|
||||
assert "5 source-addressed" in row and "first 3" in row and "omitted 2" in row
|
||||
assert "source0" in row and "source2" in row and "source3" not in row
|
||||
assert "memory/knowledge_history.jsonl" in row
|
||||
|
||||
|
||||
def test_malformed_nomination_keeps_its_position_and_cannot_retire_another_entry(tmp_path):
|
||||
from ouroboros.memory_nomination_receipts import prepare, settle
|
||||
|
||||
meta = {}
|
||||
ids = prepare(meta, "source", [({}, [None, {"topic": "people/alex", "content": "Valid"}])])
|
||||
outcomes = c._write_knowledge_entries(tmp_path / "memory" / "knowledge", [None,
|
||||
{"topic": "people/alex", "content": "Valid"}])
|
||||
assert len(outcomes) == 2 and outcomes[0]["reason"] == "malformed_nomination"
|
||||
assert outcomes[1]["ok"]
|
||||
settle(meta, ids, outcomes)
|
||||
assert [row["id"] for row in meta["pending_knowledge_nominations"]] == ["source:0:0"]
|
||||
|
||||
|
||||
def test_corrupt_obligation_index_refuses_replacement(tmp_path, fit):
|
||||
chat, blocks, meta = _paths(tmp_path)
|
||||
_write_chat(chat, count=100, text_size=0)
|
||||
meta.parent.mkdir(parents=True, exist_ok=True)
|
||||
c.atomic_write_json(meta, {"pending_knowledge_nominations": {"not": "a list"}})
|
||||
ctx = ToolContext(repo_dir=tmp_path, drive_root=tmp_path, task_id="corrupt")
|
||||
with pytest.raises(ValueError, match="refusing to replace"):
|
||||
c.consolidate(chat, blocks, meta, _Nominating(), knowledge_context=ctx)
|
||||
assert not blocks.exists()
|
||||
assert json.loads(meta.read_text())["pending_knowledge_nominations"] == {"not": "a list"}
|
||||
|
||||
|
||||
@pytest.mark.parametrize("bad_bytes", [b'{"pending_knowledge_nominations":[{"id":"old"}]',
|
||||
b'["wrong top-level type"]',
|
||||
b'{"pending_knowledge_nominations":[],"pending_knowledge_nominations":[]}'])
|
||||
def test_unreadable_existing_meta_cannot_erase_obligations(tmp_path, fit, bad_bytes):
|
||||
chat, blocks, meta = _paths(tmp_path)
|
||||
_write_chat(chat, count=100, text_size=0)
|
||||
meta.parent.mkdir(parents=True, exist_ok=True)
|
||||
meta.write_bytes(bad_bytes)
|
||||
ctx = ToolContext(repo_dir=tmp_path, drive_root=tmp_path, task_id="corrupt")
|
||||
with pytest.raises(ValueError):
|
||||
c.should_consolidate(meta, chat)
|
||||
with pytest.raises(ValueError):
|
||||
c.consolidate(chat, blocks, meta, _Nominating(), knowledge_context=ctx)
|
||||
assert meta.read_bytes() == bad_bytes
|
||||
assert not blocks.exists()
|
||||
assert any("DIALOGUE META UNREADABLE" in line for line in
|
||||
context_health._memory_health_lines(_health_env(tmp_path)))
|
||||
|
||||
|
||||
def test_pending_health_disambiguates_two_entries_from_one_source(tmp_path):
|
||||
env = _health_env(tmp_path)
|
||||
c.atomic_write_json(tmp_path / "memory" / "dialogue_meta.json", {
|
||||
"pending_knowledge_nominations": [
|
||||
{"id": "a" * 64 + f":{index}:0", "reason": "revision_conflict"}
|
||||
for index in (0, 1)]})
|
||||
row = next(line for line in context_health._memory_health_lines(env)
|
||||
if "KNOWLEDGE PUBLICATION OPEN" in line)
|
||||
assert "aaaaaaaaaaaa:0:0" in row and "aaaaaaaaaaaa:1:0" in row
|
||||
|
|
|
|||
|
|
@ -229,5 +229,7 @@ def test_era_compression_cannot_erase_unpublished_knowledge_proposals(tmp_path,
|
|||
# The era object carries no knowledge_writes, so the batch receipt lives in meta:
|
||||
# without it the incomplete publication would vanish from every resident surface.
|
||||
assert "knowledge_writes" not in saved[0]
|
||||
receipt = json.loads(meta.read_text())["last_unpublished_nominations"]
|
||||
assert receipt == {"entry_id": nominations["entry_id"], "failed": 11, "total": 11}
|
||||
pending = json.loads(meta.read_text())["pending_knowledge_nominations"]
|
||||
assert len(pending) == 11
|
||||
assert all(row["id"].startswith(nominations["entry_id"] + ":") for row in pending)
|
||||
assert all(row["reason"] == "revision_required" for row in pending)
|
||||
|
|
|
|||
|
|
@ -1,223 +1,92 @@
|
|||
"""CPL4-C16 pins (owner batch №8, 4A): old journal snapshots go digest-only.
|
||||
"""Old memory-journal snapshots remain readable after every maintenance pass.
|
||||
|
||||
Fresh entries keep their full old/new text; entries older than GC retention
|
||||
keep only sha256 + length and gain ``content_digested``. Unparseable lines and
|
||||
rows without a readable ``ts`` survive byte-identical; the consciousness
|
||||
observation inbox is out of scope.
|
||||
|
||||
Audit #15-11 corrective lane: this compactor is the one sweep that destroys
|
||||
CONTENT, so it also pins that a stored digest is verified before the text it
|
||||
describes is deleted, that the rewrite publishes only a source nothing else
|
||||
touched, and that the journal is never loaded whole.
|
||||
Previously this startup sweep digested old knowledge, identity and Pattern
|
||||
Register old/new contents. A digest cannot restore the complete source after
|
||||
retention; the compatibility entry point is intentionally non-destructive.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import hashlib
|
||||
import inspect
|
||||
import json
|
||||
import os
|
||||
import pathlib
|
||||
|
||||
import pytest
|
||||
|
||||
from ouroboros import memory_journal_compaction as mjc
|
||||
from ouroboros.memory_journal_compaction import compact_memory_journal_snapshots
|
||||
from ouroboros.utils import utc_now_iso
|
||||
|
||||
_OLD_TS = "2020-01-01T00:00:00+00:00"
|
||||
_JOURNALS = (
|
||||
"memory/identity_journal.jsonl",
|
||||
"memory/knowledge_history.jsonl",
|
||||
"memory/knowledge/patterns_history.jsonl",
|
||||
"projects/example/knowledge_history.jsonl",
|
||||
"memory/scratchpad_journal.jsonl",
|
||||
)
|
||||
|
||||
|
||||
def _journal(tmp_path, rel):
|
||||
path = tmp_path / rel
|
||||
@pytest.mark.parametrize("journal", _JOURNALS)
|
||||
def test_old_journal_bytes_remain_complete_through_repeated_maintenance(tmp_path, journal):
|
||||
path = tmp_path / journal
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
return path
|
||||
row = {"ts": "2020-01-01T00:00:00+00:00", "old_content": "old\nwith Unicode Я",
|
||||
"new_content": "new\nwith Unicode Ё", "old_sha256": "legacy-mismatch"}
|
||||
original = (json.dumps(row, ensure_ascii=False) + "\n{legacy broken row\n").encode("utf-8")
|
||||
path.write_bytes(original)
|
||||
|
||||
for _ in range(2):
|
||||
report = compact_memory_journal_snapshots(tmp_path, retention_days=0,
|
||||
now=2_000_000_000.0)
|
||||
assert report["digested"] == report["digest_mismatch"] == {}
|
||||
assert report["errors"] == []
|
||||
assert report["journal_bytes"].get(journal) == (len(original) if journal in
|
||||
("memory/identity_journal.jsonl", "memory/knowledge_history.jsonl",
|
||||
"memory/knowledge/patterns_history.jsonl") else None)
|
||||
assert path.read_bytes() == original
|
||||
# Content digested in an older release cannot be restored, but is not deleted either.
|
||||
assert b"old_content" in path.read_bytes() and b"new_content" in path.read_bytes()
|
||||
|
||||
|
||||
def test_old_rows_digested_fresh_rows_kept_full(tmp_path):
|
||||
path = _journal(tmp_path, "memory/identity_journal.jsonl")
|
||||
old_row = {
|
||||
"ts": _OLD_TS, "old_content": "I was v1", "new_content": "I am v2",
|
||||
"old_sha256": hashlib.sha256(b"I was v1").hexdigest(), "old_len": 8,
|
||||
}
|
||||
fresh_row = {"ts": utc_now_iso(), "old_content": "I am v2", "new_content": "I am v3"}
|
||||
broken_line = "{not json at all\n"
|
||||
path.write_text(
|
||||
json.dumps(old_row) + "\n" + broken_line + json.dumps(fresh_row) + "\n",
|
||||
encoding="utf-8",
|
||||
)
|
||||
|
||||
report = compact_memory_journal_snapshots(tmp_path)
|
||||
|
||||
lines = path.read_text(encoding="utf-8").splitlines()
|
||||
digested = json.loads(lines[0])
|
||||
assert "old_content" not in digested and "new_content" not in digested
|
||||
assert digested["content_digested"] is True
|
||||
assert digested["old_sha256"] == hashlib.sha256(b"I was v1").hexdigest()
|
||||
assert digested["new_sha256"] == hashlib.sha256(b"I am v2").hexdigest()
|
||||
assert digested["new_len"] == len("I am v2")
|
||||
assert lines[1] == broken_line.rstrip("\n") # unreadable: byte-identical
|
||||
kept = json.loads(lines[2])
|
||||
assert kept["old_content"] == "I am v2" and kept["new_content"] == "I am v3"
|
||||
assert report["digested"] == {"memory/identity_journal.jsonl": 1}
|
||||
assert not report["digest_mismatch"] and not report["errors"]
|
||||
def test_maintenance_does_not_create_missing_journals_or_directories(tmp_path):
|
||||
root = tmp_path / "absent"
|
||||
compact_memory_journal_snapshots(root)
|
||||
assert not root.exists()
|
||||
|
||||
|
||||
@pytest.mark.parametrize("false_fact", [
|
||||
{"old_sha256": "pinned-old-hash"},
|
||||
{"old_len": 999},
|
||||
])
|
||||
def test_a_false_stored_digest_never_costs_the_text(tmp_path, false_fact):
|
||||
"""Audit #15-11: the compactor used ``setdefault``, so a stored digest that
|
||||
contradicted its own text was KEPT while the only correct copy of the
|
||||
content was deleted — the lie became the whole record. The pre-fix pin in
|
||||
this file asserted exactly that behavior (``old_sha256 == "pinned-old-hash"``
|
||||
survives the deletion of ``old_content``); it was cementing the defect and
|
||||
is reshaped above to a truthful stored digest.
|
||||
def test_startup_prune_still_reaches_compatibility_entry_point():
|
||||
import ouroboros.server_maintenance as maintenance
|
||||
|
||||
A row whose stored fact does not match its text now keeps its FULL content
|
||||
and is reported as a typed fact."""
|
||||
path = _journal(tmp_path, "memory/identity_journal.jsonl")
|
||||
row = {"ts": _OLD_TS, "old_content": "I was v1", "new_content": "I am v2", **false_fact}
|
||||
original = json.dumps(row) + "\n"
|
||||
path.write_text(original, encoding="utf-8")
|
||||
|
||||
report = compact_memory_journal_snapshots(tmp_path)
|
||||
|
||||
assert path.read_text(encoding="utf-8") == original # byte-identical
|
||||
assert not report["digested"]
|
||||
assert report["digest_mismatch"] == {"memory/identity_journal.jsonl": 1}
|
||||
assert "compact_memory_journal_snapshots" in inspect.getsource(maintenance._startup_prune_sweeps)
|
||||
|
||||
|
||||
def test_a_concurrent_append_is_never_dropped_by_the_rewrite(tmp_path):
|
||||
"""``append_jsonl`` appends WITHOUT the sidecar lock once its own
|
||||
acquisition times out, so a row can land mid-rewrite. Whether this pass
|
||||
carries it over or abandons the rewrite, the row must survive."""
|
||||
path = _journal(tmp_path, "memory/knowledge_history.jsonl")
|
||||
old_row = {"ts": _OLD_TS, "old_content": "a", "new_content": "b"}
|
||||
path.write_text(json.dumps(old_row) + "\n", encoding="utf-8")
|
||||
racing = {"ts": utc_now_iso(), "old_content": "b", "new_content": "c"}
|
||||
real_digest_line = mjc._digest_line
|
||||
fired = {"done": False}
|
||||
|
||||
def racing_digest_line(raw, cutoff):
|
||||
result = real_digest_line(raw, cutoff)
|
||||
if not fired["done"]:
|
||||
fired["done"] = True
|
||||
with path.open("ab") as unlocked_appender:
|
||||
unlocked_appender.write(json.dumps(racing).encode("utf-8") + b"\n")
|
||||
return result
|
||||
|
||||
mjc._digest_line = racing_digest_line
|
||||
def test_size_observation_does_not_follow_a_journal_symlink(tmp_path):
|
||||
target = tmp_path / "elsewhere"
|
||||
target.write_bytes(b"secret data")
|
||||
link = tmp_path / "memory" / "knowledge_history.jsonl"
|
||||
link.parent.mkdir()
|
||||
try:
|
||||
compact_memory_journal_snapshots(tmp_path)
|
||||
finally:
|
||||
mjc._digest_line = real_digest_line
|
||||
|
||||
rows = [json.loads(line) for line in path.read_text(encoding="utf-8").splitlines()]
|
||||
assert len(rows) == 2
|
||||
assert rows[1] == racing # the concurrent row is intact whichever branch ran
|
||||
|
||||
|
||||
@pytest.mark.skipif(os.name == "nt", reason="the reader holds the journal open; Windows refuses to unlink it (no FILE_SHARE_DELETE)")
|
||||
def test_a_replaced_source_aborts_the_publish(tmp_path):
|
||||
"""Identity, not just size: if the journal is swapped for a different file
|
||||
under the rewrite, the finished temp must be dropped, not published over
|
||||
whatever now lives there."""
|
||||
path = _journal(tmp_path, "memory/knowledge/patterns_history.jsonl")
|
||||
path.write_text(json.dumps({"ts": _OLD_TS, "old_content": "a", "new_content": "b"}) + "\n",
|
||||
encoding="utf-8")
|
||||
replacement = json.dumps({"ts": _OLD_TS, "topic": "someone else's file"}) + "\n"
|
||||
real_digest_line = mjc._digest_line
|
||||
fired = {"done": False}
|
||||
|
||||
def swapping_digest_line(raw, cutoff):
|
||||
result = real_digest_line(raw, cutoff)
|
||||
if not fired["done"]:
|
||||
fired["done"] = True
|
||||
path.unlink()
|
||||
path.write_text(replacement, encoding="utf-8")
|
||||
return result
|
||||
|
||||
mjc._digest_line = swapping_digest_line
|
||||
try:
|
||||
report = compact_memory_journal_snapshots(tmp_path)
|
||||
finally:
|
||||
mjc._digest_line = real_digest_line
|
||||
|
||||
assert path.read_text(encoding="utf-8") == replacement
|
||||
assert not report["digested"]
|
||||
assert report["errors"] == [
|
||||
{"journal": "memory/knowledge/patterns_history.jsonl", "error": "source_changed"},
|
||||
]
|
||||
assert not list(path.parent.glob("*.compact.tmp"))
|
||||
|
||||
|
||||
def test_the_journal_is_never_loaded_whole(tmp_path, monkeypatch):
|
||||
"""Bounded/streaming (audit #15-11c): these journals are the worst byte
|
||||
offenders in the memory plane; a whole-file read is the thing being fixed.
|
||||
Poison the whole-file readers and the compaction must still work."""
|
||||
path = _journal(tmp_path, "memory/identity_journal.jsonl")
|
||||
rows = [{"ts": _OLD_TS, "old_content": f"o{i}", "new_content": f"n{i}"} for i in range(50)]
|
||||
path.write_text("".join(json.dumps(row) + "\n" for row in rows), encoding="utf-8")
|
||||
|
||||
def _boom(self, *args, **kwargs):
|
||||
raise AssertionError(f"whole-file read of {self}")
|
||||
|
||||
monkeypatch.setattr(pathlib.Path, "read_bytes", _boom)
|
||||
link.symlink_to(target)
|
||||
except (OSError, NotImplementedError):
|
||||
pytest.skip("symlink creation unavailable")
|
||||
report = compact_memory_journal_snapshots(tmp_path)
|
||||
|
||||
assert report["digested"] == {"memory/identity_journal.jsonl": 50}
|
||||
assert report["journal_bytes"]["memory/knowledge_history.jsonl"] is None
|
||||
assert "memory/knowledge_history.jsonl: not_regular" in report["errors"]
|
||||
assert target.read_bytes() == b"secret data"
|
||||
|
||||
|
||||
def test_the_rewrite_takes_an_unstealable_lock():
|
||||
"""Owner-aware stale: elapsed time alone must never hand a second writer
|
||||
the journal this destructive rewrite is holding."""
|
||||
import inspect
|
||||
|
||||
src = inspect.getsource(mjc._compact_one)
|
||||
assert "owner_aware_stale=True" in src
|
||||
|
||||
|
||||
def test_patterns_history_gains_derived_digests(tmp_path):
|
||||
path = _journal(tmp_path, "memory/knowledge/patterns_history.jsonl")
|
||||
path.write_text(json.dumps({
|
||||
"ts": _OLD_TS, "task_id": "t", "markers": ["m"],
|
||||
"old_content": "old body", "new_content": "new body\n",
|
||||
}) + "\n", encoding="utf-8")
|
||||
|
||||
compact_memory_journal_snapshots(tmp_path)
|
||||
|
||||
row = json.loads(path.read_text(encoding="utf-8"))
|
||||
assert row["old_sha256"] == hashlib.sha256(b"old body").hexdigest()
|
||||
assert row["new_len"] == len("new body\n")
|
||||
assert "old_content" not in row and row["content_digested"] is True
|
||||
|
||||
|
||||
def test_row_without_readable_ts_keeps_full_text(tmp_path):
|
||||
path = _journal(tmp_path, "memory/knowledge_history.jsonl")
|
||||
original = json.dumps({"topic": "x", "old_content": "a", "new_content": "b"}) + "\n"
|
||||
path.write_text(original, encoding="utf-8")
|
||||
def test_startup_event_publishes_normal_journal_sizes(tmp_path, monkeypatch):
|
||||
import ouroboros.server_maintenance as maintenance
|
||||
import supervisor.state as state
|
||||
|
||||
journal = tmp_path / "memory" / "knowledge_history.jsonl"
|
||||
journal.parent.mkdir()
|
||||
journal.write_bytes(b"full historical text\n")
|
||||
rows = []
|
||||
monkeypatch.setattr(maintenance, "DATA_DIR", tmp_path)
|
||||
monkeypatch.setattr(state, "append_jsonl", lambda _path, row: rows.append(row))
|
||||
report = compact_memory_journal_snapshots(tmp_path)
|
||||
|
||||
assert path.read_text(encoding="utf-8") == original
|
||||
assert not report["digested"] and not report["errors"]
|
||||
|
||||
|
||||
def test_observation_inbox_is_out_of_scope(tmp_path):
|
||||
inbox = tmp_path / "state" / "consciousness_observations.jsonl"
|
||||
inbox.parent.mkdir(parents=True)
|
||||
original = json.dumps({"ts": _OLD_TS, "op": "enqueue", "payload": "keep me"}) + "\n"
|
||||
inbox.write_text(original, encoding="utf-8")
|
||||
|
||||
compact_memory_journal_snapshots(tmp_path)
|
||||
|
||||
assert inbox.read_text(encoding="utf-8") == original
|
||||
|
||||
|
||||
def test_startup_prune_sweeps_run_the_compaction():
|
||||
import inspect
|
||||
|
||||
import ouroboros.server_maintenance as sm
|
||||
|
||||
assert "compact_memory_journal_snapshots" in inspect.getsource(sm._startup_prune_sweeps)
|
||||
maintenance._prune_event("memory_journal_observation", ("journal_bytes", "errors"), report=report)
|
||||
assert len(rows) == 1
|
||||
assert rows[0]["report"]["journal_bytes"]["memory/knowledge_history.jsonl"] == len(b"full historical text\n")
|
||||
assert journal.read_bytes() == b"full historical text\n"
|
||||
# Pin the real startup caller, not only this unit invocation.
|
||||
source = inspect.getsource(maintenance._startup_prune_sweeps)
|
||||
assert '_prune_event("memory_journal_observation", ("journal_bytes", "errors")' in source
|
||||
|
|
|
|||
|
|
@ -571,7 +571,10 @@ def scan_data_paths(root: pathlib.Path = REPO) -> frozenset[str]:
|
|||
# one rebuildable projection per conversation written by presence_runner at the end of an executed
|
||||
# turn; it has its own row in section 2.
|
||||
# 293 -> 295: the disposable test-environment caches (``cache/pip``, ``cache/uv``; test root only).
|
||||
EXPECTED_SCAN_PATHS = 295
|
||||
# 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
|
||||
|
||||
# Scanned paths that must always be present — guards the scanner itself
|
||||
# against a silent regression that would shrink coverage while keeping counts
|
||||
|
|
|
|||
|
|
@ -162,5 +162,6 @@ def test_unread_correction_failure_remains_visible_after_dialogue_publication(tm
|
|||
assert stored["knowledge_writes"][0]["reason"] == "revision_required"
|
||||
state = json.loads(meta.read_text(encoding="utf-8"))
|
||||
assert state["last_consolidated_offset"] == 100
|
||||
assert state["last_unpublished_nominations"]["failed"] == 1
|
||||
assert len(state["pending_knowledge_nominations"]) == 1
|
||||
assert state["pending_knowledge_nominations"][0]["reason"] == "revision_required"
|
||||
assert k.read_knowledge_note(address).raw == original.raw
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue