Fix background observation settlement and dedupe

Co-authored-by: Ouroboros <311266734+ouroboros-agent@users.noreply.github.com>
This commit is contained in:
Ouroboros 2026-08-22 00:21:55 +03:00
parent f692d019ea
commit 824c16638f
3 changed files with 533 additions and 97 deletions

View file

@ -4,6 +4,8 @@ from __future__ import annotations
import concurrent.futures
import hashlib
import contextlib
import inspect
import json
import logging
import os
@ -16,31 +18,36 @@ from typing import TYPE_CHECKING, Any, Callable, Dict, List, Optional, Sequence
if TYPE_CHECKING:
from ouroboros.tools.registry import ToolRegistry
from ouroboros.loop_tool_execution import StatefulToolExecutor, _truncate_tool_result
from ouroboros.utils import (
append_jsonl,
emit_log_event,
read_text,
sanitize_tool_args_for_log,
sanitize_tool_result_for_log,
truncate_for_log,
utc_now_iso,
)
from ouroboros.config import get_consciousness_model, resolve_effort
from ouroboros.pricing import infer_provider_from_model
from ouroboros.llm import LLMClient, add_usage
from ouroboros.memory import Memory
from ouroboros.context import (
build_runtime_section, build_memory_sections,
build_recent_sections, build_health_invariants,
build_knowledge_sections, build_governance_sections, safe_read,
build_governance_sections,
build_health_invariants,
build_knowledge_sections,
build_memory_sections,
build_recent_sections,
build_runtime_section,
safe_read,
)
from ouroboros.context_budget import (
BG_CONTEXT_MAX_CHARS,
BG_CONTEXT_WARN_CHARS,
BG_STATE_JSON_WARN_CHARS,
)
from ouroboros.llm import LLMClient, add_usage
from ouroboros.loop_tool_execution import StatefulToolExecutor, _truncate_tool_result
from ouroboros.memory import Memory
from ouroboros.platform_layer import acquire_exclusive_file_lock, release_exclusive_file_lock
from ouroboros.pricing import infer_provider_from_model
from ouroboros.utils import (
append_jsonl,
emit_log_event,
jsonl_append_lock_path,
read_text,
sanitize_tool_args_for_log,
sanitize_tool_result_for_log,
truncate_for_log,
utc_now_iso,
)
_OBSERVATIONS_REL = pathlib.Path("state") / "consciousness_observations.jsonl"
_OBSERVATION_SOURCE_REF = (
@ -103,36 +110,129 @@ class BackgroundConsciousness:
@property
def is_running(self) -> bool:
return self._running and self._thread is not None and self._thread.is_alive()
thread = getattr(self, "_thread", None)
return bool(getattr(self, "_running", False) and thread is not None and thread.is_alive())
@property
def is_paused(self) -> bool:
return self._paused
return bool(getattr(self, "_paused", False))
def _observation_lock_for_instance(self) -> threading.RLock:
"""Lazily restore observation fields for object.__new__ overlap tests."""
lock = getattr(self, "_observation_lock", None)
if lock is None:
lock = threading.RLock()
self._observation_lock = lock
return lock
def _stop_requested(self) -> bool:
event = getattr(self, "_stop_event", None)
return bool(event is not None and event.is_set())
@contextlib.contextmanager
def _observation_writer_lock(self, path: pathlib.Path):
"""Use the same sidecar lock seam as append_jsonl for store transactions."""
lock_path = jsonl_append_lock_path(path)
lock_fd = acquire_exclusive_file_lock(
lock_path,
timeout_sec=2.0,
stale_sec=10.0,
poll_sec=0.01,
)
if lock_fd is None:
yield False
return
try:
yield True
finally:
release_exclusive_file_lock(lock_path, lock_fd)
@staticmethod
def _append_observation_line_locked(path: pathlib.Path, row: Dict[str, Any]) -> bool:
"""Append one row while the shared JSONL writer lock is held."""
try:
data = (json.dumps(row, ensure_ascii=False) + "\n").encode("utf-8")
path.parent.mkdir(parents=True, exist_ok=True)
fd = os.open(str(path), os.O_WRONLY | os.O_CREAT | os.O_APPEND, 0o644)
try:
view = memoryview(data)
while view:
written = os.write(fd, view)
if written <= 0:
return False
view = view[written:]
finally:
os.close(fd)
return True
except Exception:
log.warning("Failed to append background observation row", exc_info=True)
return False
def _mark_cycle_settlement_gap(self, reason: str) -> None:
self._cycle_settlement_failed = True
# A missing/unknown durable receipt means this cognition cycle cannot
# truthfully settle its observation snapshot, even if the model later
# returns a final message.
self._cycle_ack_allowed = False
reasons = getattr(self, "_cycle_settlement_reasons", None)
if reasons is None:
reasons = []
self._cycle_settlement_reasons = reasons
if reason not in reasons:
reasons.append(reason)
log.error("Background consciousness durable settlement gap: %s", reason)
def _append_cycle_receipt(
self,
path: pathlib.Path,
row: Dict[str, Any],
*,
label: str,
) -> bool:
"""Persist a cycle receipt and latch an explicit writer failure."""
try:
written = append_jsonl(path, row)
except Exception as exc:
self._mark_cycle_settlement_gap(f"{label}: {type(exc).__name__}")
return False
# A receipt is authoritative only when the writer reports truthy
# success. ``None`` from a callback is an unknown write, not proof.
if not written:
self._mark_cycle_settlement_gap(f"{label}: append_jsonl returned false")
return False
return True
@property
def _model(self) -> str:
return get_consciousness_model()
def status_snapshot(self) -> Dict[str, Any]:
pending = self._snapshot_pending_observations()
oldest = pending[0].get("time") if pending else ""
with self._observation_lock:
gaps = list((self._read_observation_state().get("gap_reasons") or ()))
with self._observation_lock_for_instance():
state = self._read_observation_state()
pending_count = int(state.get("pending_count") or 0)
oldest = str(state.get("oldest_pending_at") or "")
gap_count = int(
state.get("gap_count", len(state.get("gap_reasons") or ())) or 0
)
return {
"running": bool(self.is_running),
"paused": bool(self._paused),
"next_wakeup_sec": int(self._next_wakeup_sec),
"last_cycle_started_at": self._last_cycle_started_at,
"last_cycle_finished_at": self._last_cycle_finished_at,
"last_idle_reason": self._last_idle_reason,
"last_error": self._last_error,
"paused": bool(self.is_paused),
"next_wakeup_sec": int(getattr(self, "_next_wakeup_sec", 300) or 300),
"last_cycle_started_at": getattr(self, "_last_cycle_started_at", ""),
"last_cycle_finished_at": getattr(self, "_last_cycle_finished_at", ""),
"last_idle_reason": getattr(self, "_last_idle_reason", "stopped"),
"last_error": getattr(self, "_last_error", ""),
# Status is deliberately content-free: it gives the operator a
# truthful horizon and a resolvable store, never observation text.
"pending_observation_count": len(pending),
"pending_observation_count": pending_count,
"oldest_observation_at": oldest or "",
"observation_source": _OBSERVATION_SOURCE_REF,
"observation_source_complete": not bool(gaps),
"observation_gap_count": len(gaps),
"observation_source_complete": gap_count == 0,
"observation_gap_count": gap_count,
}
def start(self) -> str:
@ -213,27 +313,30 @@ class BackgroundConsciousness:
"ref": ref,
}
path = self._observation_store_path()
with self._observation_lock:
state = self._read_observation_state()
if identifier in state["rows"]:
return False
try:
written = append_jsonl(path, {"op": "enqueue", **row})
except (TypeError, ValueError) as exc:
log.error("Background observation %s is not JSON-serializable: %s", identifier, exc)
return False
if not written:
log.error("Failed to durably enqueue background observation %s", identifier)
return False
state["rows"][identifier] = row
self._refresh_observation_signature(path)
with self._observation_lock_for_instance():
# The process lock protects this instance; the shared sidecar lock
# protects sibling processes. Re-read while holding both so the
# stable-ID check and append are one transaction.
with self._observation_writer_lock(path) as locked:
if not locked:
log.error("Failed to lock background observation store %s", path)
return False
state = self._read_observation_state(force=True)
if identifier in state["rows"]:
return False
if not self._append_observation_line_locked(path, {"op": "enqueue", **row}):
log.error("Failed to durably enqueue background observation %s", identifier)
return False
state["rows"][identifier] = row
self._refresh_observation_summary(state)
self._refresh_observation_signature(path)
self._wakeup_event.set()
return True
def _observation_store_path(self) -> pathlib.Path:
return self._drive_root / _OBSERVATIONS_REL
def _read_observation_state(self) -> Dict[str, Any]:
def _read_observation_state(self, *, force: bool = False) -> Dict[str, Any]:
"""Read append-only observations into a cached stable-ID index.
A process rebuilds from the durable source once, then only performs a
@ -242,11 +345,10 @@ class BackgroundConsciousness:
"""
path = self._observation_store_path()
signature = self._observation_signature(path)
if (
self._observation_state_cache is not None
and signature == self._observation_store_signature
):
return self._observation_state_cache
cached = getattr(self, "_observation_state_cache", None)
cached_signature = getattr(self, "_observation_store_signature", None)
if not force and cached is not None and signature == cached_signature:
return cached
rows: Dict[str, Dict[str, Any]] = {}
acked = set()
@ -271,10 +373,7 @@ class BackgroundConsciousness:
if not isinstance(item, dict):
gap_reasons.append(f"non-object observation row at line {line_no}")
continue
if item.get("op") in (None, "enqueue", "observation") and not (
item.get("id") or item.get("observation_id")
):
gap_reasons.append(f"observation row missing id at line {line_no}")
if not self._validate_observation_row(item, line_no, gap_reasons):
continue
self._index_observation_row(item, rows, acked)
@ -282,10 +381,70 @@ class BackgroundConsciousness:
"rows": rows,
"acked": acked,
"gap_reasons": gap_reasons,
"gap_count": len(gap_reasons),
}
self._refresh_observation_summary(self._observation_state_cache)
self._observation_store_signature = signature
return self._observation_state_cache
@staticmethod
def _validate_observation_row(
item: Dict[str, Any],
line_no: int,
gap_reasons: List[str],
) -> bool:
"""Validate the small durable row contract before indexing it."""
op = item.get("op")
if op not in (None, "enqueue", "observation", "ack"):
gap_reasons.append(f"unknown observation op at line {line_no}")
return False
identifier = item.get("id", item.get("observation_id"))
if not isinstance(identifier, str) or not identifier.strip():
gap_reasons.append(f"observation row missing id at line {line_no}")
return False
if op == "ack":
return True
# Canonical enqueue rows are strict. Legacy/old observation rows retain
# aliases and an optional ref so existing producers remain readable.
if op == "enqueue":
required = ("source", "kind", "time", "payload", "ref")
missing = [key for key in required if key not in item]
if missing:
gap_reasons.append(
f"enqueue row missing {','.join(missing)} at line {line_no}"
)
return False
elif op == "observation":
if not str(item.get("source") or "").strip():
gap_reasons.append(f"observation row missing source at line {line_no}")
return False
if not str(item.get("kind") or "").strip():
gap_reasons.append(f"observation row missing kind at line {line_no}")
return False
if not str(item.get("time") or item.get("observed_at") or "").strip():
gap_reasons.append(f"observation row missing time at line {line_no}")
return False
if "payload" not in item and "text" not in item:
gap_reasons.append(f"observation row missing payload at line {line_no}")
return False
return True
@staticmethod
def _refresh_observation_summary(state: Dict[str, Any]) -> None:
rows = state.get("rows") or {}
acked = state.get("acked") or set()
pending_count = 0
oldest = ""
for identifier, row in rows.items():
if identifier in acked:
continue
pending_count += 1
if not oldest:
oldest = str(row.get("time") or "")
state["pending_count"] = pending_count
state["oldest_pending_at"] = oldest
@staticmethod
def _index_observation_row(
item: Dict[str, Any],
@ -326,7 +485,7 @@ class BackgroundConsciousness:
def _snapshot_pending_observations(self) -> List[Dict[str, Any]]:
"""Return a non-destructive, insertion-ordered pending snapshot."""
with self._observation_lock:
with self._observation_lock_for_instance():
state = self._read_observation_state()
return [row for identifier, row in state["rows"].items()
if identifier not in state["acked"]]
@ -335,24 +494,32 @@ class BackgroundConsciousness:
"""Append acknowledgements after a settled successful cognition cycle."""
if not observations:
return True
if getattr(self, "_cycle_settlement_failed", False):
log.error("Cannot acknowledge observations while cycle settlement has a durable gap")
return False
path = self._observation_store_path()
with self._observation_lock:
state = self._read_observation_state()
if state.get("gap_reasons"):
log.error("Cannot acknowledge observations with unreadable inbox rows: %s", state["gap_reasons"][:3])
return False
for observation in observations:
identifier = str(observation.get("id") or "")
if not identifier or identifier in state["acked"]:
continue
if not append_jsonl(path, {
"op": "ack",
"id": identifier,
"time": utc_now_iso(),
}):
log.error("Failed to acknowledge background observation %s", identifier)
with self._observation_lock_for_instance():
with self._observation_writer_lock(path) as locked:
if not locked:
log.error("Failed to lock background observation store %s", path)
return False
state["acked"].add(identifier)
state = self._read_observation_state(force=True)
if state.get("gap_reasons"):
log.error("Cannot acknowledge observations with unreadable inbox rows: %s", state["gap_reasons"][:3])
return False
for observation in observations:
identifier = str(observation.get("id") or "")
if not identifier or identifier in state["acked"]:
continue
if not self._append_observation_line_locked(path, {
"op": "ack",
"id": identifier,
"time": utc_now_iso(),
}):
log.error("Failed to acknowledge background observation %s", identifier)
return False
state["acked"].add(identifier)
self._refresh_observation_summary(state)
self._refresh_observation_signature(path)
return True
@ -362,7 +529,7 @@ class BackgroundConsciousness:
shown = list(observations[-_OBSERVATION_RENDER_LIMIT:])
omitted = max(0, total - len(shown))
source = _OBSERVATION_SOURCE_REF
with self._observation_lock:
with self._observation_lock_for_instance():
gaps = list((self._read_observation_state().get("gap_reasons") or ()))
lines = [
f"## Pending observations (total={total}; showing={len(shown)}; "
@ -510,9 +677,14 @@ class BackgroundConsciousness:
def _think_scoped(self) -> bool:
"""Run one context/LLM/tools cycle; False preserves skip/error status."""
self._cycle_settlement_failed = False
self._cycle_settlement_reasons = []
self._cycle_ack_allowed = True
if not hasattr(self, "_deferred_events"):
self._deferred_events = []
observation_snapshot = self._snapshot_pending_observations()
try:
context = self._build_context(observations=observation_snapshot)
context = self._build_cycle_context(observation_snapshot)
except OverflowError as exc:
# P1: skip the cycle rather than silently truncating cognitive context.
log.warning("consciousness: wakeup cycle skipped: %s", exc)
@ -547,7 +719,8 @@ class BackgroundConsciousness:
if not _use_local_consciousness else None
)
for round_idx in range(1, self._max_bg_rounds + 1):
if self._paused:
if self.is_paused:
self._cycle_ack_allowed = False
break
if target is not None:
from ouroboros.openai_chat_dispatch import projected_context_size_bytes
@ -565,11 +738,11 @@ class BackgroundConsciousness:
"Groom memory to continue."
)
self._last_idle_reason = "context_overflow"
append_jsonl(self._drive_root / "logs" / "events.jsonl", {
self._append_cycle_receipt(self._drive_root / "logs" / "events.jsonl", {
"ts": utc_now_iso(),
"type": "consciousness_context_overflow",
"error": error,
})
}, label="context overflow")
return False
if physical_chars > BG_CONTEXT_WARN_CHARS:
log.warning(
@ -621,11 +794,12 @@ class BackgroundConsciousness:
if not self._check_budget():
self._last_idle_reason = "budget_blocked"
append_jsonl(self._drive_root / "logs" / "events.jsonl", {
self._cycle_ack_allowed = False
self._append_cycle_receipt(self._drive_root / "logs" / "events.jsonl", {
"ts": utc_now_iso(),
"type": "bg_budget_exceeded_mid_cycle",
"round": round_idx,
})
}, label="budget blocked")
break
if self._event_queue is not None:
@ -663,7 +837,8 @@ class BackgroundConsciousness:
self._emit_progress(content)
if self._paused:
if self.is_paused:
self._cycle_ack_allowed = False
break
if content and not tool_calls:
@ -688,13 +863,13 @@ class BackgroundConsciousness:
break
if all_pending_events and self._event_queue is not None:
if self._paused:
if self.is_paused:
self._deferred_events.extend(all_pending_events)
else:
for evt in all_pending_events:
self._event_queue.put(evt)
thought_written = append_jsonl(self._drive_root / "logs" / "events.jsonl", {
thought_written = self._append_cycle_receipt(self._drive_root / "logs" / "events.jsonl", {
"ts": utc_now_iso(),
"type": "consciousness_thought",
"thought_preview": (final_content or "")[:300],
@ -710,23 +885,31 @@ class BackgroundConsciousness:
)
if key in cycle_usage
},
})
}, label="thought receipt")
# A paused/stopped cycle may have performed no complete cognition;
# leave its snapshot pending for the next wake. Acknowledgement is
# the final durable step, after the thought receipt and all tool
# state writes above have settled.
if not self._paused and not self._stop_event.is_set():
if not thought_written or not self._ack_observations(observation_snapshot):
if self.is_paused or self._stop_requested():
self._cycle_ack_allowed = False
if not final_content.strip():
self._cycle_ack_allowed = False
if not self.is_paused and not self._stop_requested():
if not thought_written or not self._cycle_ack_allowed or not self._ack_observations(observation_snapshot):
self._last_idle_reason = "observation_ack_pending"
return False
elif observation_snapshot:
self._last_idle_reason = "observation_ack_pending"
return False
except Exception as e:
self._cycle_ack_allowed = False
self._emit_live_log("llm_round_error", round=round_idx, model=model, error=repr(e))
append_jsonl(self._drive_root / "logs" / "events.jsonl", {
self._append_cycle_receipt(self._drive_root / "logs" / "events.jsonl", {
"ts": utc_now_iso(),
"type": "consciousness_llm_error",
"error": repr(e),
})
}, label="llm error")
self._last_idle_reason = "llm_error"
# Back off persistent provider/tool failures.
self._next_wakeup_sec = min(self._next_wakeup_sec * 2, self._wakeup_max)
@ -751,7 +934,7 @@ class BackgroundConsciousness:
persist_locally = self._event_queue is None or chat_id is None
if self._event_queue is not None and chat_id is not None:
try:
if self._paused:
if self.is_paused:
self._deferred_events.append(entry)
else:
self._event_queue.put(entry)
@ -761,6 +944,31 @@ class BackgroundConsciousness:
if persist_locally:
append_jsonl(self._drive_root / "logs" / "progress.jsonl", entry)
def _build_cycle_context(
self,
observations: Sequence[Dict[str, Any]],
) -> str:
"""Call context builders across the pre-observation compatibility seam.
A few overlap tests construct this class with ``object.__new__`` and
provide an older zero-argument builder. Inspect the callable before
invoking it so a real ``TypeError`` raised *inside* a current builder
is never mistaken for an unsupported keyword.
"""
builder = self._build_context
try:
parameters = inspect.signature(builder).parameters.values()
except (TypeError, ValueError):
parameters = ()
accepts_observations = any(
parameter.name == "observations"
or parameter.kind is inspect.Parameter.VAR_KEYWORD
for parameter in parameters
)
if accepts_observations:
return builder(observations=observations)
return builder()
def _load_bg_prompt(self) -> str:
"""Load consciousness system prompt."""
prompt_path = self._repo_dir / "prompts" / "CONSCIOUSNESS.md"
@ -933,7 +1141,7 @@ class BackgroundConsciousness:
def _build_registry(self) -> "ToolRegistry":
"""Create a ToolRegistry scoped to background-allowed tools."""
from ouroboros.tools.registry import ToolRegistry, ToolEntry
from ouroboros.tools.registry import ToolEntry, ToolRegistry
registry = ToolRegistry(repo_dir=self._repo_dir, drive_root=self._drive_root)
@ -1034,6 +1242,7 @@ class BackgroundConsciousness:
except (TimeoutError, concurrent.futures.TimeoutError):
self._tool_executor.reset()
timed_out = True
self._cycle_ack_allowed = False
result = f"[TIMEOUT after {timeout_sec}s]"
self._emit_live_log(
"tool_call_timeout",
@ -1041,14 +1250,15 @@ class BackgroundConsciousness:
args=sanitize_tool_args_for_log(fn_name, args if isinstance(args, dict) else {}),
timeout_sec=timeout_sec,
)
append_jsonl(self._drive_root / "logs" / "events.jsonl", {
self._append_cycle_receipt(self._drive_root / "logs" / "events.jsonl", {
"ts": utc_now_iso(),
"type": "consciousness_tool_timeout",
"tool": fn_name,
"timeout_sec": timeout_sec,
})
}, label=f"tool timeout:{fn_name}")
if error is not None:
self._cycle_ack_allowed = False
self._emit_live_log(
"tool_call_finished",
tool=fn_name,
@ -1056,12 +1266,12 @@ class BackgroundConsciousness:
is_error=True,
result_preview=repr(error),
)
append_jsonl(self._drive_root / "logs" / "events.jsonl", {
self._append_cycle_receipt(self._drive_root / "logs" / "events.jsonl", {
"ts": utc_now_iso(),
"type": "consciousness_tool_error",
"tool": fn_name,
"error": repr(error),
})
}, label=f"tool error:{fn_name}")
result = f"Error: {repr(error)}"
for evt in self._registry._ctx.pending_events:
@ -1100,12 +1310,12 @@ class BackgroundConsciousness:
is_error=False,
result_preview=sanitize_tool_result_for_log(truncate_for_log(result_str, 500)),
)
append_jsonl(self._drive_root / "logs" / "tools.jsonl", {
self._append_cycle_receipt(self._drive_root / "logs" / "tools.jsonl", {
"ts": utc_now_iso(),
"tool": fn_name,
"source": "consciousness",
"args": args_for_log,
"result_preview": sanitize_tool_result_for_log(truncate_for_log(result_str, 2000)),
})
}, label=f"tool receipt:{fn_name}")
return result_str

