Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
52 changes: 46 additions & 6 deletions agent_core/core/impl/action/cancellation.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,9 @@
- ``mark_subprocess`` / ``unmark_subprocess``: pid marker FILES under the
system temp dir, used by code running in a DIFFERENT process (the
sandboxed-action pool worker) where no in-memory registry can be shared.
A marker records the pid AND its start time: a marker can outlive its
process (crash, reboot), and by then the OS may have handed the pid to an
unrelated program. A marker is only acted on while both still match.

``kill_session_processes(session_id)`` kills both kinds, entire process
trees included, and is safe to call at any time (missing/exited processes
Expand All @@ -28,12 +31,13 @@

from __future__ import annotations

import json
import os
import subprocess
import tempfile
import threading
from pathlib import Path
from typing import Dict
from typing import Dict, Optional

from agent_core.utils.logger import logger

Expand Down Expand Up @@ -71,6 +75,41 @@ def unregister_process(session_id: str, proc: subprocess.Popen) -> None:

# ─────────────────────── Cross-process pid markers ───────────────────────

# psutil's create_time can wobble by a fraction of a second between reads.
_START_TIME_TOLERANCE_S = 1.0


def _start_time(pid: int) -> Optional[float]:
"""When `pid` started, or None if it is gone or can't be inspected."""
try:
import psutil

return psutil.Process(pid).create_time()
except Exception:
return None


def _marker_still_ours(marker: Path) -> Optional[int]:
"""The marker's pid if that pid is still the process that was marked."""
try:
pid = int(marker.stem)
raw = marker.read_text(encoding="utf-8").strip()
mtime = marker.stat().st_mtime
except (ValueError, OSError):
return None
actual = _start_time(pid)
if actual is None:
return None
try:
recorded = float(json.loads(raw).get("started"))
except Exception:
recorded = None
if recorded is not None:
return pid if abs(actual - recorded) <= _START_TIME_TOLERANCE_S else None
# Legacy marker (bare pid): it was written after the child spawned, so a
# process that started after the marker is a recycled pid, not ours.
return pid if actual <= mtime + _START_TIME_TOLERANCE_S else None


def mark_subprocess(session_id: str, pid: int) -> None:
"""Record a child pid from ANOTHER process (e.g. a pool worker).
Expand All @@ -83,7 +122,9 @@ def mark_subprocess(session_id: str, pid: int) -> None:
try:
d = _marker_dir(session_id)
d.mkdir(parents=True, exist_ok=True)
(d / f"{pid}.pid").write_text(str(pid), encoding="utf-8")
(d / f"{pid}.pid").write_text(
json.dumps({"pid": pid, "started": _start_time(pid)}), encoding="utf-8"
)
except Exception:
pass # markers are best-effort; never fail the action over them

Expand Down Expand Up @@ -149,11 +190,10 @@ def kill_session_processes(session_id: str) -> int:
d = _marker_dir(session_id)
if d.is_dir():
for marker in d.glob("*.pid"):
try:
_kill_tree(int(marker.stem))
pid = _marker_still_ours(marker)
if pid is not None:
_kill_tree(pid)
killed += 1
except ValueError:
pass
marker.unlink(missing_ok=True)
except Exception as e:
logger.debug(f"[CANCEL] Marker sweep failed for {session_id}: {e}")
Expand Down
38 changes: 13 additions & 25 deletions app/agent_app/lifecycle/provisioner.py
Original file line number Diff line number Diff line change
Expand Up @@ -25,10 +25,8 @@
"""

import json
import os
import re
import shutil
import signal
import subprocess
import time
from dataclasses import dataclass
Expand All @@ -43,6 +41,7 @@
logger = logging.getLogger(__name__)

from app.agent_app.instances import Instance
from app.process_ledger import get_ledger, kill_tree

# Content-addressed artifacts to keep per project (the newest is usually the
# only one that matters; a couple of spares make flip-flopping edits cheap).
Expand Down Expand Up @@ -183,9 +182,9 @@ def _prune_builds(self, builds: Path) -> int:
def reap_dirs(self) -> int:
"""Startup dir sweep: delete every shadow boot dir/build cache and the
retired dev-copy root (`_staging/project`). Process kills are NOT done
here — the manager kills leftovers by the ports it OWNS (both ranges),
which is verified against its own records instead of a stored pid that
may have been reused. `_staging/wizard` is the wizard's attachment
here — the manager reaps leftovers from the owned-process ledger,
which verifies pid AND start time, so a pid that may have been reused
is never killed. `_staging/wizard` is the wizard's attachment
staging and is never ours to touch."""
reaped = 0
for root in (self.root, self.agent_app_dir / "_staging" / "project"):
Expand Down Expand Up @@ -219,32 +218,21 @@ def _guarded_rmtree(self, target: Path) -> None:
shutil.rmtree(resolved)

def _kill(self, process=None, pid: Optional[int] = None) -> None:
target_pid = pid if pid else (process.pid if process is not None else None)
"""Stop a shadow PocketBase and its tree (a bare terminate strands
grandchildren that keep the port bound). A live handle is ours by
construction; a bare pid — all that survives a CraftBot restart — is
killed only if the ledger recorded exactly that process, so a pid the
OS has since recycled is never touched."""
ledger = get_ledger()
try:
if os.name == "nt" and target_pid:
# PocketBase (and any node child) is a tree; a bare terminate
# strands grandchildren that keep the port bound.
subprocess.run(
["taskkill", "/T", "/F", "/PID", str(target_pid)],
capture_output=True,
creationflags=getattr(subprocess, "CREATE_NO_WINDOW", 0),
)
return
if process is not None:
process.terminate()
kill_tree(process.pid)
try:
process.wait(timeout=5)
except Exception:
process.kill()
ledger.forget(process.pid)
elif pid:
os.kill(pid, signal.SIGTERM)
time.sleep(0.5)
try:
os.kill(pid, 0)
except OSError:
return # already gone
os.kill(pid, signal.SIGKILL)
except ProcessLookupError:
pass
ledger.kill_pid(pid)
except Exception as e:
logger.warning(f"[AGENT_APP:SHADOW] kill failed: {e}")
Loading