-
Notifications
You must be signed in to change notification settings - Fork 281
fix(learn): fail-closed episodic mirror before candidate publish #71
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: master
Are you sure you want to change the base?
Changes from all commits
30e03cc
0a9f48e
7a477dd
949aa74
be684ff
409e741
4b56ff0
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change | ||||
|---|---|---|---|---|---|---|
|
|
@@ -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(): | ||||||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
Suggested change
Prompt To Fix With AIThis is a comment left during a code review.
Path: .agent/harness/hooks/_episodic_io.py
Line: 59
Comment:
**Malformed bytes block staging.** If an older episodic entry contains a non-UTF-8 byte, decoding the whole log raises `UnicodeDecodeError` before the corrupt-line handler runs. Manual staging now scans that log, so an unrelated malformed entry stops staging with an uncaught traceback and no candidate is published.
```suggestion
for line in raw.decode("utf-8", errors="replace").splitlines():
```
---
For each issue above, determine whether it is valid and should be fixed. If so, fix it directly.
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. @greptileai fixed this in 91c3dff There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Confirmed fixed in |
||||||
| 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) | ||||||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -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() |
Uh oh!
There was an error while loading. Please reload this page.