fix(code-execution): remote kernel reap/evict skip kernels with a running cell
Same attached-cell guard as the local kernel host (#101861): a remote kernel mid-cell is never reaped or cap-evicted, so a fan-out never has its runner killed under a live poll loop.
This commit is contained in:
@@ -244,6 +244,36 @@ class TestIdleReapAndCapEviction(RemoteKernelBase):
|
||||
self.assertIn("owner-1", owners)
|
||||
self.assertIn("owner-2", owners)
|
||||
|
||||
def test_eviction_skips_kernels_with_a_running_cell(self):
|
||||
"""Cap eviction must never kill a kernel mid-cell (the local-kernel
|
||||
race from hermes-agent#101861): a busy kernel stays put and a
|
||||
settled one goes instead, even if the busy one is older."""
|
||||
import threading
|
||||
|
||||
gate = threading.Event()
|
||||
|
||||
def slow_cat(command):
|
||||
gate.wait(10)
|
||||
return {"output": json.dumps(_cell()), "returncode": 0}
|
||||
|
||||
busy_env = ScriptedEnv([
|
||||
("nohup", lambda c: {"output": "PID:4242\n", "returncode": 0}),
|
||||
("kill -0", lambda c: {"output": "ALIVE\n", "returncode": 0}),
|
||||
("cat ", slow_cat),
|
||||
])
|
||||
with patch("tools.code_kernel._lifecycle_limits", return_value=(1, 1800)):
|
||||
worker = threading.Thread(target=_run, args=(busy_env,), kwargs={"task": "busy"})
|
||||
worker.start()
|
||||
while not any(k.attached for k in _REMOTE_KERNELS.values()):
|
||||
pass
|
||||
env = ScriptedEnv(_spawn_ok_handlers([_cell()]))
|
||||
_run(env, task="settled")
|
||||
owners = {key[0] for key in _REMOTE_KERNELS}
|
||||
self.assertIn("busy", owners)
|
||||
gate.set()
|
||||
worker.join(10)
|
||||
self.assertFalse(any("kill 4242" in c for c in busy_env.commands))
|
||||
|
||||
|
||||
class TestDispatchIntegration(unittest.TestCase):
|
||||
"""_execute_remote prefers the kernel and falls open to per-call."""
|
||||
|
||||
@@ -156,6 +156,10 @@ class RemoteKernel:
|
||||
last_used: float = field(default_factory=time.monotonic)
|
||||
execution_count: int = 0
|
||||
cell_seq: int = 0
|
||||
# Cells currently running on this kernel. Reap/evict skip attached
|
||||
# kernels: killing one mid-cell tears the runner out from under a live
|
||||
# poll loop (same guard as tools.code_kernel, hermes-agent#101861).
|
||||
attached: int = 0
|
||||
|
||||
|
||||
def _kernel_key(owner: str, env_type: str, task_env_id: str) -> Tuple:
|
||||
@@ -233,7 +237,7 @@ def _reap_unlocked(idle_timeout: int) -> List["RemoteKernel"]:
|
||||
doomed = [
|
||||
key
|
||||
for key, kernel in _REMOTE_KERNELS.items()
|
||||
if now - kernel.last_used > idle_timeout
|
||||
if kernel.attached == 0 and now - kernel.last_used > idle_timeout
|
||||
]
|
||||
return [_REMOTE_KERNELS.pop(key) for key in doomed]
|
||||
|
||||
@@ -250,7 +254,7 @@ def _evict_over_cap_unlocked(keep: Tuple) -> List["RemoteKernel"]:
|
||||
if len(_REMOTE_KERNELS) <= cap:
|
||||
return []
|
||||
by_age = sorted(
|
||||
(key for key in _REMOTE_KERNELS if key != keep),
|
||||
(key for key in _REMOTE_KERNELS if key != keep and _REMOTE_KERNELS[key].attached == 0),
|
||||
key=lambda key: _REMOTE_KERNELS[key].last_used,
|
||||
)
|
||||
doomed = by_age[: len(_REMOTE_KERNELS) - cap]
|
||||
@@ -355,11 +359,6 @@ def execute_in_remote_kernel(
|
||||
the ``kernel`` sub-dict, matching the local kernel's result shape.
|
||||
"""
|
||||
from tools.code_kernel import _resolve_owner
|
||||
from tools.code_execution_tool import (
|
||||
_rpc_poll_loop,
|
||||
_ship_file_to_remote,
|
||||
)
|
||||
from tools.thread_context import propagate_context_to_thread
|
||||
|
||||
owner = _resolve_owner(task_env_id)
|
||||
key = _kernel_key(owner, env_type, task_env_id)
|
||||
@@ -401,9 +400,43 @@ def execute_in_remote_kernel(
|
||||
|
||||
kernel.last_used = time.monotonic()
|
||||
with _REMOTE_KERNELS_LOCK:
|
||||
kernel.attached += 1
|
||||
evicted = _evict_over_cap_unlocked(keep=key)
|
||||
for doomed in evicted:
|
||||
_kill(doomed)
|
||||
try:
|
||||
return _run_remote_cell(
|
||||
kernel, key, code, env=env, task_env_id=task_env_id,
|
||||
sandbox_tools=sandbox_tools, timeout=timeout,
|
||||
max_tool_calls=max_tool_calls, reused=reused,
|
||||
state_reset=state_reset, state_lost=state_lost,
|
||||
)
|
||||
finally:
|
||||
with _REMOTE_KERNELS_LOCK:
|
||||
kernel.attached -= 1
|
||||
kernel.last_used = time.monotonic()
|
||||
|
||||
|
||||
def _run_remote_cell(
|
||||
kernel: RemoteKernel,
|
||||
key: Tuple,
|
||||
code: str,
|
||||
*,
|
||||
env,
|
||||
task_env_id: str,
|
||||
sandbox_tools: frozenset,
|
||||
timeout: int,
|
||||
max_tool_calls: int,
|
||||
reused: bool,
|
||||
state_reset: bool,
|
||||
state_lost: bool,
|
||||
) -> Dict[str, Any]:
|
||||
from tools.code_execution_tool import (
|
||||
_rpc_poll_loop,
|
||||
_ship_file_to_remote,
|
||||
)
|
||||
from tools.thread_context import propagate_context_to_thread
|
||||
|
||||
kernel.cell_seq += 1
|
||||
seq = f"{kernel.cell_seq:06d}"
|
||||
q_cells = shlex.quote(f"{kernel.kernel_dir}/cells")
|
||||
|
||||
Reference in New Issue
Block a user