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:<job>:<uuid>` while their session_id is `cron_<job>_<ts>`, 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.
This commit is contained in:
@@ -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:<job>:<uuid>``, 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:
|
||||
|
||||
@@ -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:<job>:<uuid>`` 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:<job>:<uuid>``, 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()
|
||||
|
||||
@@ -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:<job>:<uuid>`` 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."""
|
||||
|
||||
Reference in New Issue
Block a user