fix: preserve complete durable learning inputs

Co-authored-by: Ouroboros <311266734+ouroboros-agent@users.noreply.github.com>
This commit is contained in:
Ouroboros 2026-08-21 19:19:16 +03:00
parent 72b502228f
commit 1b7f94973c
7 changed files with 353 additions and 60 deletions

View file

@ -3,6 +3,7 @@
from __future__ import annotations
import concurrent.futures
import hashlib
import json
import logging
import os
@ -73,6 +74,8 @@ class BackgroundConsciousness:
self._observations: queue.Queue = queue.Queue(maxsize=100)
self._deferred_events: list = []
self._tool_executor = StatefulToolExecutor()
self._identity_source_requirements: Dict[str, pathlib.Path] = {}
self._identity_source_reads: Dict[str, str] = {}
self._bg_spent_usd: float = 0.0
self._bg_budget_pct: float = float(
@ -527,6 +530,10 @@ class BackgroundConsciousness:
env = Env(repo_dir=self._repo_dir, drive_root=self._drive_root)
memory = Memory(drive_root=self._drive_root, repo_dir=self._repo_dir)
bg_task = {"id": "bg-consciousness", "type": "consciousness"}
# Per-cycle proof only. A read from an older decision envelope cannot
# authorize a later identity rewrite after the source changed.
self._identity_source_requirements = {}
self._identity_source_reads = {}
parts = [self._load_bg_prompt()]
@ -547,12 +554,43 @@ class BackgroundConsciousness:
)
try:
from ouroboros.improvement_backlog import format_backlog_digest
from ouroboros.improvement_backlog import (
backlog_path, format_backlog_digest, load_backlog_items,
)
full_backlog_digest = format_backlog_digest(
self._drive_root, limit=8, max_chars=2_147_483_647,
)
backlog_digest = format_backlog_digest(self._drive_root, limit=8, max_chars=4000)
if backlog_digest:
open_items = [
item for item in load_backlog_items(self._drive_root)
if str(item.get("status") or "open").lower() == "open"
]
if len(open_items) > 8 or backlog_digest != full_backlog_digest:
self._identity_source_requirements["improvement-backlog"] = backlog_path(
self._drive_root
)
backlog_digest += (
"\n- complete_source: call knowledge_read(topic=\"improvement-backlog\") "
"and receive the complete current record before update_identity; "
"if unavailable, abstain from that rewrite"
)
parts.append(backlog_digest)
except Exception:
try:
from ouroboros.improvement_backlog import backlog_path
path = backlog_path(self._drive_root)
if path.exists():
self._identity_source_requirements["improvement-backlog"] = path
parts.append(
"## Improvement Backlog — source unavailable\n\n"
"The named current source could not be materialized. "
"Abstain from update_identity in this cycle."
)
except Exception:
pass
log.debug("Failed to include improvement backlog in consciousness context", exc_info=True)
health_section = build_health_invariants(env)
@ -669,6 +707,23 @@ class BackgroundConsciousness:
except (json.JSONDecodeError, ValueError):
return "Failed to parse arguments."
if fn_name == "update_identity":
pending = []
for topic, path in self._identity_source_requirements.items():
try:
current_sha = hashlib.sha256(path.read_bytes()).hexdigest()
except Exception:
current_sha = ""
if not current_sha or self._identity_source_reads.get(topic) != current_sha:
pending.append(topic)
if pending:
return (
"⚠️ IDENTITY_UPDATE_ABSTAINED: direct update_identity authority remains "
"available, but this decision envelope omitted named source(s) that are "
"not completely materialized in this cycle: " + ", ".join(sorted(pending))
+ ". Read each complete current source with its existing reader, or abstain."
)
self._emit_live_log(
"tool_call_started",
tool=fn_name,
@ -745,6 +800,24 @@ class BackgroundConsciousness:
tool_args=args if isinstance(args, dict) else {},
)
if (
fn_name == "knowledge_read"
and error is None
and not timed_out
and isinstance(args, dict)
):
topic = str(args.get("topic") or "").strip()
path = self._identity_source_requirements.get(topic)
if path is not None:
try:
current = path.read_text(encoding="utf-8")
if str(result) == current and result_str == current:
self._identity_source_reads[topic] = hashlib.sha256(
current.encode("utf-8")
).hexdigest()
except Exception:
pass
args_for_log = sanitize_tool_args_for_log(fn_name, args)
if error is None and result is not None and not timed_out:
self._emit_live_log(

View file

@ -1468,7 +1468,7 @@ def _capture_context_core(
from ouroboros.improvement_backlog import format_backlog_digest
backlog_digest = format_backlog_digest(canonical_root)
if backlog_digest:
if backlog_digest and str(task.get("type") or "") in {"evolution", "deep_self_review"}:
dynamic_parts.append(backlog_digest)
except Exception:
log.debug("Failed to build improvement backlog digest", exc_info=True)

View file

@ -235,13 +235,19 @@ def _semantic_redirect_fingerprints(
out: List[Dict[str, Any]] = []
for item in items:
summary = _sanitize(item.get("summary", ""), 260)
raw_summary = str(item.get("summary") or "")
raw_category = str(item.get("category") or "process")
raw_source = str(item.get("source") or "task")
summary = _sanitize(raw_summary, 260)
if not summary:
out.append(item)
continue
category = _sanitize(item.get("category", "process"), 60) or "process"
source = _sanitize(item.get("source", "task"), 60) or "task"
fingerprint = str(item.get("fingerprint") or _stable_fingerprint(summary, category, source))
category = _sanitize(raw_category, 60) or "process"
source = _sanitize(raw_source, 60) or "task"
fingerprint = str(
item.get("fingerprint")
or _stable_fingerprint(raw_summary, raw_category, raw_source)
)
if fingerprint in existing_fps:
out.append(item) # exact hit — the locked pass bumps it, no LLM needed
continue
@ -292,12 +298,18 @@ def append_backlog_items(drive_root: Any, items: List[Dict[str, Any]]) -> int:
changed = 0
for item in items:
summary = _sanitize(item.get("summary", ""), 260)
raw_summary = str(item.get("summary") or "")
raw_category = str(item.get("category") or "process")
raw_source = str(item.get("source") or "task")
summary = _sanitize(raw_summary, 260)
if not summary:
continue
category = _sanitize(item.get("category", "process"), 60) or "process"
source = _sanitize(item.get("source", "task"), 60) or "task"
fingerprint = str(item.get("fingerprint") or _stable_fingerprint(summary, category, source))
category = _sanitize(raw_category, 60) or "process"
source = _sanitize(raw_source, 60) or "task"
fingerprint = str(
item.get("fingerprint")
or _stable_fingerprint(raw_summary, raw_category, raw_source)
)
if fingerprint in fp_to_key:
ex = by_key[fp_to_key[fingerprint]]
ex["count"] = str(_count_of(ex) + 1)
@ -379,10 +391,13 @@ def merge_backlog_text(drive_root: Any, text: str) -> int:
existing_fps = set()
fresh: List[Dict[str, Any]] = []
for it in items:
summary = _sanitize(it.get("summary", ""), 260)
category = _sanitize(it.get("category", "process"), 60) or "process"
source = _sanitize(it.get("source", "task"), 60) or "task"
fingerprint = str(it.get("fingerprint") or _stable_fingerprint(summary, category, source))
raw_summary = str(it.get("summary") or "")
raw_category = str(it.get("category") or "process")
raw_source = str(it.get("source") or "task")
fingerprint = str(
it.get("fingerprint")
or _stable_fingerprint(raw_summary, raw_category, raw_source)
)
if fingerprint in existing_fps:
continue
fresh.append(it)
@ -497,8 +512,11 @@ def groom_backlog(drive_root: Any, *, cap: int = _GROOM_CAP) -> int:
path = backlog_path(drive_root)
if not path.exists():
return 0
with _locked_text_file(path, mode="r", shared=True) as fh:
snapshot_text = fh.read()
try:
with _locked_text_file(path, mode="r", shared=True) as fh:
snapshot_text = fh.read()
except Exception:
return 0
items = _parse_backlog_items(snapshot_text)
if len(items) <= cap:
return 0
@ -514,16 +532,12 @@ def groom_backlog(drive_root: Any, *, cap: int = _GROOM_CAP) -> int:
from ouroboros.llm import LLMClient
from ouroboros.llm_observability import chat_observed
compact = [
{
k: it.get(k, "")
for k in ("id", "status", "priority", "kind", "summary", "category",
"source", "task_id", "requires_plan_review", "count",
"created_at", "fingerprint")
}
for it in fp_items
]
prompt = _GROOM_PROMPT.format(cap=cap, items_json=_json.dumps(compact, ensure_ascii=False))
# One destructive grooming call sees every complete stored record,
# including evidence/context/custom fields and immutable manual items.
# If this full prompt cannot be served, the exception path preserves the
# file; there is no smaller second call or prefix-authorized rewrite.
complete = [dict(it) for it in items]
prompt = _GROOM_PROMPT.format(cap=cap, items_json=_json.dumps(complete, ensure_ascii=False))
client = LLMClient()
resp, usage = chat_observed(
client,

View file

@ -143,7 +143,7 @@ Return ONLY a JSON object:
Rules: set promote=true ONLY when there is a concrete, high-value, self-contained code/process improvement worth a reviewed cycle right now. Prefer items already in the backlog, and weigh the solve-capability history: objective classes that historically got ABSORBED are better bets than classes that kept ending no_op/abandoned. Bias toward SMALL, TARGETED objectives that directly improve the ability to solve tasks (a sharper tool, a fixed failure mode, a removed bottleneck) over broad refactors or speculative platform work — small reviewed wins absorb; sprawling objectives historically die as no_op. Do NOT propose anything in the CLOSED / DROPPED list, or a restatement of the ACTIVE CAMPAIGN OBJECTIVE — those are already handled; if the only candidates are closed/active, return promote=false. If nothing is clearly worthwhile, return promote=false. {force_note}"""
def _closed_objectives_digest(drive_root: pathlib.Path, *, limit: int = 12, max_entries: int = 200) -> str:
def _closed_objectives_digest(drive_root: pathlib.Path) -> Optional[str]:
"""BUG3 Layer A: the objectives the chooser must NOT re-propose.
Built from the STRUCTURED ledger (state/evolution_checkpoints.jsonl), not patterns.md prose,
@ -157,19 +157,21 @@ def _closed_objectives_digest(drive_root: pathlib.Path, *, limit: int = 12, max_
from ouroboros.evolution_fingerprint import canonical_objective_fingerprint
path = pathlib.Path(drive_root) / "state" / "evolution_checkpoints.jsonl"
try:
lines = path.read_text(encoding="utf-8").splitlines()[-max_entries:]
except Exception:
if not path.exists():
return ""
try:
lines = path.read_text(encoding="utf-8").splitlines()
except Exception:
return None
by_task: Dict[str, Dict[str, Any]] = {}
order: list = []
for line in lines:
try:
row = _json.loads(line)
except Exception:
continue
return None
if not isinstance(row, dict):
continue
return None
task_id = str(row.get("task_id") or "")
if not task_id:
continue
@ -204,11 +206,7 @@ def _closed_objectives_digest(drive_root: pathlib.Path, *, limit: int = 12, max_
continue
seen.add(fp)
tag = "BLOCKED" if blocked else (str(info.get("cycle_outcome") or "DROPPED").upper())
if len(objective) > 110:
objective = objective[:110] + " …[truncated; full objective in the ledger]"
out.append(f"- [{tag}] {objective}")
if len(out) >= limit:
break
return "\n".join(out)
@ -238,7 +236,10 @@ def _decide_promotion(env: Any, task: Dict[str, Any], reflection_entry: Optional
# BUG3 Layer A: give the chooser the objectives it must NOT re-propose (the missing input
# that let a CLOSED objective be re-suggested 4-5x). Sourced from the structured ledger
# (not patterns.md prose) plus the campaign-local dropped set, deduped by fingerprint.
closed = truncate_review_artifact(_closed_objectives_digest(drive_root), 1500)
closed = _closed_objectives_digest(drive_root)
if closed is None:
log.warning("post_task_evolution: closed-objective history unavailable; abstaining")
return None
active_objective = _active_campaign_objective()
force_note = (
"The cadence already decided WHEN to evolve; choose the single most valuable "

View file

@ -665,21 +665,11 @@ def _update_patterns(drive_root: pathlib.Path, entry: Dict[str, Any]) -> None:
else:
current = _PATTERNS_HEADER
# The register is bounded by the prompt contract (max 20 rows), which fits
# well under this cap — the old 3000-char cut fed the LLM a PARTIAL table
# and the full-replace write then dropped every unseen row (memory loss).
# The cap remains only as a backstop against a pathologically bloated file.
_register_cap = 16_000
current_truncated = _truncate_with_notice(current, _register_cap)
prompt = _PATTERNS_PROMPT.format(
current_patterns=(
current_truncated
+ (
"\n\n[IMPORTANT: The current register was compacted for prompt size. "
"Preserve existing rows unless you are intentionally merging or updating them.]"
if len(current) > _register_cap else ""
)
),
# This call replaces the whole file, so its decision input must be the
# complete current register. Provider overflow/error is handled by the
# caller as an abstention; a prefix can never authorize the rewrite.
current_patterns=current,
goal=_truncate_with_notice(entry.get("goal", "?"), 200),
markers=", ".join(entry.get("key_markers", [])),
reflection=_truncate_with_notice(entry.get("reflection", ""), 500),
@ -714,6 +704,17 @@ def _update_patterns(drive_root: pathlib.Path, entry: Dict[str, Any]) -> None:
if not updated.startswith("#"):
updated = "# Pattern Register\n\n" + updated
# The lock-free LLM call may race another reflection. That makes the prompt
# source stale/incomplete, so preserve the newer register for the next pass.
try:
latest = patterns_path.read_text(encoding="utf-8") if patterns_path.exists() else _PATTERNS_HEADER
except Exception:
log.warning("Pattern register source became unavailable; preserving it")
return
if latest != current:
log.info("Pattern register changed during update; preserving the newer source")
return
append_jsonl(drive_root / "memory" / "knowledge" / "patterns_history.jsonl", {
"ts": utc_now_iso(),
"task_id": str(entry.get("task_id") or ""),

View file

@ -791,7 +791,7 @@ class TestAdvisoryReviewStatusInContext:
assert "third evidence" in dynamic_text
def test_runtime_section_includes_improvement_backlog_digest(tmp_path):
def test_improvement_backlog_digest_is_actor_scoped(tmp_path):
from ouroboros.context import build_llm_messages
from ouroboros.memory import Memory
@ -832,14 +832,24 @@ def test_runtime_section_includes_improvement_backlog_digest(tmp_path):
encoding="utf-8",
)
messages, _ = build_llm_messages(
env=FakeEnv(),
memory=Memory(drive_root=tmp_path),
task={"id": "task-a", "type": "task", "text": "hello"},
)
dynamic_text = messages[0]["content"][2]["text"]
assert "## Improvement Backlog" in dynamic_text
assert "Reduce recurring task friction around REVIEW_BLOCKED" in dynamic_text
ordinary = ({"id": "main", "type": "task", "_is_direct_chat": True},
{"id": "project", "type": "task", "project_id": "p1"},
{"id": "managed", "type": "task"},
{"id": "child", "type": "task", "delegation_role": "subagent"})
for task in ordinary:
task["text"] = "hello"
messages, _ = build_llm_messages(
env=FakeEnv(), memory=Memory(drive_root=tmp_path), task=task,
)
assert "## Improvement Backlog" not in messages[0]["content"][2]["text"]
for task_type in ("evolution", "deep_self_review"):
messages, _ = build_llm_messages(
env=FakeEnv(), memory=Memory(drive_root=tmp_path),
task={"id": task_type, "type": task_type, "text": "improve"})
dynamic_text = messages[0]["content"][2]["text"]
assert "## Improvement Backlog" in dynamic_text
assert "Reduce recurring task friction around REVIEW_BLOCKED" in dynamic_text
class TestRuntimeEnvSection:

View file

@ -0,0 +1,194 @@
"""Phase 5A: destructive learning sees complete source state or abstains."""
from __future__ import annotations
import json
import pathlib
def test_pattern_register_rewrite_receives_complete_tail(tmp_path, monkeypatch):
from ouroboros import reflection
knowledge = tmp_path / "memory" / "knowledge"
knowledge.mkdir(parents=True)
tail = "DECISIVE_PATTERN_TAIL_MUST_SURVIVE"
current = reflection._PATTERNS_HEADER + ("| old | 1 | cause | fix | open |\n" * 600) + tail
path = knowledge / "patterns.md"
path.write_text(current, encoding="utf-8")
captured = {}
monkeypatch.setattr("ouroboros.config.get_light_model", lambda: "light")
monkeypatch.setattr("ouroboros.llm.LLMClient", lambda: object())
def fake_chat(*args, **kwargs):
captured["prompt"] = kwargs["messages"][0]["content"]
return ({"content": current}, {})
monkeypatch.setattr("ouroboros.llm_observability.chat_observed", fake_chat)
reflection._update_patterns(tmp_path, {
"task_id": "task-pattern", "goal": "keep complete patterns",
"key_markers": ["TOOL_ERROR"], "reflection": "A new occurrence.",
})
assert tail in captured["prompt"]
assert tail in path.read_text(encoding="utf-8")
def test_backlog_fingerprint_uses_unsanitized_canonical_fields(tmp_path, monkeypatch):
from ouroboros.improvement_backlog import append_backlog_items, load_backlog_items
monkeypatch.setattr("ouroboros.semantic_dedup.find_semantic_duplicate_id", lambda *a, **k: None)
prefix = "x" * 300
assert append_backlog_items(tmp_path, [{
"summary": prefix + "A" * 40, "category": "process", "source": "reflection",
}]) == 1
assert append_backlog_items(tmp_path, [{
"summary": prefix + "B" * 40, "category": "process", "source": "reflection",
}]) == 1
items = load_backlog_items(tmp_path)
assert len(items) == 2
assert len({item["fingerprint"] for item in items}) == 2
def test_groom_receives_complete_records_and_preserves_on_unavailable(tmp_path, monkeypatch):
from ouroboros import improvement_backlog as ib
monkeypatch.setattr("ouroboros.semantic_dedup.find_semantic_duplicate_id", lambda *a, **k: None)
for idx in range(35):
ib.append_backlog_items(tmp_path, [{
"id": f"ibl-{idx}", "fingerprint": f"fp-{idx}", "summary": f"item {idx}",
"category": "process", "source": "reflection",
"evidence": f"complete-evidence-{idx}", "context": f"complete-context-{idx}",
"proposed_next_step": f"complete-next-step-{idx}",
}])
items = ib.load_backlog_items(tmp_path)
keep = [{"id": i["id"], "fingerprint": i["fingerprint"], "summary": i["summary"]}
for i in items[:20]]
captured = {}
monkeypatch.setattr("ouroboros.config.get_light_model", lambda: "light")
monkeypatch.setattr("ouroboros.llm.LLMClient", lambda: object())
def fake_chat(*args, **kwargs):
captured["prompt"] = kwargs["messages"][0]["content"]
return ({"content": json.dumps(keep)}, {})
monkeypatch.setattr("ouroboros.llm_observability.chat_observed", fake_chat)
assert ib.groom_backlog(tmp_path, cap=30) == 20
for expected in ("complete-evidence-34", "complete-context-34", "complete-next-step-34"):
assert expected in captured["prompt"]
before = ib.backlog_path(tmp_path).read_text(encoding="utf-8")
real_locked = ib._locked_text_file
def unavailable(path, mode, *, shared=False):
if mode == "r":
raise PermissionError("backlog unavailable")
return real_locked(path, mode, shared=shared)
monkeypatch.setattr(ib, "_locked_text_file", unavailable)
assert ib.groom_backlog(tmp_path, cap=10) == 0
assert ib.backlog_path(tmp_path).read_text(encoding="utf-8") == before
def test_closed_objective_before_old_horizon_reaches_chooser(tmp_path, monkeypatch):
from ouroboros import post_task_evolution as pte
state = tmp_path / "state"
state.mkdir(parents=True)
old_objective = "OLD CLOSED OBJECTIVE MUST NOT BE PROMOTED AGAIN"
rows = [{"task_id": "old", "kind": "cycle_outcome", "cycle_outcome": "absorbed",
"campaign_objective": old_objective}]
rows.extend({"task_id": f"new-{i}", "kind": "cycle_outcome", "cycle_outcome": "absorbed",
"campaign_objective": f"new objective {i}"} for i in range(230))
(state / "evolution_checkpoints.jsonl").write_text(
"\n".join(json.dumps(row) for row in rows), encoding="utf-8",
)
captured = {}
def fake_chat(*args, **kwargs):
captured["prompt"] = kwargs["messages"][0]["content"]
return ({"content": '{"promote": false, "objective": ""}'}, {})
monkeypatch.setattr("ouroboros.llm_observability.chat_observed", fake_chat)
monkeypatch.setattr(pte, "_active_campaign_objective", lambda: "")
env = type("Env", (), {"drive_root": tmp_path})()
decision = pte._decide_promotion(env, {"id": "root"}, {"reflection": "done"}, object(), force=False)
assert decision and decision["promote"] is False
assert old_objective in captured["prompt"]
def test_closed_objective_unavailable_abstains_before_chooser(tmp_path, monkeypatch):
from ouroboros import post_task_evolution as pte
state = tmp_path / "state"
state.mkdir(parents=True)
ledger = state / "evolution_checkpoints.jsonl"
ledger.write_text("{}\n", encoding="utf-8")
real_read_text = pathlib.Path.read_text
def unreadable(self, *args, **kwargs):
if self == ledger:
raise PermissionError("ledger unavailable")
return real_read_text(self, *args, **kwargs)
monkeypatch.setattr(pathlib.Path, "read_text", unreadable)
called = []
def chooser(*args, **kwargs):
called.append(True)
return ({"content": '{"promote": true, "objective": "unsafe"}'}, {})
monkeypatch.setattr("ouroboros.llm_observability.chat_observed", chooser)
env = type("Env", (), {"drive_root": tmp_path})()
assert pte._decide_promotion(env, {"id": "root"}, {"reflection": "done"}, object(), force=False) is None
assert called == []
def _bg_fixture(tmp_path):
from ouroboros.consciousness import BackgroundConsciousness
from ouroboros.improvement_backlog import append_backlog_items
repo_dir = pathlib.Path(__file__).parents[1]
(tmp_path / "logs").mkdir(parents=True)
(tmp_path / "state").mkdir(parents=True)
(tmp_path / "state" / "state.json").write_text("{}", encoding="utf-8")
for idx in range(10):
append_backlog_items(tmp_path, [{
"id": f"ibl-bg-{idx}", "fingerprint": f"fp-bg-{idx}",
"summary": f"background item {idx}", "category": "identity", "source": "reflection",
}])
return BackgroundConsciousness(tmp_path, repo_dir, None, lambda: None)
def _tool_call(name, args, call_id):
return {"id": call_id, "function": {"name": name, "arguments": json.dumps(args)}}
def test_bgc_direct_identity_update_requires_complete_named_omission(tmp_path):
bc = _bg_fixture(tmp_path)
try:
context = bc._build_context()
assert "knowledge_read" in context and "improvement-backlog" in context
content = "I remain directly self-authoring after complete source materialization."
blocked = bc._execute_tool(_tool_call("update_identity", {"content": content}, "u1"), [])
assert "IDENTITY_UPDATE_ABSTAINED" in blocked
read = bc._execute_tool(_tool_call("knowledge_read", {"topic": "improvement-backlog"}, "r1"), [])
assert "background item 9" in read
updated = bc._execute_tool(_tool_call("update_identity", {"content": content}, "u2"), [])
assert updated.startswith("OK: identity updated")
journal = tmp_path / "memory" / "identity_journal.jsonl"
assert journal.exists() and content in journal.read_text(encoding="utf-8")
finally:
bc._tool_executor.shutdown(wait=False, cancel_futures=True)
def test_bgc_unavailable_named_omission_abstains_without_approval_flow(tmp_path):
bc = _bg_fixture(tmp_path)
try:
bc._build_context()
(tmp_path / "memory" / "knowledge" / "improvement-backlog.md").unlink()
content = "I retain direct authority but abstain when the named source is unavailable."
result = bc._execute_tool(_tool_call("update_identity", {"content": content}, "u1"), [])
assert "IDENTITY_UPDATE_ABSTAINED" in result
assert "approval" not in result.lower()
assert not (tmp_path / "memory" / "identity_journal.jsonl").exists()
finally:
bc._tool_executor.shutdown(wait=False, cancel_futures=True)