diff --git a/.agent/harness/hooks/_episodic_io.py b/.agent/harness/hooks/_episodic_io.py index e434908..dab741f 100644 --- a/.agent/harness/hooks/_episodic_io.py +++ b/.agent/harness/hooks/_episodic_io.py @@ -39,7 +39,104 @@ def append_jsonl(path: str, entry: dict) -> dict: try: f.write(payload) f.flush() + # Durability must finish before the flock drops. auto_dream + # rewrites this file under the same lock as soon as it can + # acquire it, so a later fsync in the caller can sync the + # wrong generation. + os.fsync(f.fileno()) finally: if _HAVE_FLOCK: fcntl.flock(f.fileno(), fcntl.LOCK_UN) return entry + + +def _first_matching_action(raw: bytes, match_action: str) -> dict | None: + """First JSON object whose action equals `match_action`. + + File order is the canonical order. Blank lines and corrupt lines + are skipped. A row with an empty timestamp does not count. + """ + for line in raw.decode("utf-8").splitlines(): + line = line.strip() + if not line: + continue + try: + row = json.loads(line) + except json.JSONDecodeError: + continue + if not isinstance(row, dict) or row.get("action") != match_action: + continue + timestamp = row.get("timestamp") + if isinstance(timestamp, str) and timestamp: + return row + return None + + +def append_jsonl_once(path: str, entry: dict, *, match_action: str) -> dict: + """Append `entry` unless `match_action` is already present. + + The read and the optional append share one `LOCK_EX` on `path`, the + same flock `append_jsonl` and `auto_dream` already take. This is + episodic serialization, not a candidate-directory lock. No second + lock file is created. + + Returns the earliest existing row with that action, or `entry` when + this call appended it. Older duplicate rows are left in place. + + Without `fcntl` the check-and-append is best-effort, matching the + pre-lock baseline. `O_APPEND` keeps the write at end of file. + """ + os.makedirs(os.path.dirname(path), exist_ok=True) + with open(path, "a+b") as handle: + if _HAVE_FLOCK: + fcntl.flock(handle.fileno(), fcntl.LOCK_EX) + try: + handle.seek(0) + existing = _first_matching_action(handle.read(), match_action) + if existing is not None: + return existing + payload = (json.dumps(entry) + "\n").encode("utf-8") + handle.write(payload) + handle.flush() + os.fsync(handle.fileno()) + return entry + finally: + if _HAVE_FLOCK: + fcntl.flock(handle.fileno(), fcntl.LOCK_UN) + + +def has_jsonl_timestamp(path: str, timestamp: str) -> bool: + """Return whether a JSONL row has this exact timestamp under LOCK_EX. + + The lock is held for the complete read. auto_dream rewrites the episodic + file while holding the same lock, so recovery must not inspect a truncated + generation between its truncate and rewrite. + + A missing file is absence, not an error. This probe does not create the + JSONL or its parent directory. FileNotFoundError during open is the same + absence. Any other OSError propagates so recovery does not treat a failed + read as missing evidence and delete a resumable temp. + + Without fcntl the lock is a no-op, matching the pre-lock baseline. + """ + if not timestamp or not os.path.isfile(path): + return False + try: + handle = open(path, "rb") + except FileNotFoundError: + return False + with handle: + if _HAVE_FLOCK: + fcntl.flock(handle.fileno(), fcntl.LOCK_EX) + try: + for line in handle: + try: + row = json.loads(line.decode("utf-8")) + except (UnicodeDecodeError, json.JSONDecodeError): + continue + if isinstance(row, dict) and row.get("timestamp") == timestamp: + return True + return False + finally: + if _HAVE_FLOCK: + fcntl.flock(handle.fileno(), fcntl.LOCK_UN) diff --git a/.agent/harness/hooks/on_failure.py b/.agent/harness/hooks/on_failure.py index 277e057..ce9d402 100644 --- a/.agent/harness/hooks/on_failure.py +++ b/.agent/harness/hooks/on_failure.py @@ -1,5 +1,5 @@ """Failures are learning. High pain score + rewrite flag after repeat offenses.""" -import json, datetime, os +import json, datetime, os, sys from ._provenance import build_source from ._episodic_io import append_jsonl @@ -73,4 +73,10 @@ def on_failure(skill_name, action, error, context="", confidence=0.9, f"Flag for rewrite." ) entry["pain_score"] = 10 - return append_jsonl(EPISODIC, entry) + try: + return append_jsonl(EPISODIC, entry) + except OSError as err: + # Failure telemetry must not mask the original failure or terminate + # the hook process when the filesystem rejects fsync. + print(f"WARNING: episodic failure log write failed: {err}", file=sys.stderr) + return entry diff --git a/.agent/harness/hooks/post_execution.py b/.agent/harness/hooks/post_execution.py index 2727b97..b5e5ab3 100644 --- a/.agent/harness/hooks/post_execution.py +++ b/.agent/harness/hooks/post_execution.py @@ -1,5 +1,5 @@ """Runs after every action. Appends a structured entry to episodic memory.""" -import datetime, os +import datetime, os, sys from ._provenance import build_source from ._episodic_io import append_jsonl @@ -31,4 +31,10 @@ def log_execution(skill_name, action, result, success, reflection="", "source": build_source(skill_name), "evidence_ids": list(evidence_ids) if evidence_ids else [], } - return append_jsonl(EPISODIC, entry) + try: + return append_jsonl(EPISODIC, entry) + except OSError as err: + # Episodic logging is observability. A sync failure must not turn + # a completed tool action into a failed hook. + print(f"WARNING: episodic log write failed: {err}", file=sys.stderr) + return entry diff --git a/.agent/harness/hooks/test_episodic_hooks.py b/.agent/harness/hooks/test_episodic_hooks.py new file mode 100644 index 0000000..b7440ba --- /dev/null +++ b/.agent/harness/hooks/test_episodic_hooks.py @@ -0,0 +1,58 @@ +"""Hook regression tests for non-fatal episodic filesystem failures.""" + +import contextlib +import io +import sys +import unittest +from pathlib import Path + +HARNESS = Path(__file__).resolve().parents[1] +if str(HARNESS) not in sys.path: + sys.path.insert(0, str(HARNESS)) + +from hooks import on_failure, post_execution # noqa: E402 + + +class EpisodicHookTest(unittest.TestCase): + def test_post_execution_sync_error_does_not_escape(self): + original_append = post_execution.append_jsonl + original_source = post_execution.build_source + post_execution.append_jsonl = lambda *_a, **_k: (_ for _ in ()).throw( + OSError("forced fsync failure") + ) + post_execution.build_source = lambda skill: {"skill": skill} + try: + stderr = io.StringIO() + with contextlib.redirect_stderr(stderr): + result = post_execution.log_execution( + "test", "action", "result", True + ) + self.assertEqual(result["result"], "success") + self.assertIn("forced fsync failure", stderr.getvalue()) + finally: + post_execution.append_jsonl = original_append + post_execution.build_source = original_source + + def test_on_failure_sync_error_does_not_escape(self): + original_append = on_failure.append_jsonl + original_source = on_failure.build_source + original_count = on_failure._count_recent_failures + on_failure.append_jsonl = lambda *_a, **_k: (_ for _ in ()).throw( + OSError("forced fsync failure") + ) + on_failure.build_source = lambda skill: {"skill": skill} + on_failure._count_recent_failures = lambda _skill: 0 + try: + stderr = io.StringIO() + with contextlib.redirect_stderr(stderr): + result = on_failure.on_failure("test", "action", "boom") + self.assertEqual(result["result"], "failure") + self.assertIn("forced fsync failure", stderr.getvalue()) + finally: + on_failure.append_jsonl = original_append + on_failure.build_source = original_source + on_failure._count_recent_failures = original_count + + +if __name__ == "__main__": + unittest.main() diff --git a/.agent/tools/learn.py b/.agent/tools/learn.py index 8cb2620..94bd0dc 100644 --- a/.agent/tools/learn.py +++ b/.agent/tools/learn.py @@ -18,7 +18,7 @@ candidate file is removed so `show.py` / `REVIEW_QUEUE.md` don't show orphaned dead-ends. """ -import argparse, datetime, json, os, subprocess, sys +import argparse, datetime, errno, json, os, subprocess, sys, tempfile if hasattr(sys.stdout, "reconfigure"): sys.stdout.reconfigure(encoding="utf-8", errors="replace") @@ -28,6 +28,7 @@ CANDIDATES = os.path.join(BASE, "memory/candidates") sys.path.insert(0, os.path.join(BASE, "harness")) sys.path.insert(0, os.path.join(BASE, "memory")) +from hooks._episodic_io import append_jsonl_once, has_jsonl_timestamp # noqa: E402 from text import word_set # noqa: E402 from cluster import pattern_id # noqa: E402 @@ -60,18 +61,23 @@ def _lesson_already_appended(cid): return False -def _append_episodic_mirror(cid, claim, ts, source="learn"): - """Mirror a manual stage into AGENT_LEARNINGS.jsonl so evidence_ids - referencing `ts` resolve to a real episodic record — matching the - auto-derived candidate path's existing behavior. Never raises; a - failure here must not block staging (same fail-open posture as - _lesson_already_appended's read-only probe). +def _append_episodic_mirror(cid, claim, source="learn"): + """Return the canonical timestamp for `manual-stage:{cid}`. + + Inserts one episodic row when that action is absent. A later call + returns the earliest existing row's timestamp and does not append. + The check and the insert share the episodic flock inside + ``append_jsonl_once``. This does not lock ``CANDIDATES``. + + Raises OSError on write failure. ``stage()`` must not publish a + candidate that references a timestamp until this returns. """ - episodic_path = os.path.join(BASE, "memory/episodic/AGENT_LEARNINGS.jsonl") + ts = datetime.datetime.now(datetime.timezone.utc).isoformat() + action = f"manual-stage:{cid}" entry = { "timestamp": ts, "skill": "learn", - "action": f"manual-stage:{cid}", + "action": action, "result": "success", "detail": f"Manually staged lesson {cid} via .agent/tools/learn.py: {claim!r}", "pain_score": 1, @@ -81,17 +87,222 @@ def _append_episodic_mirror(cid, claim, ts, source="learn"): "source": {"skill": "learn", "profile": "manual", "run_id": f"manual_{cid[:6]}"}, "evidence_ids": [ts], } + canonical = append_jsonl_once(_episodic_path(), entry, match_action=action) + canonical_ts = canonical.get("timestamp") + if not isinstance(canonical_ts, str) or not canonical_ts: + raise OSError(f"episodic mirror for {cid} has no timestamp") + return canonical_ts + + +def _episodic_path(): + return os.path.join(BASE, "memory/episodic/AGENT_LEARNINGS.jsonl") + + +def _fsync_dir(directory): + """Best-effort directory fsync. Portability failures are ignored. + + Directory durability is not part of the evidence invariant. EINVAL, + ENOTSUP, EBADF, and EPERM (Windows and some filesystems) must not + decide whether a candidate is published. + """ + try: + fd = os.open(directory, os.O_RDONLY) + except OSError as err: + if err.errno in (errno.EINVAL, errno.ENOTSUP, errno.EBADF, errno.EPERM): + return + raise + try: + try: + os.fsync(fd) + except OSError as err: + if err.errno not in (errno.EINVAL, errno.ENOTSUP, errno.EBADF, errno.EPERM): + raise + finally: + os.close(fd) + + +def _leftover_temps(cid): + """Absolute paths of `.{cid}.*.tmp` files in CANDIDATES, sorted.""" + if not os.path.isdir(CANDIDATES): + return [] + prefix = f".{cid}." + found = [] + for name in os.listdir(CANDIDATES): + if name.startswith(prefix) and name.endswith(".tmp"): + found.append(os.path.join(CANDIDATES, name)) + return sorted(found) + + +def _evidence_landed(episodic_path, timestamp): + """True when a parsed JSONL row has this exact timestamp field. + + Raw substring search is intentionally not used. A timestamp that + appears only inside another string must not count as evidence. The + complete read holds the episodic LOCK_EX. Without fcntl that lock + is a no-op. + """ + return has_jsonl_timestamp(episodic_path, timestamp) + + +def _load_json_object(path): + """Load a JSON object. ``None`` means corrupt or the wrong shape. + + ``OSError`` propagates. A transient read error must not look like + corrupt JSON, or recovery will delete a fsynced temp. + """ + try: + with open(path, encoding="utf-8") as stream: + payload = json.load(stream) + except json.JSONDecodeError: + return None + if not isinstance(payload, dict): + return None + return payload + + +def _one_evidence_id(payload): + evidence = payload.get("evidence_ids") + if ( + not isinstance(evidence, list) + or len(evidence) != 1 + or not isinstance(evidence[0], str) + or not evidence[0] + ): + return None + return evidence[0] + + +def _resumable_evidence(temp_path, cid): + """Evidence timestamp if this temp is safe to publish, else None.""" + payload = _load_json_object(temp_path) + if payload is None or payload.get("id") != cid: + return None + evidence = _one_evidence_id(payload) + if evidence is None: + return None + if not _evidence_landed(_episodic_path(), evidence): + return None + return evidence + + +def _remove_or_raise(temp_path): + try: + os.remove(temp_path) + except OSError as err: + raise OSError(f"failed to remove temp {temp_path}: {err}") from err + + +def _resume_temp(temp_path, path): try: - with open(episodic_path, "a", encoding="utf-8") as f: - f.write(json.dumps(entry) + "\n") + os.replace(temp_path, path) + except OSError as err: + raise OSError( + f"{err}; candidate publish failed; fsynced temp kept at {temp_path}" + ) from err + # The rename already published the candidate. A directory fsync + # failure must not make the caller retry and append a second mirror. + try: + _fsync_dir(CANDIDATES) except OSError: - pass # fail-open: staging must succeed even if the mirror write fails + pass + + +def _published_evidence(path, cid): + if not os.path.isfile(path): + return None + payload = _load_json_object(path) + if payload is None or payload.get("id") != cid: + return None + return _one_evidence_id(payload) + + +def _shared_resumable_evidence(leftovers, cid): + """One timestamp when every temp is resumable with that same value.""" + stamps = [] + for temp_path in leftovers: + evidence = _resumable_evidence(temp_path, cid) + if evidence is None: + return None + stamps.append(evidence) + if len(set(stamps)) != 1: + return None + return stamps[0] + + +def _publish_shared_temps(leftovers, cid, path, evidence): + published = _published_evidence(path, cid) + if published is not None and published >= evidence: + for temp_path in leftovers: + _remove_or_raise(temp_path) + return True + _resume_temp(leftovers[0], path) + for temp_path in leftovers[1:]: + _remove_or_raise(temp_path) + return True + + +def _resolve_leftovers(cid, path): + """Finish or discard a prior transaction before a new one starts. + + Returns True when `{cid}.json` is already the file the caller should + return. Returns False when the caller must start a fresh publish. + + A resumable temp (valid JSON, matching id, evidence timestamp present + as a JSONL `timestamp` field) is published with `os.replace` and no + second mirror. If the published evidence timestamp is the same or + newer, that temp is stale and is removed instead. An evidence-less + or corrupt temp is + never published. + + The episodic flock makes the mirror row single-flight for one action. + Temp files are not locked. Two callers can still each leave a temp + with the same evidence timestamp. Those twins are one transaction: + publish one and delete the rest. Temps whose evidence differs stay + fail-closed. This function does not add a lock. + """ + leftovers = _leftover_temps(cid) + if not leftovers: + return False + if len(leftovers) > 1: + shared = _shared_resumable_evidence(leftovers, cid) + if shared is None: + raise OSError( + "ambiguous leftover temps for {cid}; refusing to publish: {paths}".format( + cid=cid, + paths=", ".join(leftovers), + ) + ) + return _publish_shared_temps(leftovers, cid, path, shared) + temp_path = leftovers[0] + evidence = _resumable_evidence(temp_path, cid) + if evidence is None: + _remove_or_raise(temp_path) + return os.path.isfile(path) + published = None + if os.path.isfile(path): + published_payload = _load_json_object(path) + if published_payload is not None and published_payload.get("id") == cid: + published = _one_evidence_id(published_payload) + # ISO-8601 timestamps from datetime.isoformat() sort lexicographically. + if published is not None and published >= evidence: + _remove_or_raise(temp_path) + return True + _resume_temp(temp_path, path) + return True def stage(claim, conditions, source="learn", importance=7): os.makedirs(CANDIDATES, exist_ok=True) cid = pattern_id(claim, conditions) - now = datetime.datetime.now(datetime.timezone.utc).isoformat() + path = os.path.join(CANDIDATES, f"{cid}.json") + # Resolve any prior temp before creating another one, so a retry + # cannot append a second mirror while the first payload is still + # recoverable. + if _resolve_leftovers(cid, path): + return cid, path + # The mirror assigns the evidence id. A repeat call returns the + # earliest row's timestamp and does not append another line. + now = _append_episodic_mirror(cid, claim, source) candidate = { "id": cid, "key": f"manual_{cid[:6]}", @@ -109,10 +320,41 @@ def stage(claim, conditions, source="learn", importance=7): "decisions": [{"ts": now, "action": "staged", "reviewer": source}], "rejection_count": 0, } - path = os.path.join(CANDIDATES, f"{cid}.json") - with open(path, "w", encoding="utf-8") as f: - json.dump(candidate, f, indent=2) - _append_episodic_mirror(cid, claim, now, source) + # Publish the candidate only after the episodic mirror succeeds, so a + # visible staged file never carries a dangling evidence_id. Temp file + # stays in CANDIDATES so os.replace stays same-filesystem. + temp_path = None + try: + with tempfile.NamedTemporaryFile( + mode="w", + encoding="utf-8", + dir=CANDIDATES, + prefix=f".{cid}.", + suffix=".tmp", + delete=False, + ) as stream: + temp_path = stream.name + json.dump(candidate, stream, indent=2) + stream.flush() + os.fsync(stream.fileno()) + _fsync_dir(CANDIDATES) + os.replace(temp_path, path) + temp_path = None + # Rename already published. Directory fsync must not fail the call. + try: + _fsync_dir(CANDIDATES) + except OSError: + pass + except BaseException as primary: + # The mirror already returned. Keep a temp that was written so + # the next call can publish it. Do not turn KeyboardInterrupt + # into OSError. + if temp_path is not None and isinstance(primary, OSError): + raise OSError( + f"{primary}; candidate publish failed; " + f"fsynced temp kept at {temp_path}" + ) from primary + raise return cid, path @@ -148,7 +390,11 @@ def main(): # in ways Codex caught. Fixed list here is the stable signature. conditions = sorted(word_set(claim)) - cid, path = stage(claim, conditions) + try: + cid, path = stage(claim, conditions) + except OSError as err: + print(f"ERROR: {err}", file=sys.stderr) + sys.exit(1) print(f"staged candidate {cid}") print(f" path: {path}") print(f" conditions: {conditions}") diff --git a/.agent/tools/test_learn_episodic_mirror.py b/.agent/tools/test_learn_episodic_mirror.py index 808e3cc..104eb09 100644 --- a/.agent/tools/test_learn_episodic_mirror.py +++ b/.agent/tools/test_learn_episodic_mirror.py @@ -24,8 +24,12 @@ def _load_learn(base_dir): """Load .agent/tools/learn.py with BASE/CANDIDATES pointed at base_dir. Sibling modules (text.word_set, cluster.pattern_id) are stubbed so the - test needs no part of the harness beyond learn.py itself. + test needs no part of the harness beyond learn.py itself. hooks._episodic_io + is imported from the real tree (stdlib-only locked append helper). """ + harness_dir = str(Path(__file__).resolve().parents[1] / "harness") + if harness_dir not in sys.path: + sys.path.insert(0, harness_dir) for name, attrs in [ ("text", {"word_set": lambda *a, **k: set()}), ("cluster", {"pattern_id": lambda claim, cond: "testcid" + str(abs(hash((claim, tuple(cond)))))[:6]}), @@ -48,37 +52,501 @@ def _load_learn(base_dir): def _episodic(base_dir): path = os.path.join(base_dir, "memory", "episodic", "AGENT_LEARNINGS.jsonl") - if not os.path.exists(path): + if not os.path.isfile(path): return [] with open(path, encoding="utf-8") as f: return [json.loads(line) for line in f if line.strip()] +def _names(candidates, suffix): + return sorted( + name for name in os.listdir(candidates) if name.endswith(suffix) + ) + + +CLAIM = "Serialize timestamps in UTC" +CONDITIONS = ["timestamps", "utc"] + + class EpisodicMirrorTest(unittest.TestCase): def setUp(self): self.tmp = tempfile.mkdtemp() def test_stage_writes_one_episodic_mirror(self): mod = _load_learn(self.tmp) - cid, _ = mod.stage("Serialize timestamps in UTC", ["timestamps", "utc"]) + cid, _ = mod.stage(CLAIM, CONDITIONS) entries = _episodic(self.tmp) mirrors = [e for e in entries if e.get("action") == f"manual-stage:{cid}"] self.assertEqual(len(mirrors), 1) + self.assertEqual(_names(mod.CANDIDATES, ".json"), [f"{cid}.json"]) + self.assertEqual(_names(mod.CANDIDATES, ".tmp"), []) def test_evidence_id_resolves_to_the_mirror(self): mod = _load_learn(self.tmp) - cid, path = mod.stage("Serialize timestamps in UTC", ["timestamps", "utc"]) + cid, path = mod.stage(CLAIM, CONDITIONS) candidate = json.loads(Path(path).read_text()) evidence_ts = candidate["evidence_ids"][0] matching = [e for e in _episodic(self.tmp) if e["timestamp"] == evidence_ts] self.assertEqual(len(matching), 1) self.assertEqual(matching[0]["evidence_ids"], [evidence_ts]) - def test_append_mirror_fails_open_on_write_error(self): + def test_success_fsyncs_jsonl_before_unlock_and_directory_twice(self): + mod = _load_learn(self.tmp) + import hooks._episodic_io as episodic_io + + if not episodic_io._HAVE_FLOCK: + self.skipTest("fcntl flock not available") + + real_fsync = episodic_io.os.fsync + + dir_calls = [] + real_dir = mod._fsync_dir + + def _record_dir(directory): + dir_calls.append(directory) + return real_dir(directory) + + order = [] + real_flock = episodic_io.fcntl.flock + + def _record_flock(fd, operation): + if operation == episodic_io.fcntl.LOCK_UN: + order.append("unlock") + else: + order.append("lock") + return real_flock(fd, operation) + + def _record_fsync(fd): + order.append("fsync") + return real_fsync(fd) + + episodic_io.os.fsync = _record_fsync + episodic_io.fcntl.flock = _record_flock + mod._fsync_dir = _record_dir + try: + mod.stage(CLAIM, CONDITIONS) + finally: + episodic_io.os.fsync = real_fsync + episodic_io.fcntl.flock = real_flock + mod._fsync_dir = real_dir + self.assertIn("lock", order) + lock_at = order.index("lock") + unlock_at = order.index("unlock", lock_at) + self.assertIn("fsync", order[lock_at + 1:unlock_at]) + self.assertEqual(dir_calls, [mod.CANDIDATES, mod.CANDIDATES]) + self.assertEqual(len(_names(mod.CANDIDATES, ".json")), 1) + self.assertEqual(_names(mod.CANDIDATES, ".tmp"), []) + self.assertEqual(len(_episodic(self.tmp)), 1) + + def test_stage_fails_closed_when_mirror_write_errors(self): + mod = _load_learn(self.tmp) + + def _boom(*_a, **_k): + raise OSError("forced mirror-write failure") + + mod._append_episodic_mirror = _boom + with self.assertRaises(OSError): + mod.stage(CLAIM, CONDITIONS) + self.assertEqual(_names(mod.CANDIDATES, ".json"), []) + self.assertEqual(_names(mod.CANDIDATES, ".tmp"), []) + self.assertEqual(_episodic(self.tmp), []) + + def test_real_append_jsonl_failure_leaves_no_temp(self): + mod = _load_learn(self.tmp) + episodic_path = os.path.join( + self.tmp, "memory", "episodic", "AGENT_LEARNINGS.jsonl") + os.mkdir(episodic_path) + with self.assertRaises(OSError): + mod.stage(CLAIM, CONDITIONS) + self.assertTrue(os.path.isdir(episodic_path)) + self.assertEqual(_names(mod.CANDIDATES, ".json"), []) + self.assertEqual(_names(mod.CANDIDATES, ".tmp"), []) + + def test_stage_keeps_fsynced_candidate_when_publish_fails(self): + mod = _load_learn(self.tmp) + original_replace = mod.os.replace + + def _boom(*_a, **_k): + raise OSError("forced publish failure") + + mod.os.replace = _boom + try: + with self.assertRaises(OSError) as caught: + mod.stage(CLAIM, CONDITIONS) + finally: + mod.os.replace = original_replace + + self.assertIn("forced publish failure", str(caught.exception)) + self.assertIn("fsynced temp kept at", str(caught.exception)) + self.assertEqual(_names(mod.CANDIDATES, ".json"), []) + staged = _names(mod.CANDIDATES, ".tmp") + self.assertEqual(len(staged), 1) + staged_path = os.path.join(mod.CANDIDATES, staged[0]) + self.assertIn(staged_path, str(caught.exception)) + self.assertIn( + '"claim": "Serialize timestamps in UTC"', + Path(staged_path).read_text(), + ) + entries = _episodic(self.tmp) + self.assertEqual(len(entries), 1) + self.assertEqual(entries[0]["result"], "success") + + def test_retry_recovers_temp_without_a_second_mirror(self): + mod = _load_learn(self.tmp) + original_replace = mod.os.replace + + def _boom(*_a, **_k): + raise OSError("forced publish failure") + + mod.os.replace = _boom + try: + with self.assertRaises(OSError): + mod.stage(CLAIM, CONDITIONS) + finally: + mod.os.replace = original_replace + + failed = _episodic(self.tmp) + self.assertEqual(len(failed), 1) + cid, path = mod.stage(CLAIM, CONDITIONS) + self.assertEqual(_names(mod.CANDIDATES, ".json"), [f"{cid}.json"]) + self.assertEqual(_names(mod.CANDIDATES, ".tmp"), []) + self.assertEqual(_episodic(self.tmp), failed) + candidate = json.loads(Path(path).read_text()) + self.assertEqual(candidate["evidence_ids"], [failed[0]["timestamp"]]) + + def test_repeat_stage_reuses_the_first_mirror(self): mod = _load_learn(self.tmp) - mod.BASE = os.path.join(self.tmp, "does-not-exist") - # Must not raise even though the target directory is missing. - mod._append_episodic_mirror("deadbeef", "claim", "2026-01-01T00:00:00+00:00") + cid, path = mod.stage(CLAIM, CONDITIONS) + first = json.loads(Path(path).read_text()) + before = _episodic(self.tmp) + self.assertEqual(len(before), 1) + cid2, path2 = mod.stage(CLAIM, CONDITIONS) + second = json.loads(Path(path2).read_text()) + self.assertEqual(cid2, cid) + self.assertEqual(second["evidence_ids"], first["evidence_ids"]) + self.assertEqual(second["evidence_ids"], [before[0]["timestamp"]]) + self.assertEqual(_episodic(self.tmp), before) + self.assertEqual(_names(mod.CANDIDATES, ".tmp"), []) + + def test_planted_newer_temp_replaces_older_json_without_a_third_mirror(self): + mod = _load_learn(self.tmp) + cid, path = mod.stage(CLAIM, CONDITIONS) + published = json.loads(Path(path).read_text()) + later = "2099-01-01T00:00:00+00:00" + newer = dict(published) + newer["evidence_ids"] = [later] + newer["staged_at"] = later + Path(os.path.join(mod.CANDIDATES, f".{cid}.later.tmp")).write_text( + json.dumps(newer)) + episodic_path = os.path.join( + self.tmp, "memory", "episodic", "AGENT_LEARNINGS.jsonl") + with open(episodic_path, "a", encoding="utf-8") as stream: + stream.write(json.dumps({ + "timestamp": later, + "action": f"manual-stage:{cid}", + "result": "success", + }) + "\n") + mod.stage(CLAIM, CONDITIONS) + self.assertEqual(json.loads(Path(path).read_text())["evidence_ids"], [later]) + self.assertEqual(_names(mod.CANDIDATES, ".tmp"), []) + self.assertEqual(len(_episodic(self.tmp)), 2) + + def test_stale_temp_is_removed_when_published_json_is_newer(self): + mod = _load_learn(self.tmp) + cid, path = mod.stage(CLAIM, CONDITIONS) + published = json.loads(Path(path).read_text()) + stale_ts = "2000-01-01T00:00:00+00:00" + stale = dict(published) + stale["evidence_ids"] = [stale_ts] + stale["staged_at"] = stale_ts + temp_path = os.path.join(mod.CANDIDATES, f".{cid}.stale.tmp") + Path(temp_path).write_text(json.dumps(stale)) + episodic_path = os.path.join( + self.tmp, "memory", "episodic", "AGENT_LEARNINGS.jsonl") + with open(episodic_path, "a", encoding="utf-8") as stream: + stream.write(json.dumps({ + "timestamp": stale_ts, + "action": f"manual-stage:{cid}", + "result": "success", + }) + "\n") + before = _episodic(self.tmp) + mod.stage(CLAIM, CONDITIONS) + self.assertEqual(json.loads(Path(path).read_text())["evidence_ids"], + published["evidence_ids"]) + self.assertEqual(_names(mod.CANDIDATES, ".tmp"), []) + self.assertEqual(_episodic(self.tmp), before) + + def test_evidence_substring_does_not_resume_temp(self): + mod = _load_learn(self.tmp) + cid = mod.pattern_id(CLAIM, CONDITIONS) + buried = "2026-01-01T00:00:00+00:00" + temp_path = os.path.join(mod.CANDIDATES, f".{cid}.buried.tmp") + Path(temp_path).write_text(json.dumps({ + "id": cid, + "claim": CLAIM, + "evidence_ids": [buried], + })) + episodic_path = os.path.join( + self.tmp, "memory", "episodic", "AGENT_LEARNINGS.jsonl") + with open(episodic_path, "a", encoding="utf-8") as stream: + stream.write(json.dumps({ + "timestamp": "1999-01-01T00:00:00+00:00", + "detail": f"see {buried} in prose only", + }) + "\n") + cid_out, path = mod.stage(CLAIM, CONDITIONS) + self.assertEqual(cid_out, cid) + published = json.loads(Path(path).read_text()) + self.assertNotEqual(published["evidence_ids"], [buried]) + self.assertEqual(_names(mod.CANDIDATES, ".tmp"), []) + mirrors = [ + row for row in _episodic(self.tmp) + if row.get("action") == f"manual-stage:{cid}" + ] + self.assertEqual(len(mirrors), 1) + self.assertEqual(mirrors[0]["timestamp"], published["evidence_ids"][0]) + + def test_corrupt_temp_is_deleted_then_fresh_stage_runs(self): + mod = _load_learn(self.tmp) + cid = mod.pattern_id(CLAIM, CONDITIONS) + temp_path = os.path.join(mod.CANDIDATES, f".{cid}.corrupt.tmp") + Path(temp_path).write_text("{not json") + mod.stage(CLAIM, CONDITIONS) + self.assertEqual(_names(mod.CANDIDATES, ".json"), [f"{cid}.json"]) + self.assertEqual(_names(mod.CANDIDATES, ".tmp"), []) + self.assertEqual(len(_episodic(self.tmp)), 1) + + def test_mirror_failure_before_temp_leaves_no_candidate(self): + mod = _load_learn(self.tmp) + + def _mirror(*_a, **_k): + raise OSError("forced mirror-write failure") + + mod._append_episodic_mirror = _mirror + with self.assertRaises(OSError) as caught: + mod.stage(CLAIM, CONDITIONS) + self.assertIn("forced mirror-write failure", str(caught.exception)) + self.assertEqual(_names(mod.CANDIDATES, ".json"), []) + self.assertEqual(_names(mod.CANDIDATES, ".tmp"), []) + self.assertEqual(_episodic(self.tmp), []) + + def test_keyboard_interrupt_from_mirror_is_not_an_oserror(self): + mod = _load_learn(self.tmp) + + def _mirror(*_a, **_k): + raise KeyboardInterrupt + + mod._append_episodic_mirror = _mirror + with self.assertRaises(KeyboardInterrupt): + mod.stage(CLAIM, CONDITIONS) + self.assertEqual(_names(mod.CANDIDATES, ".tmp"), []) + self.assertEqual(_names(mod.CANDIDATES, ".json"), []) + + def test_ambiguous_leftovers_fail_closed(self): + mod = _load_learn(self.tmp) + cid = mod.pattern_id(CLAIM, CONDITIONS) + for suffix in ("a", "b"): + Path(os.path.join(mod.CANDIDATES, f".{cid}.{suffix}.tmp")).write_text("{}") + with self.assertRaises(OSError) as caught: + mod.stage(CLAIM, CONDITIONS) + message = str(caught.exception) + self.assertIn("ambiguous leftover temps", message) + self.assertIn(f".{cid}.a.tmp", message) + self.assertIn(f".{cid}.b.tmp", message) + self.assertEqual(_names(mod.CANDIDATES, ".json"), []) + self.assertEqual(len(_names(mod.CANDIDATES, ".tmp")), 2) + self.assertEqual(_episodic(self.tmp), []) + + def test_fsync_failure_after_mirror_line_retries_without_a_second_row(self): + mod = _load_learn(self.tmp) + real_append = mod._append_episodic_mirror + + def _append(*args, **kwargs): + real_fsync = mod.os.fsync + + def _boom(_fd): + raise OSError("forced fsync failure") + + mod.os.fsync = _boom + try: + return real_append(*args, **kwargs) + finally: + mod.os.fsync = real_fsync + + mod._append_episodic_mirror = _append + with self.assertRaises(OSError) as caught: + mod.stage(CLAIM, CONDITIONS) + self.assertIn("forced fsync failure", str(caught.exception)) + self.assertEqual(_names(mod.CANDIDATES, ".json"), []) + self.assertEqual(_names(mod.CANDIDATES, ".tmp"), []) + self.assertEqual(len(_episodic(self.tmp)), 1) + mod._append_episodic_mirror = real_append + cid, path = mod.stage(CLAIM, CONDITIONS) + self.assertEqual(_names(mod.CANDIDATES, ".json"), [f"{cid}.json"]) + self.assertEqual(_names(mod.CANDIDATES, ".tmp"), []) + self.assertEqual(len(_episodic(self.tmp)), 1) + candidate = json.loads(Path(path).read_text()) + self.assertEqual( + candidate["evidence_ids"], + [_episodic(self.tmp)[0]["timestamp"]], + ) + + def test_published_read_error_does_not_clobber_json(self): + mod = _load_learn(self.tmp) + cid, path = mod.stage(CLAIM, CONDITIONS) + original = Path(path).read_text() + before = _episodic(self.tmp) + temp_path = os.path.join(mod.CANDIDATES, f".{cid}.newer.tmp") + Path(temp_path).write_text(original) + real_load = mod._load_json_object + + def _load(candidate_path): + if os.path.abspath(candidate_path) == os.path.abspath(path): + raise OSError("forced published read failure") + return real_load(candidate_path) + + mod._load_json_object = _load + with self.assertRaises(OSError) as caught: + mod.stage(CLAIM, CONDITIONS) + self.assertIn("forced published read failure", str(caught.exception)) + self.assertEqual(Path(path).read_text(), original) + self.assertTrue(os.path.isfile(temp_path)) + self.assertEqual(_episodic(self.tmp), before) + + def test_corrupt_temp_beside_published_json_does_not_remirror(self): + mod = _load_learn(self.tmp) + cid, path = mod.stage(CLAIM, CONDITIONS) + published = Path(path).read_text() + before = _episodic(self.tmp) + Path(os.path.join(mod.CANDIDATES, f".{cid}.corrupt.tmp")).write_text("{not json") + mod.stage(CLAIM, CONDITIONS) + self.assertEqual(Path(path).read_text(), published) + self.assertEqual(_names(mod.CANDIDATES, ".tmp"), []) + self.assertEqual(_episodic(self.tmp), before) + + def test_keyboard_interrupt_after_mirror_write_retries_without_second_mirror(self): + mod = _load_learn(self.tmp) + real_append = mod._append_episodic_mirror + + def _append(*args, **kwargs): + real_append(*args, **kwargs) + raise KeyboardInterrupt + + mod._append_episodic_mirror = _append + with self.assertRaises(KeyboardInterrupt): + mod.stage(CLAIM, CONDITIONS) + self.assertEqual(_names(mod.CANDIDATES, ".tmp"), []) + self.assertEqual(len(_episodic(self.tmp)), 1) + mod._append_episodic_mirror = real_append + cid, resumed = mod.stage(CLAIM, CONDITIONS) + self.assertEqual(_names(mod.CANDIDATES, ".json"), [f"{cid}.json"]) + self.assertEqual(_names(mod.CANDIDATES, ".tmp"), []) + self.assertEqual(len(_episodic(self.tmp)), 1) + candidate = json.loads(Path(resumed).read_text()) + self.assertEqual( + candidate["evidence_ids"], + [_episodic(self.tmp)[0]["timestamp"]], + ) + + def test_append_jsonl_once_reuses_the_first_row(self): + import hooks._episodic_io as episodic_io + + path = os.path.join(self.tmp, "once.jsonl") + first = {"timestamp": "t1", "action": "manual-stage:abc"} + second = {"timestamp": "t2", "action": "manual-stage:abc"} + self.assertEqual( + episodic_io.append_jsonl_once( + path, first, match_action="manual-stage:abc")["timestamp"], + "t1", + ) + self.assertEqual( + episodic_io.append_jsonl_once( + path, second, match_action="manual-stage:abc")["timestamp"], + "t1", + ) + episodic_io.append_jsonl(path, second) + rows = [ + json.loads(line) + for line in Path(path).read_text().splitlines() + if line.strip() + ] + self.assertEqual([row["timestamp"] for row in rows], ["t1", "t2"]) + + def test_identical_resumable_temps_publish_once(self): + mod = _load_learn(self.tmp) + cid, path = mod.stage(CLAIM, CONDITIONS) + published = Path(path).read_text() + os.remove(path) + for suffix in ("a", "b"): + Path(os.path.join(mod.CANDIDATES, f".{cid}.{suffix}.tmp")).write_text( + published) + before = _episodic(self.tmp) + mod.stage(CLAIM, CONDITIONS) + self.assertEqual(Path(path).read_text(), published) + self.assertEqual(_names(mod.CANDIDATES, ".tmp"), []) + self.assertEqual(_episodic(self.tmp), before) + + def test_evidence_landed_reads_under_exclusive_lock(self): + mod = _load_learn(self.tmp) + import hooks._episodic_io as episodic_io + + path = os.path.join( + self.tmp, "memory", "episodic", "AGENT_LEARNINGS.jsonl" + ) + row = {"timestamp": "2026-09-26T00:00:00+00:00", "action": "test"} + Path(path).write_text(json.dumps(row) + "\n") + if not episodic_io._HAVE_FLOCK: + self.skipTest("fcntl flock not available") + + order = [] + real_flock = episodic_io.fcntl.flock + + def _record_flock(fd, operation): + if operation == episodic_io.fcntl.LOCK_UN: + order.append("unlock") + else: + order.append("lock") + return real_flock(fd, operation) + + episodic_io.fcntl.flock = _record_flock + try: + self.assertTrue(mod._evidence_landed(path, row["timestamp"])) + finally: + episodic_io.fcntl.flock = real_flock + + self.assertEqual(order, ["lock", "unlock"]) + + def test_evidence_read_error_keeps_resumable_temp(self): + mod = _load_learn(self.tmp) + cid = mod.pattern_id(CLAIM, CONDITIONS) + temp_path = os.path.join(mod.CANDIDATES, f".{cid}.pending.tmp") + timestamp = "2026-09-26T00:00:00+00:00" + Path(temp_path).write_text(json.dumps({ + "id": cid, + "claim": CLAIM, + "evidence_ids": [timestamp], + })) + real_check = mod.has_jsonl_timestamp + mod.has_jsonl_timestamp = lambda *_a, **_k: (_ for _ in ()).throw( + OSError("forced locked-read failure") + ) + try: + with self.assertRaises(OSError) as caught: + mod.stage(CLAIM, CONDITIONS) + finally: + mod.has_jsonl_timestamp = real_check + self.assertIn("forced locked-read failure", str(caught.exception)) + self.assertTrue(os.path.isfile(temp_path)) + self.assertEqual(_episodic(self.tmp), []) + + def test_missing_episodic_file_is_not_created(self): + import hooks._episodic_io as episodic_io + + path = os.path.join(self.tmp, "memory", "episodic", "absent.jsonl") + self.assertFalse(os.path.exists(path)) + self.assertFalse( + episodic_io.has_jsonl_timestamp(path, "2026-09-26T00:00:00+00:00") + ) + self.assertFalse(os.path.exists(path)) if __name__ == "__main__":