refactor(process-registry): fold kill_started_since into kill_all via exclude_ids

kill_started_since duplicated kill_all's collect-under-lock/kill-outside-lock
loop line for line; it is now a thin delegate through new kill_all kwargs
(exclude_ids, source, consume_output). Public signatures unchanged — existing
callers and test monkeypatch seams keep working. kill_process's docstring now
names the deliberate consume_output=True exception for abandoned-turn reaping
so the deviation isn't 'fixed' later.
This commit is contained in:
kshitijk4poor 2026-08-02 14:10:05 +05:30 committed by kshitij
parent 8f0f55eaac
commit cd6585abf8
2 changed files with 58 additions and 28 deletions

View File

@ -92,6 +92,33 @@ def test_kill_started_since_preserves_preexisting_and_foreign_processes(registry
]
def test_kill_all_backward_compat_and_exclude_ids(registry):
"""kill_all keeps its historical default behavior (kill everything for
the task, consume_output=False, source='kill_all') and honors the new
exclude_ids kwarg that kill_started_since delegates through (#76188)."""
a = _make_session(sid="proc_a", task_id="session-a")
b = _make_session(sid="proc_b", task_id="session-a")
registry._running[a.id] = a
registry._running[b.id] = b
calls = []
def fake_kill(session_id, **kwargs):
calls.append((session_id, kwargs))
return {"status": "killed"}
registry.kill_process = fake_kill
assert registry.kill_all("session-a", exclude_ids=frozenset({"proc_a"})) == 1
assert calls == [
("proc_b", {"source": "kill_all", "consume_output": False})
]
calls.clear()
assert registry.kill_all("session-a") == 2
assert sorted(c[0] for c in calls) == ["proc_a", "proc_b"]
def _wait_until(predicate, timeout: float = 5.0, interval: float = 0.05) -> bool:
"""Poll a predicate until it returns truthy or the timeout elapses."""
deadline = time.monotonic() + timeout

View File

@ -1617,7 +1617,10 @@ class ProcessRegistry:
``consume_output`` is true for explicit tool/RPC kills because their
caller observes the returned output. Bulk cleanup passes false: it
discards each result and therefore must not suppress an autonomous
output-bearing completion notification.
output-bearing completion notification. Exception: abandoned-turn
reaping (``kill_started_since``) is bulk cleanup that deliberately
passes true a killed abandoned process must not enqueue a synthetic
follow-up that revives work the timeout/interrupt stopped.
"""
from tools.ansi_strip import strip_ansi
@ -1951,13 +1954,34 @@ class ProcessRegistry:
*,
source: str,
) -> int:
"""Kill processes created for ``task_id`` after a prior snapshot."""
baseline = frozenset(baseline_ids or ())
"""Kill processes created for ``task_id`` after a prior snapshot.
``consume_output`` is forced on: abandoned-turn output must not
enqueue a synthetic follow-up that revives work the timeout
deliberately stopped.
"""
return self.kill_all(
task_id,
exclude_ids=frozenset(baseline_ids or ()),
source=source,
consume_output=True,
)
def kill_all(
self,
task_id: Optional[str] = None,
*,
exclude_ids: frozenset = frozenset(),
source: str = "kill_all",
consume_output: bool = False,
) -> int:
"""Kill all running processes, optionally filtered by task_id. Returns count killed."""
with self._lock:
targets = [
s
for s in self._running.values()
if s.task_id == task_id and s.id not in baseline and not s.exited
s for s in self._running.values()
if (task_id is None or s.task_id == task_id)
and s.id not in exclude_ids
and not s.exited
]
killed = 0
@ -1965,28 +1989,7 @@ class ProcessRegistry:
result = self.kill_process(
session.id,
source=source,
# Abandoned-turn output must not enqueue a synthetic follow-up
# that revives work the timeout deliberately stopped.
consume_output=True,
)
if result.get("status") in {"killed", "already_exited"}:
killed += 1
return killed
def kill_all(self, task_id: str = None) -> int:
"""Kill all running processes, optionally filtered by task_id. Returns count killed."""
with self._lock:
targets = [
s for s in self._running.values()
if (task_id is None or s.task_id == task_id) and not s.exited
]
killed = 0
for session in targets:
result = self.kill_process(
session.id,
source="kill_all",
consume_output=False,
consume_output=consume_output,
)
if result.get("status") in {"killed", "already_exited"}:
killed += 1