From f7a422ee4a0df12f468b06ef8f08d5696247cd91 Mon Sep 17 00:00:00 2001 From: teknium1 <127238744+teknium1@users.noreply.github.com> Date: Fri, 18 Sep 2026 00:51:07 -0700 Subject: [PATCH] fix(file-state): close() releases file state for every task id the agent ran; drop the writer TTL knob Follow-up to the salvaged #114470 commit. `AIAgent.close()` hands `_close_task_resources()` the agent's session_id, but the file tools key `FileStateRegistry` by the per-turn task_id: cron runs use `cron::` while their session_id is `cron__`, and delegate children use `subagent-N-xxxx` against a fresh uuid session. `cleanup_vm(session_id)` therefore never reached `forget_task()` for the id that owns the read stamps and writer claims, so the purge added to `forget_task()` had no production caller for exactly the lifecycles the issue describes. `close()` now runs `clear_file_ops_cache()` for every id in `_process_owner_task_ids` (the set the turn context already maintains for process ownership) before dropping the session. Dropped from #114470: the `HERMES_FILE_STATE_WRITER_TTL` env knob and the time-based eviction in `check_stale()`. With the lifecycle end actually releasing the finished task's claims, a TTL only weakens the concurrent case the guard exists for (a live sibling's hour-old write is still a real conflict), and behavioural env vars are not a config surface. The module-level `forget_task()` wrapper and its `__all__` entry from the salvaged commit are kept. Tests trimmed to two invariants in the mirroring file: `forget_task()` purges the finished task's writer claims (a live sibling still fires), and `AIAgent.close()` releases the file state of every task id it ran even though it receives the session_id. --- agent/client_lifecycle.py | 12 +++++- tests/tools/test_file_state_registry.py | 54 ++++++++++++------------- tools/file_state.py | 30 ++++---------- 3 files changed, 45 insertions(+), 51 deletions(-) diff --git a/agent/client_lifecycle.py b/agent/client_lifecycle.py index 8d1b76337f..e2722dad25 100644 --- a/agent/client_lifecycle.py +++ b/agent/client_lifecycle.py @@ -123,7 +123,17 @@ class ClientLifecycleMixin: from tools.computer_use.tool import release_computer_use_session release_computer_use_session(task_id) - for step in (kill_processes, lambda: cleanup_vm(task_id), lambda: cleanup_browser(task_id), release_computer_use): + def forget_file_state() -> None: + # File tools key their read stamps / writer claims by the per-turn task_id (cron: + # ``cron::``, subagents: ``subagent-N-xxxx``), which differs from session_id; + # cleanup_vm(session_id) alone leaves a finished run looking like a live sibling (#114446). + from tools.file_tools import clear_file_ops_cache + for owner in getattr(self, "_process_owner_task_ids", ()): + if owner and owner != task_id: + clear_file_ops_cache(owner) + + for step in (kill_processes, lambda: cleanup_vm(task_id), lambda: cleanup_browser(task_id), + release_computer_use, forget_file_state): _quietly(step) def _client_log_context(self) -> str: diff --git a/tests/tools/test_file_state_registry.py b/tests/tools/test_file_state_registry.py index bf46723eb0..44c42f3c8f 100644 --- a/tests/tools/test_file_state_registry.py +++ b/tests/tools/test_file_state_registry.py @@ -22,6 +22,7 @@ import tempfile import threading import time import unittest +from unittest.mock import patch from tools import file_state from tools.file_tools import ( @@ -173,44 +174,43 @@ class FileStateRegistryUnitTests(unittest.TestCase): self.assertNotIn(task_id, rt._patch_failure_tracker) def test_forget_task_clears_last_writer_claims(self): - """Regression test for issue #114446: forget_task must prune _last_writer - entries owned by the ended task so sequential runs don't false-positive as - concurrent sibling conflicts.""" + """A finished task is not a concurrent sibling: forget_task must drop its writer + claims so the next run of the same job (fresh ``cron::`` id) can write + the same scratch path without a "modified by sibling subagent" refusal.""" p = self._mk() file_state.note_write("cron:JOB:run1", p) registry = file_state.get_registry() - self.assertIn(p, registry._last_writer) self.assertEqual(registry._last_writer[p][0], "cron:JOB:run1") - # Now task lifecycle ends file_state.forget_task("cron:JOB:run1") - # The writer claim must be gone self.assertNotIn(p, registry._last_writer) + self.assertIsNone(file_state.check_stale("cron:JOB:run2", p)) + # A sibling that has NOT ended still triggers the guard. + file_state.note_write("subagent-1-live", p) + self.assertIn("sibling subagent 'subagent-1-live'", file_state.check_stale("cron:JOB:run2", p)) - # Next sequential run touching the file should not trigger a sibling warning - warn = file_state.check_stale("cron:JOB:run2", p) - self.assertIsNone(warn) - - def test_clear_file_ops_cache_clears_last_writer_claims(self): - """Ensure file_tools.clear_file_ops_cache propagates forget_task to _last_writer.""" + def test_agent_close_forgets_every_task_id_it_ran(self): + """``AIAgent.close()`` receives the session_id, but file tools key the registry by + the per-turn task_id (cron ``cron::``, subagent ``subagent-N-xxxx``). + close() must release the file state of every task id the agent ran.""" p = self._mk() - file_state.note_write("worker-1", p) - clear_file_ops_cache("worker-1") - warn = file_state.check_stale("worker-2", p) - self.assertIsNone(warn) + file_state.record_read("cron:JOB:run1", p) + file_state.note_write("cron:JOB:run1", p) + with patch("run_agent.AIAgent.__init__", return_value=None): + from run_agent import AIAgent + agent = AIAgent.__new__(AIAgent) + agent.session_id = "cron_JOB_20260918_060000" + agent._process_owner_task_ids = {"cron:JOB:run1"} + agent._active_children = [] + agent._active_children_lock = threading.Lock() + agent.client = None + with patch("run_agent.cleanup_vm"), patch("run_agent.cleanup_browser"), \ + patch("tools.computer_use.tool.release_computer_use_session"): + agent.close() - def test_last_writer_ttl_expiration(self): - """Entries older than TTL must not report stale conflicts for long-lived processes.""" - p = self._mk() - # Simulate a write from 2 hours ago - old_time = time.time() - 7200 - with file_state.get_registry()._state_lock: - file_state.get_registry()._last_writer[p] = ("old-task", old_time) - - warn = file_state.check_stale("new-task", p) - self.assertIsNone(warn) - self.assertNotIn(p, file_state.get_registry()._last_writer) + self.assertEqual(file_state.known_reads("cron:JOB:run1"), []) + self.assertIsNone(file_state.check_stale("cron:JOB:run2", p)) def test_kill_switch_env_var(self): p = self._mk() diff --git a/tools/file_state.py b/tools/file_state.py index f4576e82cb..0c5e57f092 100644 --- a/tools/file_state.py +++ b/tools/file_state.py @@ -38,17 +38,6 @@ def guard_disabled() -> bool: return _disabled() -def _writer_ttl_seconds() -> float: - # TTL for _last_writer entries to bound concurrent conflict detection window - raw = os.environ.get("HERMES_FILE_STATE_WRITER_TTL") - if raw: - try: - return float(raw) - except ValueError: - pass - return 3600.0 # default: 1 hour - - def _mtime_or_none(resolved: str) -> Optional[float]: try: return os.path.getmtime(resolved) @@ -143,11 +132,6 @@ class FileStateRegistry: with self._state_lock: stamp = self._reads.get(task_id, {}).get(resolved) last_writer = self._last_writer.get(resolved) - if last_writer is not None: - ttl = _writer_ttl_seconds() - if ttl > 0 and (time.time() - last_writer[1]) > ttl: - self._last_writer.pop(resolved, None) - last_writer = None if stamp is None and last_writer is None: # net-new file / first touch return None @@ -213,15 +197,15 @@ class FileStateRegistry: return list(self._reads.get(task_id, {}).keys()) def forget_task(self, task_id: str) -> None: - """Release read stamps and writer claims owned by a task after its lifecycle ends.""" + """Release read stamps and writer claims owned by a task after its lifecycle ends. + + A finished task is not a concurrent sibling: leaving its writer claims behind makes + the next run of the same job (a fresh ``cron::`` id) refuse to write the + same scratch path as "modified by sibling subagent" hours after the writer exited.""" with self._state_lock: self._reads.pop(task_id, None) - stale_paths = [ - p for p, (writer_tid, _) in self._last_writer.items() - if writer_tid == task_id - ] - for p in stale_paths: - self._last_writer.pop(p, None) + for p in [p for p, (writer_tid, _ts) in self._last_writer.items() if writer_tid == task_id]: + del self._last_writer[p] def clear(self) -> None: """Reset all state. Intended for tests only."""