View file

@ -541,7 +541,11 @@ def test_background_preserves_receipt_and_aggregates_round_disclosures(
monkeypatch.setattr(observed, "chat_observed", chat_observed)
monkeypatch.setattr(consciousness, "get_consciousness_model", lambda: "anthropic/claude-future")
monkeypatch.setattr(consciousness, "append_jsonl", lambda _path, row: events.append(row))
monkeypatch.setattr(
consciousness,
"append_jsonl",
lambda _path, row: (events.append(row) or True),
)
monkeypatch.setattr(consciousness.BackgroundConsciousness, "_build_context", lambda _self: "ctx")
monkeypatch.setattr(consciousness.BackgroundConsciousness, "_tool_schemas", lambda _self: [])
monkeypatch.setattr(consciousness.BackgroundConsciousness, "_check_budget", lambda _self: True)

View file

@ -4,6 +4,7 @@ from __future__ import annotations
import concurrent.futures
import json
import multiprocessing
import pathlib
import queue
from unittest.mock import MagicMock, patch
@ -27,6 +28,30 @@ def _make(tmp_path):
), drive
def _enqueue_from_process(drive_str, result_queue, start_event):
"""Exercise the durable writer seam from an independent process."""
from ouroboros.consciousness import BackgroundConsciousness
drive = pathlib.Path(drive_str)
repo = drive.parent / "repo-process"
repo.mkdir(exist_ok=True)
with patch.object(BackgroundConsciousness, "_build_registry", return_value=MagicMock()):
instance = BackgroundConsciousness(
drive_root=drive,
repo_dir=repo,
event_queue=queue.Queue(),
owner_chat_id_fn=lambda: None,
)
start_event.wait(10)
result_queue.put(instance.inject_observation(
"process-payload",
observation_id="process-stable-id",
source="process",
kind="trace",
ref={"worker": "process"},
))
def test_observations_are_append_only_and_deduplicated_over_100_rows(tmp_path):
bc, drive = _make(tmp_path)
for index in range(125):
@ -61,6 +86,54 @@ def test_concurrent_enqueue_same_id_has_one_durable_row(tmp_path):
assert [row["id"] for row in rows] == ["same-id"]
def test_cross_instance_enqueue_same_id_is_atomically_deduplicated(tmp_path):
first, drive = _make(tmp_path)
second, _ = _make(tmp_path)
def enqueue(instance):
return instance.inject_observation(
"same", observation_id="cross-instance-id", source="instance"
)
with concurrent.futures.ThreadPoolExecutor(max_workers=2) as pool:
results = list(pool.map(enqueue, (first, second)))
assert sum(results) == 1
rows = [
json.loads(line)
for line in (drive / "state" / "consciousness_observations.jsonl").read_text().splitlines()
]
assert [row["id"] for row in rows] == ["cross-instance-id"]
def test_cross_process_enqueue_same_id_is_atomically_deduplicated(tmp_path):
try:
context = multiprocessing.get_context("fork")
except ValueError:
pytest.skip("fork process context is unavailable on this platform")
_, drive = _make(tmp_path)
start_event = context.Event()
result_queue = context.Queue()
workers = [
context.Process(
target=_enqueue_from_process,
args=(str(drive), result_queue, start_event),
)
for _ in range(2)
]
for worker in workers:
worker.start()
start_event.set()
for worker in workers:
worker.join(10)
assert worker.exitcode == 0
assert sorted(result_queue.get(timeout=2) for _ in workers) == [False, True]
rows = [
json.loads(line)
for line in (drive / "state" / "consciousness_observations.jsonl").read_text().splitlines()
]
assert [row["id"] for row in rows] == ["process-stable-id"]
def test_context_snapshot_is_bounded_and_status_has_no_payload(tmp_path):
bc, _ = _make(tmp_path)
for index in range(15):
@ -95,6 +168,20 @@ def test_status_uses_cached_projection_after_initial_rebuild(tmp_path, monkeypat
assert calls["read"] == 0
def test_status_snapshot_never_materializes_pending_payloads(tmp_path, monkeypatch):
bc, _ = _make(tmp_path)
bc.inject_observation("payload", observation_id="status-id")
def fail_if_snapshot_called():
raise AssertionError("status_snapshot must use cached counters")
monkeypatch.setattr(bc, "_snapshot_pending_observations", fail_if_snapshot_called)
first = bc.status_snapshot()
second = bc.status_snapshot()
assert first["pending_observation_count"] == 1
assert second["oldest_observation_at"] == first["oldest_observation_at"]
def test_malformed_store_row_is_disclosed_and_blocks_ack(tmp_path):
bc, drive = _make(tmp_path)
bc.inject_observation("known", observation_id="known")
@ -109,6 +196,41 @@ def test_malformed_store_row_is_disclosed_and_blocks_ack(tmp_path):
assert [row["id"] for row in restarted._snapshot_pending_observations()] == ["known"]
def test_unknown_and_malformed_rows_create_gaps_and_block_ack(tmp_path):
bc, drive = _make(tmp_path)
store = drive / "state" / "consciousness_observations.jsonl"
store.parent.mkdir(parents=True, exist_ok=True)
rows = [
{
"op": "enqueue",
"id": "valid",
"source": "test",
"kind": "trace",
"time": "2026-08-22T00:00:00Z",
"payload": "kept",
"ref": None,
},
{"op": "future-op", "id": "unknown", "payload": "ignored"},
{"op": "ack"},
{"op": "ack", "id": 42},
{
"op": "enqueue",
"id": "missing-ref",
"source": "test",
"kind": "trace",
"time": "2026-08-22T00:00:01Z",
"payload": "ignored",
},
]
store.write_text("".join(json.dumps(row) + "\n" for row in rows), encoding="utf-8")
restarted, _ = _make(tmp_path)
status = restarted.status_snapshot()
assert status["observation_source_complete"] is False
assert status["observation_gap_count"] == 4
assert [row["id"] for row in restarted._snapshot_pending_observations()] == ["valid"]
assert restarted._ack_observations(restarted._snapshot_pending_observations()) is False
def test_success_acks_only_cycle_snapshot_and_later_rows_stay_pending(tmp_path, monkeypatch):
bc, _ = _make(tmp_path)
bc.inject_observation("first", observation_id="first")
@ -155,6 +277,106 @@ def test_failed_cycle_keeps_observations_pending(tmp_path, monkeypatch, failure)
assert [row["id"] for row in bc._snapshot_pending_observations()] == ["replay"]
def test_tool_receipt_write_failure_forbids_cycle_ack(tmp_path, monkeypatch):
from types import SimpleNamespace
from ouroboros import consciousness
bc, _ = _make(tmp_path)
bc.inject_observation("must settle", observation_id="tool-gap")
bc._build_context = lambda **_: "context"
bc._tool_schemas = lambda: [{
"type": "function",
"function": {"name": "read_file", "parameters": {"type": "object"}},
}]
bc._llm = SimpleNamespace(_resolve_remote_target=lambda _model: {
"provider": "openai", "resolved_model": "test-model",
})
bc._registry.get_timeout.return_value = 1
bc._registry.execute.return_value = "tool-result"
bc._registry._ctx = SimpleNamespace(pending_events=[])
responses = [
({"content": "", "tool_calls": [{
"id": "tool-call",
"function": {"name": "read_file", "arguments": "{}"},
}]}, {"cost": 0.0}),
({"content": "settled"}, {"cost": 0.0}),
]
monkeypatch.setattr(
"ouroboros.llm_observability.chat_observed",
lambda *args, **kwargs: responses.pop(0),
)
monkeypatch.setattr(consciousness, "get_consciousness_model", lambda: "openai/test-model")
monkeypatch.setattr(consciousness, "resolve_effort", lambda _slot: "medium")
monkeypatch.setattr(consciousness, "append_jsonl", lambda path, row: path.name != "tools.jsonl")
assert bc._think_scoped() is False
assert [row["id"] for row in bc._snapshot_pending_observations()] == ["tool-gap"]
def test_budget_refusal_after_response_forbids_cycle_ack(tmp_path, monkeypatch):
from types import SimpleNamespace
from ouroboros import consciousness
bc, drive = _make(tmp_path)
bc.inject_observation("budget retry", observation_id="budget-gap")
bc._build_context = lambda **_: "context"
bc._tool_schemas = lambda: []
bc._llm = SimpleNamespace(_resolve_remote_target=lambda _model: {
"provider": "openai", "resolved_model": "test-model",
})
seen = {"chat": 0}
def fake_chat(*args, **kwargs):
seen["chat"] += 1
return {"content": "would settle"}, {"cost": 0.0}
monkeypatch.setattr("ouroboros.llm_observability.chat_observed", fake_chat)
monkeypatch.setattr(consciousness, "get_consciousness_model", lambda: "openai/test-model")
monkeypatch.setattr(consciousness, "resolve_effort", lambda _slot: "medium")
monkeypatch.setattr(bc, "_check_budget", lambda: False)
assert bc._think_scoped() is False
assert seen["chat"] == 1
assert [row["id"] for row in bc._snapshot_pending_observations()] == ["budget-gap"]
events = [
json.loads(line)
for line in (drive / "logs" / "events.jsonl").read_text().splitlines()
]
assert any(row.get("type") == "bg_budget_exceeded_mid_cycle" for row in events)
def test_budget_receipt_write_failure_latches_settlement_gap(tmp_path, monkeypatch):
from types import SimpleNamespace
from ouroboros import consciousness
bc, _ = _make(tmp_path)
bc.inject_observation("budget receipt retry", observation_id="budget-receipt-gap")
bc._build_context = lambda **_: "context"
bc._tool_schemas = lambda: []
bc._llm = SimpleNamespace(_resolve_remote_target=lambda _model: {
"provider": "openai", "resolved_model": "test-model",
})
monkeypatch.setattr(
"ouroboros.llm_observability.chat_observed",
lambda *args, **kwargs: ({"content": "would settle"}, {"cost": 0.0}),
)
monkeypatch.setattr(consciousness, "get_consciousness_model", lambda: "openai/test-model")
monkeypatch.setattr(consciousness, "resolve_effort", lambda _slot: "medium")
monkeypatch.setattr(bc, "_check_budget", lambda: False)
monkeypatch.setattr(
consciousness,
"append_jsonl",
lambda _path, row: row.get("type") != "bg_budget_exceeded_mid_cycle",
)
assert bc._think_scoped() is False
assert bc._cycle_settlement_failed is True
assert [row["id"] for row in bc._snapshot_pending_observations()] == ["budget-receipt-gap"]
def test_restart_replays_unacknowledged_observation_by_id(tmp_path):
bc, _ = _make(tmp_path)
bc.inject_observation("survive", observation_id="restart-id")