fix(cron): isolate per-execution working directories

This commit is contained in:
Andrew Bagrin
2026-08-27 14:32:50 -04:00
committed by Teknium
parent 89bcad4d0c
commit b7c59bda54
14 changed files with 201 additions and 663 deletions
+30 -261
View File
@@ -395,20 +395,6 @@ def _summarize_cron_failure_for_delivery(job: dict, error: str | None) -> str:
"quiet. Full details saved in cron output."
)
# Sibling scheduler-side timeout (#79768): the TERMINAL_CWD lock-wait
# abort also phrases itself with "Timed out ..." and would fall through
# to the generic provider-timeout branch below. Like the inactivity
# watchdog above, it is entirely scheduler-internal — no provider or
# fallback chain involved — so classify it before the generic match.
if "terminal_cwd" in lower and ("lock" in lower or "timed out" in lower):
return (
f"⚠️ Cron '{job_name}' failed: could not acquire the scheduler's "
"working-directory lock — another cron job (a workdir writer or "
"long-running readers) held it too long. Not a provider or "
"fallback-chain issue; stagger the holder's schedule or remove "
"its workdir. Full details saved in cron output."
)
if provider_reachable and (
"readtimeout" in lower or "timed out" in lower or "timeout" in lower
):
@@ -1390,120 +1376,6 @@ def _consume_interrupted_flag(job_id: str, token: Optional[object] = None) -> bo
return hit
# Sequential (env-mutating) cron jobs — workdir jobs that touch
# process-global runtime state — must run one at a time, but must NOT block the
# ticker thread. A persistent single-thread executor preserves ordering across
# ticks while keeping dispatch fire-and-forget, the same as the parallel pool.
_sequential_pool: Optional[concurrent.futures.ThreadPoolExecutor] = None
class _ReadWriteLock:
"""Writer-preferring readers-writer lock.
Guards the process-global ``os.environ["TERMINAL_CWD"]`` override that a
workdir cron job applies for the whole of its agent run. Workdir jobs are
writers: they mutate the shared env and need exclusive access. Workdir-less
jobs are readers: they only observe ``TERMINAL_CWD`` (indirectly, via the
terminal / file / code-exec tools), so any number of them may run
concurrently with each other, but none may run alongside a writer — that is
exactly what stops a workdir-less job from picking up another job's workdir
override and running its commands in the wrong directory.
Writer preference bounds the wait for a workdir job (dispatched on the
single-thread sequential pool) so a stream of workdir-less readers cannot
starve it.
"""
def __init__(self) -> None:
self._cond = threading.Condition(threading.Lock())
self._readers = 0
self._writer_active = False
self._writers_waiting = 0
def acquire_read(self, timeout: float | None = None) -> bool:
"""Acquire a read lock.
Returns ``True`` if the lock was acquired, ``False`` on timeout.
A timed-out caller proceeds without the lock (degraded mode) —
see the call-site in ``run_job`` for the logging / trade-off.
"""
deadline = (
time.monotonic() + timeout if timeout is not None else None
)
with self._cond:
while self._writer_active or self._writers_waiting > 0:
if deadline is not None:
remaining = deadline - time.monotonic()
if remaining <= 0:
self._cond.notify_all()
return False
self._cond.wait(timeout=remaining)
else:
self._cond.wait()
self._readers += 1
return True
def release_read(self) -> None:
with self._cond:
self._readers -= 1
if self._readers == 0:
self._cond.notify_all()
def acquire_write(self, timeout: float | None = None) -> bool:
"""Acquire a write lock.
Returns ``True`` if the lock was acquired, ``False`` on timeout.
A timed-out caller proceeds without the lock (degraded mode).
"""
deadline = (
time.monotonic() + timeout if timeout is not None else None
)
with self._cond:
self._writers_waiting += 1
try:
while self._writer_active or self._readers > 0:
if deadline is not None:
remaining = deadline - time.monotonic()
if remaining <= 0:
self._cond.notify_all()
return False
self._cond.wait(timeout=remaining)
else:
self._cond.wait()
finally:
self._writers_waiting -= 1
self._writer_active = True
return True
def release_write(self) -> None:
with self._cond:
self._writer_active = False
self._cond.notify_all()
# Serializes the per-job TERMINAL_CWD override against every other concurrently
# running cron job. See _ReadWriteLock and run_job for the usage contract.
_terminal_cwd_lock = _ReadWriteLock()
# Ceiling on how long a cron job waits for the TERMINAL_CWD lock before
# FAILING (fail-closed, #79768). Derived from the cron inactivity limit
# (HERMES_CRON_TIMEOUT, default 600s): a wedged lock holder stops touching
# its activity clock, so the inactivity monitor usually reaps it and the
# lock is released within roughly that limit. The bound is measured from
# the WAITER's arrival, so a holder that wedges late (or hangs in pre-agent
# setup before the monitor arms) can still outlive it — waiters then fail
# loudly rather than run: proceeding without the lock lets the holder's
# process-global TERMINAL_CWD override leak into this job's shell/file/
# code-exec commands (wrong-directory execution — the exact corruption
# _ReadWriteLock exists to prevent, see
# test_reader_never_observes_writer_override). A healthy long-running
# workdir job past the bound also fails its waiters loudly rather than
# corrupting them silently; the failure names the holder pattern so the fix
# (stagger schedules / drop the workdir) is actionable.
_CWD_LOCK_TIMEOUT_FLOOR_SECONDS = 120.0
_CWD_LOCK_TIMEOUT_MARGIN_SECONDS = 60.0
def _cron_inactivity_seconds() -> float:
"""Parse HERMES_CRON_TIMEOUT (seconds). 0 = unlimited; bad input = 600.
@@ -1522,17 +1394,6 @@ def _cron_inactivity_seconds() -> float:
return 600.0
def _cwd_lock_timeout_seconds() -> float:
"""Bound for the TERMINAL_CWD lock wait: inactivity limit + margin."""
inactivity = _cron_inactivity_seconds()
if inactivity <= 0: # 0 = unlimited job runtime; keep the wait bounded.
inactivity = 600.0
return (
max(inactivity, _CWD_LOCK_TIMEOUT_FLOOR_SECONDS)
+ _CWD_LOCK_TIMEOUT_MARGIN_SECONDS
)
def _get_parallel_pool(max_workers: Optional[int]) -> concurrent.futures.ThreadPoolExecutor:
"""Return (or create) the persistent parallel pool."""
global _parallel_pool, _parallel_pool_max_workers
@@ -1547,33 +1408,14 @@ def _get_parallel_pool(max_workers: Optional[int]) -> concurrent.futures.ThreadP
return _parallel_pool
def _get_sequential_pool() -> concurrent.futures.ThreadPoolExecutor:
"""Return (or create) the persistent single-thread sequential pool.
A single worker guarantees env-mutating jobs never overlap, even
across ticks: a job queued by a newer tick waits for the previous tick's
sequential jobs to finish rather than corrupting their os.environ
state.
"""
global _sequential_pool
if _sequential_pool is None:
_sequential_pool = concurrent.futures.ThreadPoolExecutor(
max_workers=1,
thread_name_prefix="cron-seq",
)
return _sequential_pool
def _shutdown_parallel_pool() -> None:
"""Shut down the persistent pools on process exit."""
global _parallel_pool, _parallel_pool_max_workers, _sequential_pool
"""Shut down the persistent pool on process exit."""
global _parallel_pool, _parallel_pool_max_workers
if _parallel_pool is not None:
_parallel_pool.shutdown(wait=True, cancel_futures=False)
_parallel_pool = None
_parallel_pool_max_workers = None
if _sequential_pool is not None:
_sequential_pool.shutdown(wait=True, cancel_futures=False)
_sequential_pool = None
atexit.register(_shutdown_parallel_pool)
@@ -5613,6 +5455,7 @@ def run_job(
defer_agent_teardown: Optional[list] = None,
extra_prompt: Optional[str] = None,
cancel_event: Optional[_CancelEventLike] = None,
execution_id: Optional[str] = None,
) -> tuple[bool, str, str, Optional[str]]:
"""
Execute a single cron job.
@@ -5977,68 +5820,25 @@ def run_job(
for _var_name in _cron_delivery_vars:
_VAR_MAP[_var_name].set("")
# Per-job working directory — _SESSION_CWD was already set via
# set_session_vars(cwd=...) above. Here we only handle the
# process-global TERMINAL_CWD env var, which is serialized by
# _terminal_cwd_lock to avoid leaking into concurrent jobs.
#
# os.environ["TERMINAL_CWD"] is process-global, so this override is
# serialized by _terminal_cwd_lock (acquired just below): a workdir job
# holds it as a writer for its whole run, excluding every other job, while
# workdir-less jobs hold it as readers and stay parallel with each other.
# The sequential pool only keeps workdir jobs from overlapping EACH OTHER;
# the lock is what additionally keeps a concurrently-firing workdir-less
# parallel-pool job from observing this override and running its shell /
# file / code-exec commands in the wrong directory. For workdir-less jobs
# we leave TERMINAL_CWD untouched — preserves the original behaviour
# (skip_context_files=True, tools use whatever cwd the scheduler has).
#
# The critical path (resolve_context_cwd / build_context_files_prompt)
# checks _SESSION_CWD first, so gateway sessions with no override see
# their own cwd, not the cron's workdir (#69396).
# Snapshot the current env value BEFORE acquiring the lock so the finally
# below can always restore it, even if an exception fires before we set the
# override inside the try. This read can't leak the lock (it precedes the
# acquire) and is a no-op for workdir-less jobs (they never mutate the env).
_prior_terminal_cwd = os.environ.get("TERMINAL_CWD", "_UNSET_")
_holds_cwd_write = _job_workdir is not None
_cwd_lock_timeout = _cwd_lock_timeout_seconds()
_cwd_lock_acquired = True
if _holds_cwd_write:
if not _terminal_cwd_lock.acquire_write(timeout=_cwd_lock_timeout):
_cwd_lock_acquired = False
else:
if not _terminal_cwd_lock.acquire_read(timeout=_cwd_lock_timeout):
_cwd_lock_acquired = False
# Everything after the acquire MUST live inside this try, so the finally
# below always releases the lock even if the env override or any later
# statement raises. A leaked writer would deadlock the whole scheduler
# (every future job blocks on acquire_*); a leaked reader blocks all
# future writers. Acquire itself can't leak (it either blocks or returns).
# Tool calls are keyed by a per-run task id. Bind the cron workdir to that
# identity instead of mutating process-global TERMINAL_CWD. The session CWD
# record is the tool-layer authority for terminal/file/code-exec/delegation,
# while the _SESSION_CWD ContextVar above remains the prompt/context-file
# authority. Both are isolated across concurrent cron runs.
_cron_task_id = (
f"cron:{job_id}:"
f"{execution_id or job.get('execution_id') or uuid.uuid4().hex}"
)
from tools.terminal_tool import clear_session_cwd as _clear_tool_session_cwd
from tools.terminal_tool import record_session_cwd as _record_tool_session_cwd
if _job_workdir:
_record_tool_session_cwd(_cron_task_id, _job_workdir)
_cron_session_var = _VAR_MAP["HERMES_CRON_SESSION"]
_cron_session_token = None
_non_dispatcher_token = None
_session_db = None
try:
if not _cwd_lock_acquired:
# Fail closed (#79768): running without the lock would let a
# concurrent workdir job's process-global TERMINAL_CWD override
# leak into this job's shell/file/code-exec commands — silent
# wrong-directory execution, the exact corruption the lock
# exists to prevent. A loud failure is recoverable (next tick /
# manual rerun); a job that ran in the wrong directory is not.
raise TimeoutError(
f"Timed out waiting for the TERMINAL_CWD "
f"{'write' if _holds_cwd_write else 'read'} lock after "
f"{_cwd_lock_timeout:.0f}s — another cron job (a workdir "
f"writer, or long-running readers) has held it for longer "
f"than the cron inactivity limit. If a workdir job is the "
f"holder, stagger its schedule or remove its workdir to "
f"unblock this job (#79768)."
)
# Scope cron approval policy to this job. Keep the token so the finally
# restores the pre-job state instead of pinning an explicit empty value,
# which would suppress the legacy os.environ fallback used by standalone
@@ -6065,8 +5865,7 @@ def run_job(
# at the run_conversation hop carries this into the agent thread.
_non_dispatcher_token = enter_non_dispatcher_owned_context()
if _job_workdir:
os.environ["TERMINAL_CWD"] = _job_workdir
logger.info("Job '%s': using workdir %s", job_id, _job_workdir)
logger.info("Job '%s': using task-scoped workdir %s", job_id, _job_workdir)
# Re-read .env and config.yaml fresh every run so provider/key
# changes take effect without a gateway restart. Route through
@@ -6696,7 +6495,12 @@ def run_job(
# Tag this fire and time the run_conversation call for the usage_audit.jsonl entry.
_audit_fire_id = uuid.uuid4().hex
_audit_t_start = time.monotonic()
_cron_future = _cron_pool.submit(_cron_context.run, agent.run_conversation, prompt)
_cron_future = _cron_pool.submit(
_cron_context.run,
agent.run_conversation,
prompt,
task_id=_cron_task_id,
)
_inactivity_timeout = False
try:
if _cron_inactivity_limit is None:
@@ -6948,23 +6752,7 @@ def run_job(
return False, output, "", error_msg
finally:
# Restore TERMINAL_CWD to whatever it was before this job ran. We
# only ever mutate it when the job has a workdir AND actually held
# the write lock — a fail-closed timeout raised before the env-set,
# so restoring there would replay a pre-wait snapshot over the
# ACTIVE holder's live override.
if _job_workdir and _cwd_lock_acquired:
if _prior_terminal_cwd == "_UNSET_":
os.environ.pop("TERMINAL_CWD", None)
else:
os.environ["TERMINAL_CWD"] = _prior_terminal_cwd
# Release the cwd lock now that the env is restored, so a waiting
# workdir job (or queued reader) can proceed without seeing the override.
if _cwd_lock_acquired:
if _holds_cwd_write:
_terminal_cwd_lock.release_write()
else:
_terminal_cwd_lock.release_read()
_clear_tool_session_cwd(_cron_task_id)
# Clean up ContextVar session/delivery state for this job.
# clear_session_vars also clears _SESSION_CWD internally, so no
# separate clear_session_cwd() call is needed.
@@ -7409,6 +7197,7 @@ def _run_one_job_body(
job,
defer_agent_teardown=_deferred_agents,
extra_prompt=extra_prompt,
execution_id=execution_id,
)
else:
success, output, final_response, error = run_job(
@@ -7416,6 +7205,7 @@ def _run_one_job_body(
defer_agent_teardown=_deferred_agents,
extra_prompt=extra_prompt,
cancel_event=fire_claim_lost,
execution_id=execution_id,
)
except BaseException:
# run_job's finally still hands back the agent when it raises; tear
@@ -8247,14 +8037,8 @@ def tick(
verbose=verbose,
)
# Partition due jobs: those with a per-job workdir mutate
# os.environ["TERMINAL_CWD"] inside run_job, which is process-global, so
# they queue on the single-thread sequential pool to run one at a time.
# That alone only keeps workdir jobs from overlapping EACH OTHER;
# run_job's _terminal_cwd_lock is what additionally stops a concurrently
# firing workdir-less parallel-pool job from observing the override.
sequential_jobs = [j for j in due_jobs if (j.get("workdir") or "").strip()]
parallel_jobs = [j for j in due_jobs if not (j.get("workdir") or "").strip()]
# Workdir is task-scoped, so every job uses the normal parallel lane.
parallel_jobs = due_jobs
_results: list = []
_all_futures: list = []
@@ -8370,21 +8154,6 @@ def tick(
_running_futures[job_id] = fut
return fut
# Sequential pass for env-mutating (workdir) jobs.
# Queued to a persistent single-thread pool so they run one at a time
# WITHOUT blocking the ticker thread — a long workdir job no
# longer starves the rest of the schedule (same fix as the parallel
# pass, just serialized). The in-flight guard prevents a still-running
# job from being re-queued on the next tick.
if sequential_jobs:
seq_pool = _get_sequential_pool()
for job in sequential_jobs:
fut = _submit_with_guard(job, seq_pool)
if fut is None:
continue
_all_futures.append(fut)
if not sync:
_results.append(True) # optimistically counted
# Parallel pass — persistent pool, non-blocking dispatch.
# Jobs that are already running (from a previous tick) are skipped.
@@ -94,22 +94,3 @@ def test_rate_limit_classification_still_takes_priority_over_inactivity_text(mon
msg = _summarize_cron_failure_for_delivery(job, error)
assert "weekly usage limit" in msg
assert "No fallback chain configured" in msg
def test_terminal_cwd_lock_timeout_is_not_reported_as_provider_timeout():
"""Sibling scheduler-internal timeout (#79768): the TERMINAL_CWD lock-wait
abort says "Timed out ..." and must not fall through to the generic
provider-timeout branch."""
job = {"name": "Workdir Job", "id": "abc123def456"}
error = (
"TimeoutError: Timed out waiting for the TERMINAL_CWD write lock "
"after 600s — another cron job (a workdir writer, or long-running "
"readers) has held it for longer than the cron inactivity limit. "
"If a workdir job is the holder, stagger its schedule or remove its "
"workdir to unblock this job (#79768)."
)
msg = _summarize_cron_failure_for_delivery(job, error)
assert "provider timeout" not in msg
assert "fallback chain" not in msg.lower()
assert "working-directory lock" in msg
assert "Workdir Job" in msg
+62 -44
View File
@@ -134,74 +134,56 @@ class TestCronjobToolWorkdir:
# ---------------------------------------------------------------------------
# scheduler.tick(): workdir partition
# scheduler.tick(): workdir jobs use the parallel execution lane
# ---------------------------------------------------------------------------
class TestTickWorkdirPartition:
"""
tick() must run workdir jobs sequentially (outside the ThreadPoolExecutor)
because run_job mutates os.environ["TERMINAL_CWD"], which is process-global.
We verify the partition without booting the real scheduler by patching the
pieces tick() calls.
"""
"""Workdir is per execution, so it must not force a global serial lane."""
def test_workdir_jobs_run_sequentially(self, tmp_path, monkeypatch):
def test_workdir_jobs_overlap_on_parallel_pool(self, tmp_path, monkeypatch):
import cron.scheduler as sched
import threading
# Two workdir jobs (both sequential) + one parallel job.
workdir_a = {"id": "a", "name": "A", "workdir": str(tmp_path)}
workdir_b = {"id": "b", "name": "B", "workdir": str(tmp_path)}
parallel_job = {"id": "c", "name": "C", "workdir": None}
monkeypatch.setattr(sched, "get_due_jobs", lambda: [workdir_a, workdir_b, parallel_job])
workdir_a = tmp_path / "a"
workdir_b = tmp_path / "b"
workdir_a.mkdir()
workdir_b.mkdir()
jobs = [
{"id": "a", "name": "A", "workdir": str(workdir_a)},
{"id": "b", "name": "B", "workdir": str(workdir_b)},
]
monkeypatch.setattr(sched, "get_due_jobs", lambda: jobs)
monkeypatch.setattr(sched, "claim_job_for_fire", lambda *_a, **_kw: True)
# Record call order / thread context.
import threading
barrier = threading.Barrier(2, timeout=5)
calls: list[tuple[str, str]] = []
order_lock = threading.Lock()
calls_lock = threading.Lock()
def fake_run_job(job, *, defer_agent_teardown=None, **_kw):
# Return a minimal tuple matching run_job's signature.
with order_lock:
with calls_lock:
calls.append((job["id"], threading.current_thread().name))
barrier.wait()
return True, "output", "response", None
monkeypatch.setattr(sched, "run_job", fake_run_job)
monkeypatch.setattr(sched, "save_job_output", lambda _jid, _o: None)
monkeypatch.setattr(sched, "mark_job_run", lambda *_a, **_kw: None)
monkeypatch.setattr(
sched, "_deliver_result", lambda *_a, **_kw: None
)
monkeypatch.setattr(sched, "_deliver_result", lambda *_a, **_kw: None)
n = sched.tick(verbose=False)
assert n == 3
ids = [c[0] for c in calls]
# Sequential workdir jobs preserve submission order relative to each
# other (single-thread pool).
assert ids.index("a") < ids.index("b")
# Workdir jobs run on the persistent single-thread cron-seq pool —
# NOT the main thread — so a long workdir job never blocks the ticker.
main_thread_name = threading.current_thread().name
for jid in ("a", "b"):
workdir_thread_name = next(t for j, t in calls if j == jid)
assert workdir_thread_name != main_thread_name
assert workdir_thread_name.startswith("cron-seq"), workdir_thread_name
par_thread_name = next(t for j, t in calls if j == "c")
assert par_thread_name.startswith("cron-parallel"), par_thread_name
assert sched.tick(verbose=False, sync=True) == 2
assert {job_id for job_id, _thread in calls} == {"a", "b"}
assert all(thread.startswith("cron-parallel") for _job, thread in calls)
# ---------------------------------------------------------------------------
# scheduler.run_job: TERMINAL_CWD + skip_context_files wiring
# scheduler.run_job: per-task cwd + skip_context_files wiring
# ---------------------------------------------------------------------------
class TestRunJobTerminalCwd:
"""
run_job sets TERMINAL_CWD + flips skip_context_files=False when workdir
is set, and restores the prior TERMINAL_CWD in finally — even on error.
We stub AIAgent so no real API call happens.
run_job binds workdir to its unique task CWD without mutating ambient
TERMINAL_CWD, and clears the task record in finally — even on error.
AIAgent is stubbed so no real API call happens.
"""
@staticmethod
@@ -219,7 +201,11 @@ class TestRunJobTerminalCwd:
"TERMINAL_CWD", "_UNSET_"
)
def run_conversation(self, *_a, **_kw):
def run_conversation(self, *_a, task_id=None, **_kw):
from tools.terminal_tool import get_session_cwd
observed["task_id"] = task_id
observed["task_cwd_during_run"] = get_session_cwd(task_id)
observed["terminal_cwd_during_run"] = os.environ.get(
"TERMINAL_CWD", "_UNSET_"
)
@@ -298,3 +284,35 @@ class TestRunJobTerminalCwd:
# And after run_job completes, it's still the sentinel (nothing
# overwrote or cleared it).
assert os.environ["TERMINAL_CWD"] == before
def test_workdir_is_bound_to_unique_task_without_mutating_process_env(
self, monkeypatch, tmp_path
):
import os
import cron.scheduler as sched
from tools.terminal_tool import get_session_cwd
baseline = str(tmp_path / "baseline")
workdir = tmp_path / "project"
(tmp_path / "baseline").mkdir()
workdir.mkdir()
monkeypatch.setenv("TERMINAL_CWD", baseline)
observed: dict = {}
self._install_stubs(monkeypatch, observed)
success, *_ = sched.run_job(
{
"id": "cwd-bound",
"name": "cwd-bound",
"workdir": str(workdir),
"schedule_display": "manual",
}
)
assert success is True
assert observed["skip_context_files"] is False
assert observed["task_id"].startswith("cron:cwd-bound:")
assert observed["task_cwd_during_run"] == str(workdir)
assert observed["terminal_cwd_during_run"] == baseline
assert os.environ["TERMINAL_CWD"] == baseline
assert get_session_cwd(observed["task_id"]) is None
+8 -5
View File
@@ -217,6 +217,7 @@ def test_run_one_job_records_running_then_terminal(monkeypatch):
import cron.scheduler as scheduler
events = []
run_execution_ids = []
monkeypatch.setattr(
scheduler,
"mark_execution_running",
@@ -230,16 +231,18 @@ def test_run_one_job_records_running_then_terminal(monkeypatch):
raising=False,
)
monkeypatch.setattr(scheduler, "claim_dispatch", lambda _job_id: True)
monkeypatch.setattr(
scheduler,
"run_job",
lambda job, *, defer_agent_teardown=None, **_kw: (True, "output", "response", None),
)
def fake_run_job(job, *, defer_agent_teardown=None, execution_id=None, **_kw):
run_execution_ids.append(execution_id)
return True, "output", "response", None
monkeypatch.setattr(scheduler, "run_job", fake_run_job)
monkeypatch.setattr(scheduler, "save_job_output", lambda *_args: None)
monkeypatch.setattr(scheduler, "_deliver_result", lambda *_args, **_kwargs: None)
monkeypatch.setattr(scheduler, "mark_job_run", lambda *_args, **_kwargs: None)
assert scheduler.run_one_job({"id": "job-3", "execution_id": "exec-3"}) is True
assert run_execution_ids == ["exec-3"]
assert events[0] == ("running", "exec-3")
assert events[-1][0:2] == ("finish", "exec-3")
assert events[-1][2]["success"] is True
@@ -112,8 +112,7 @@ class TestDispatchFailurePathsClearClaim:
raise RuntimeError("cannot schedule new futures")
pool = _ExplodingPool()
with patch.object(sched, "_get_parallel_pool", return_value=pool), \
patch.object(sched, "_get_sequential_pool", return_value=pool):
with patch.object(sched, "_get_parallel_pool", return_value=pool):
self._tick_one(job)
reloaded = [j for j in jobs_mod.load_jobs() if j["id"] == job["id"]][0]
assert reloaded.get("run_claim") is None
+5 -25
View File
@@ -274,21 +274,15 @@ class TestSyncMode:
sched._shutdown_parallel_pool()
class TestSequentialPool:
"""Sequential (workdir) jobs use the persistent cron-seq pool.
class TestWorkdirParallelPool:
"""Task-scoped workdir jobs use the normal persistent parallel pool."""
Verifies the follow-up fix: env-mutating jobs no longer run inline
in the ticker thread, so a long workdir job can't starve the
schedule the same way the parallel path used to.
"""
def test_sequential_job_does_not_block_ticker(self, tmp_path, monkeypatch):
def test_workdir_job_does_not_block_ticker(self, tmp_path, monkeypatch):
"""sync=False returns immediately even when a workdir job is slow."""
import cron.scheduler as sched
sched._parallel_pool = None
sched._parallel_pool_max_workers = None
sched._sequential_pool = None
sched._running_job_ids.clear()
job = {
@@ -299,7 +293,7 @@ class TestSequentialPool:
"enabled": True,
"next_run_at": "2020-01-01T00:00:00",
"deliver": "local",
"workdir": str(tmp_path), # makes it sequential
"workdir": str(tmp_path),
}
barrier = threading.Barrier(2, timeout=5)
@@ -326,13 +320,12 @@ class TestSequentialPool:
time.sleep(0.1)
sched._shutdown_parallel_pool()
def test_sequential_running_guard_prevents_double_dispatch(self, tmp_path, monkeypatch):
def test_workdir_running_guard_prevents_double_dispatch(self, tmp_path, monkeypatch):
"""A workdir job already in _running_job_ids is skipped on next tick."""
import cron.scheduler as sched
sched._parallel_pool = None
sched._parallel_pool_max_workers = None
sched._sequential_pool = None
sched._running_job_ids.clear()
job = {
@@ -364,19 +357,6 @@ class TestSequentialPool:
sched._running_job_ids.discard("guard-seq")
sched._shutdown_parallel_pool()
def test_get_sequential_pool_is_persistent(self):
"""_get_sequential_pool returns the same single-thread pool."""
import cron.scheduler as sched
sched._sequential_pool = None
pool1 = sched._get_sequential_pool()
pool2 = sched._get_sequential_pool()
assert pool1 is pool2
sched._shutdown_parallel_pool()
assert sched._sequential_pool is None
class TestTickBatchAdvance:
"""The tick's pre-dispatch advance must go through advance_next_runs
exactly once with the whole due set — a revert to the per-job loop
+3 -1
View File
@@ -1395,9 +1395,11 @@ class TestRunJobSkillBacked:
register_env_passthrough(["NOTION_API_KEY"])
return json.dumps({"success": True, "content": "# notion\nUse Notion."})
def _run_conversation(prompt):
def _run_conversation(prompt, *, task_id=None):
from tools.env_passthrough import get_all_passthrough
assert isinstance(task_id, str)
assert task_id.startswith("cron:skill-env-job:")
assert "NOTION_API_KEY" in get_all_passthrough()
return {"final_response": "ok"}
@@ -36,7 +36,9 @@ class _FakeCronAgent:
def __init__(self, *args, **kwargs):
self.kwargs = kwargs
def run_conversation(self, prompt):
def run_conversation(self, prompt, *, task_id=None):
assert isinstance(task_id, str)
assert task_id.startswith("cron:ctx-isolation:")
result = approval_module.check_execute_code_guard(
"import os; print(1)", "local"
)
-300
View File
@@ -1,300 +0,0 @@
"""Tests for the TERMINAL_CWD readers-writer lock in cron/scheduler.py.
Workdir cron jobs override the process-global ``os.environ["TERMINAL_CWD"]``
for their whole agent run. Workdir-less jobs run concurrently on a separate
pool and read that same global (via the terminal / file / code-exec tools), so
without serialization they execute commands in another job's workdir.
``_ReadWriteLock`` models workdir jobs as writers (exclusive) and workdir-less
jobs as readers (concurrent with each other, excluded from a writer's run).
These tests assert that contract.
"""
import os
import threading
import time
def _lock():
import cron.scheduler as sched
return sched._ReadWriteLock()
def test_multiple_readers_run_concurrently():
"""Workdir-less jobs (readers) hold the lock simultaneously."""
lock = _lock()
# Barrier of 3 only releases if both reader threads hold the read lock at
# the same time as the main thread waits — proving readers are concurrent.
barrier = threading.Barrier(3, timeout=5)
def reader():
lock.acquire_read()
try:
barrier.wait()
finally:
lock.release_read()
threads = [threading.Thread(target=reader) for _ in range(2)]
for t in threads:
t.start()
# Does not raise BrokenBarrierError -> both readers were holding at once.
barrier.wait(timeout=5)
for t in threads:
t.join(timeout=5)
assert not t.is_alive()
def test_writer_waits_for_active_reader():
"""A workdir job (writer) cannot acquire while a reader holds the lock."""
lock = _lock()
order = []
reader_holding = threading.Event()
let_reader_go = threading.Event()
def reader():
lock.acquire_read()
try:
reader_holding.set()
let_reader_go.wait(timeout=5)
order.append("reader-release")
finally:
lock.release_read()
def writer():
reader_holding.wait(timeout=5)
lock.acquire_write()
try:
order.append("writer-acquire")
finally:
lock.release_write()
rt = threading.Thread(target=reader)
wt = threading.Thread(target=writer)
rt.start()
wt.start()
# Give the writer time to try (and block) while the reader still holds.
reader_holding.wait(timeout=5)
let_reader_go.set()
rt.join(timeout=5)
wt.join(timeout=5)
assert not rt.is_alive() and not wt.is_alive()
# The writer only ran after the reader released — never alongside it.
assert order == ["reader-release", "writer-acquire"]
def test_writer_acquire_times_out_behind_active_reader():
"""Bounded acquire prevents a stuck reader from blocking a writer forever."""
lock = _lock()
lock.acquire_read()
started = time.monotonic()
try:
assert lock.acquire_write(timeout=0.01) is False
finally:
lock.release_read()
assert time.monotonic() - started < 1.0
def test_reader_acquire_times_out_behind_active_writer():
"""Bounded acquire prevents a stuck writer from blocking readers forever."""
lock = _lock()
lock.acquire_write()
started = time.monotonic()
try:
assert lock.acquire_read(timeout=0.01) is False
finally:
lock.release_write()
assert time.monotonic() - started < 1.0
def test_reader_never_observes_writer_override():
"""Regression: the cross-pool TERMINAL_CWD corruption.
A workdir job (writer) overriding the shared cwd must never be observed by
a concurrent workdir-less job (reader). ``shared["cwd"]`` stands in for
``os.environ["TERMINAL_CWD"]``: the reader, even though it starts while the
writer holds the override, must block until the writer restores the value.
"""
lock = _lock()
shared = {"cwd": "<scheduler>"}
observations = []
writer_holding = threading.Event()
release_writer = threading.Event()
def writer():
lock.acquire_write()
try:
shared["cwd"] = "/project/A"
writer_holding.set()
release_writer.wait(timeout=5)
finally:
shared["cwd"] = "<scheduler>"
lock.release_write()
def reader():
# Start only once the writer holds the lock and has applied the
# override — the exact window the old code corrupted.
writer_holding.wait(timeout=5)
lock.acquire_read()
try:
observations.append(shared["cwd"])
finally:
lock.release_read()
wt = threading.Thread(target=writer)
rt = threading.Thread(target=reader)
wt.start()
rt.start()
# The reader is now blocked on the writer; let the writer finish.
writer_holding.wait(timeout=5)
release_writer.set()
wt.join(timeout=5)
rt.join(timeout=5)
assert not wt.is_alive() and not rt.is_alive()
# The reader saw the restored value, never the writer's /project/A override.
assert observations == ["<scheduler>"]
def test_run_job_releases_cwd_lock_when_body_raises(tmp_path):
"""A workdir job whose run_job body raises must still RELEASE the writer lock.
Regression for the leak that made the fix "still broken": the acquire was
placed before the try whose finally releases, so an exception in the
unprotected window (or anywhere in the body) leaked the writer lock and
deadlocked the whole scheduler. This asserts the lock is free again after a
raising run — acquire_write() must not block.
"""
from unittest.mock import MagicMock, patch
import cron.scheduler as sched
workdir = tmp_path / "proj"
workdir.mkdir()
job = {"id": "boom-job", "name": "boom", "prompt": "hi", "workdir": str(workdir)}
# Force a raise in the WINDOW BETWEEN acquire and the try body — the exact
# spot the buggy placement left unprotected. With the fix these statements
# are inside the try (finally releases); with the bug the lock leaks.
# logger.info(...) fires right after os.environ["TERMINAL_CWD"] is set for a
# workdir job, in that window, so making it raise exercises the leak path.
real_info = sched.logger.info
def _raise_on_workdir_log(msg, *args, **kwargs):
if isinstance(msg, str) and "using workdir" in msg:
raise RuntimeError("boom")
return real_info(msg, *args, **kwargs)
with patch("cron.scheduler._hermes_home", tmp_path), \
patch("cron.scheduler._resolve_origin", return_value=None), \
patch("hermes_cli.env_loader.load_hermes_dotenv"), \
patch("hermes_cli.env_loader.reset_secret_source_cache"), \
patch.object(sched.logger, "info", side_effect=_raise_on_workdir_log), \
patch("hermes_state.SessionDB", return_value=MagicMock()):
# run_job catches its own body exceptions and returns (False, ...);
# it must not propagate, and it must release the lock either way.
success, _out, _final, _err = sched.run_job(job)
assert success is False
# If the writer lock leaked, this acquire would block forever. Prove it's
# free by acquiring as a writer from another thread under a short timeout.
acquired = threading.Event()
def try_acquire():
sched._terminal_cwd_lock.acquire_write()
try:
acquired.set()
finally:
sched._terminal_cwd_lock.release_write()
t = threading.Thread(target=try_acquire, daemon=True)
t.start()
assert acquired.wait(timeout=5), "writer lock was leaked by run_job on exception"
t.join(timeout=5)
def test_run_job_fails_fast_when_cwd_lock_is_stuck(monkeypatch):
"""Fail-closed (#79768): a reader blocked past the bound FAILS as a
normal cron error instead of proceeding without the lock (which would
expose the stuck writer's TERMINAL_CWD override to its commands)."""
import cron.scheduler as sched
monkeypatch.setattr(sched, "_cwd_lock_timeout_seconds", lambda: 0.05)
sched._terminal_cwd_lock.acquire_write()
try:
start = time.monotonic()
success, _out, _final, error = sched.run_job(
{"id": "blocked-reader", "name": "blocked-reader", "prompt": "hi"}
)
finally:
sched._terminal_cwd_lock.release_write()
assert success is False
assert "TERMINAL_CWD read lock" in (error or "")
assert time.monotonic() - start < 10.0
def test_run_job_writer_fails_fast_and_never_sets_env(monkeypatch, tmp_path):
"""A WRITER that times out must fail before touching TERMINAL_CWD —
the fail-open design let it clobber the active holder's override."""
import cron.scheduler as sched
monkeypatch.setattr(sched, "_cwd_lock_timeout_seconds", lambda: 0.05)
monkeypatch.setenv("TERMINAL_CWD", "/holder/dir")
sched._terminal_cwd_lock.acquire_write()
try:
success, _out, _final, error = sched.run_job(
{
"id": "blocked-writer",
"name": "blocked-writer",
"prompt": "hi",
"workdir": str(tmp_path),
}
)
observed_during = os.environ.get("TERMINAL_CWD")
finally:
sched._terminal_cwd_lock.release_write()
assert success is False
assert "TERMINAL_CWD write lock" in (error or "")
assert observed_during == "/holder/dir", (
"timed-out writer mutated the active holder's TERMINAL_CWD override"
)
assert os.environ.get("TERMINAL_CWD") == "/holder/dir"
def test_lock_wait_shorter_than_bound_still_succeeds(monkeypatch, tmp_path):
"""A waiter whose holder finishes inside the bound proceeds normally —
the fail-closed path only fires past the derived ceiling."""
import cron.scheduler as sched
lock = sched._terminal_cwd_lock
lock.acquire_write()
releaser = threading.Timer(0.1, lock.release_write)
releaser.start()
try:
assert lock.acquire_read(timeout=5.0) is True
lock.release_read()
finally:
releaser.cancel()
def test_cwd_lock_timeout_derivation(monkeypatch):
"""The bound tracks HERMES_CRON_TIMEOUT (+margin) with a floor, and
stays finite when the job runtime is unlimited (0)."""
import cron.scheduler as sched
monkeypatch.delenv("HERMES_CRON_TIMEOUT", raising=False)
assert sched._cwd_lock_timeout_seconds() == 660.0
monkeypatch.setenv("HERMES_CRON_TIMEOUT", "1800")
assert sched._cwd_lock_timeout_seconds() == 1860.0
monkeypatch.setenv("HERMES_CRON_TIMEOUT", "30")
assert sched._cwd_lock_timeout_seconds() == 180.0
monkeypatch.setenv("HERMES_CRON_TIMEOUT", "0")
assert sched._cwd_lock_timeout_seconds() == 660.0
monkeypatch.setenv("HERMES_CRON_TIMEOUT", "bogus")
assert sched._cwd_lock_timeout_seconds() == 660.0
+70
View File
@@ -241,6 +241,76 @@ def test_host_local_background_command_bypasses_configured_backend(tmp_path, mon
assert calls[0][1]["task_id"] == f"host-local-{task_id}"
def test_concurrent_commands_record_their_own_observed_cwd(monkeypatch, tmp_path):
"""A shared env's mutable cwd must not cross-write concurrent sessions."""
import threading
cwd_a = tmp_path / "a"
cwd_b = tmp_path / "b"
cwd_a.mkdir()
cwd_b.mkdir()
a_set = threading.Event()
b_set = threading.Event()
class SharedEnv:
env = {}
cwd = str(tmp_path)
def execute(self, command, **kwargs):
observed = str(cwd_a if command == "session-a" else cwd_b)
if command == "session-a":
self.cwd = observed
a_set.set()
assert b_set.wait(timeout=5)
else:
assert a_set.wait(timeout=5)
self.cwd = observed
b_set.set()
return {
"output": "ok",
"returncode": 0,
"cwd_observed": True,
"cwd": observed,
}
shared = SharedEnv()
monkeypatch.setattr(terminal_tool, "_active_environments", {"default": shared})
monkeypatch.setattr(terminal_tool, "_last_activity", {})
monkeypatch.setattr(terminal_tool, "_task_env_overrides", {})
monkeypatch.setattr(terminal_tool, "_session_cwd", {})
monkeypatch.setattr(terminal_tool, "_resolve_container_task_id", lambda _value: "default")
monkeypatch.setattr(terminal_tool, "_get_env_config", lambda: _minimal_terminal_config())
monkeypatch.setattr(
terminal_tool,
"_check_all_guards",
lambda command, env_type, **kwargs: {"approved": True},
)
terminal_tool.record_session_cwd("task-a", str(cwd_a))
terminal_tool.record_session_cwd("task-b", str(cwd_b))
results = {}
def run(name, task_id):
results[task_id] = json.loads(
terminal_tool.terminal_tool(command=name, task_id=task_id)
)
threads = [
threading.Thread(target=run, args=("session-a", "task-a")),
threading.Thread(target=run, args=("session-b", "task-b")),
]
for thread in threads:
thread.start()
for thread in threads:
thread.join(timeout=5)
assert not thread.is_alive()
assert results["task-a"]["exit_code"] == 0
assert results["task-b"]["exit_code"] == 0
assert terminal_tool.get_session_cwd("task-a") == str(cwd_a)
assert terminal_tool.get_session_cwd("task-b") == str(cwd_b)
def test_safe_getcwd_falls_back_to_home_when_no_terminal_cwd(monkeypatch):
def _boom():
raise FileNotFoundError()
+4
View File
@@ -1389,6 +1389,10 @@ class BaseEnvironment(ABC):
if cwd_path:
self.cwd = cwd_path
result["cwd_observed"] = True
# Keep the observation on this command's result as well as on the
# shared environment. Concurrent callers must not read self.cwd
# after another command has already updated it.
result["cwd"] = cwd_path
# Strip the marker line AND the \n we injected before it.
# The wrapper emits: printf '\n__MARKER__%s__MARKER__\n'
+2
View File
@@ -2322,6 +2322,7 @@ class LocalEnvironment(BaseEnvironment):
normalized = _msys_to_windows_path(self.cwd) if _IS_WINDOWS else self.cwd
if normalized and os.path.isdir(normalized):
self.cwd = normalized
result["cwd"] = normalized
else:
# Stale / non-existent path — keep previous cwd; _run_bash
# will resolve a safe fallback on the next call if needed.
@@ -2329,6 +2330,7 @@ class LocalEnvironment(BaseEnvironment):
# so it is not attributable to this command's session either.
self.cwd = prev_cwd
result.pop("cwd_observed", None)
result.pop("cwd", None)
def cleanup(self):
"""Clean up temp files."""
+11 -3
View File
@@ -3639,8 +3639,16 @@ def terminal_tool(
# the last command to FINISH left there — on a shared env, that is
# another session's directory. Recording it silently re-homes this
# session into a directory the user never opened.
if not workdir and (result or {}).get("cwd_observed"):
record_session_cwd(session_key, getattr(env, "cwd", None))
observed_cwd = None
if (result or {}).get("cwd_observed"):
# New/current environments return the CWD observed by THIS
# command. The env field is shared mutable compatibility state
# and may already belong to a concurrent command. Keep the
# fallback for third-party providers that only implement the
# older cwd_observed + env.cwd contract.
observed_cwd = (result or {}).get("cwd") or getattr(env, "cwd", None)
if not workdir and observed_cwd:
record_session_cwd(session_key, observed_cwd)
# Extract output
output = result.get("output", "")
@@ -3768,7 +3776,7 @@ def terminal_tool(
# and tells the model it moved to a directory another session
# opened.
try:
post_cwd = getattr(env, "cwd", None) if (result or {}).get("cwd_observed") else None
post_cwd = observed_cwd
if post_cwd and command_cwd and os.path.realpath(str(post_cwd)) != os.path.realpath(str(command_cwd)):
result_dict["cwd"] = str(post_cwd)
except Exception:
+2 -2
View File
@@ -187,8 +187,8 @@ When `workdir` is set:
- The path must be an absolute directory that exists — relative paths and missing directories are rejected at create / update time
- Pass `--workdir ""` (or `workdir=""` via the tool) on edit to clear it and restore the old behaviour
:::note Serialization
Jobs with a `workdir` run sequentially on the scheduler tick, not in the parallel pool. This is deliberate: the cron worker applies the job workdir through process-global terminal state, so two workdir jobs running at the same time would corrupt each other's cwd. Workdir-less jobs still run in parallel as before.
:::note Isolation
Each agent run binds its `workdir` to that run's unique task identity. Workdir jobs therefore use the normal parallel pool without mutating process-global terminal state or leaking paths between concurrent runs. Set `cron.max_parallel_jobs` if you want to limit total cron concurrency.
:::
## Editing jobs