From 8aa2c2d19216e3c443d76bd12dd892ce8432376a Mon Sep 17 00:00:00 2001 From: ahmad-ajmal Date: Tue, 22 Sep 2026 14:21:49 +0100 Subject: [PATCH] fix: process killing logic --- agent_core/core/impl/action/cancellation.py | 52 ++- app/agent_app/lifecycle/provisioner.py | 38 +- app/agent_app/manager.py | 374 ++++++----------- app/agent_app/runner.py | 9 + app/process_ledger.py | 433 ++++++++++++++++++++ app/ui_layer/local_llm_setup.py | 49 ++- craftbot.py | 100 +++-- environment.yml | 1 + installer/helpers.py | 88 +++- main.py | 158 +------ requirements.txt | 1 + run.py | 104 +++-- tests/test_agent_app_undeployed.py | 6 +- tests/test_pid_identity.py | 139 +++++++ tests/test_process_ledger.py | 192 +++++++++ 15 files changed, 1206 insertions(+), 538 deletions(-) create mode 100644 app/process_ledger.py create mode 100644 tests/test_pid_identity.py create mode 100644 tests/test_process_ledger.py diff --git a/agent_core/core/impl/action/cancellation.py b/agent_core/core/impl/action/cancellation.py index 445e0c80..27b5918f 100644 --- a/agent_core/core/impl/action/cancellation.py +++ b/agent_core/core/impl/action/cancellation.py @@ -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 @@ -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 @@ -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). @@ -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 @@ -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}") diff --git a/app/agent_app/lifecycle/provisioner.py b/app/agent_app/lifecycle/provisioner.py index 0676662d..51a11736 100644 --- a/app/agent_app/lifecycle/provisioner.py +++ b/app/agent_app/lifecycle/provisioner.py @@ -25,10 +25,8 @@ """ import json -import os import re import shutil -import signal import subprocess import time from dataclasses import dataclass @@ -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). @@ -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"): @@ -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}") diff --git a/app/agent_app/manager.py b/app/agent_app/manager.py index 1b24c92a..fa17bccb 100644 --- a/app/agent_app/manager.py +++ b/app/agent_app/manager.py @@ -11,6 +11,7 @@ """ import asyncio +import errno import json import os import re @@ -27,9 +28,16 @@ from dataclasses import dataclass, field from datetime import datetime from pathlib import Path -from typing import Dict, List, Optional, Any, Set, Tuple, TYPE_CHECKING +from typing import Dict, List, Optional, Any, Tuple, TYPE_CHECKING from app import node_runtime +from app.process_ledger import ( + ROLE_AGENT_APP, + ROLE_TUNNEL, + get_ledger, + kill_tree, + listening_pids, +) from app.agent_app import marketplace_source try: @@ -1066,125 +1074,22 @@ def _can_bind(port: int) -> bool: return False def _is_port_in_use(self, port: int) -> bool: - """Check if a port is actually in use on the system.""" - with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as s: - s.settimeout(0.5) - return s.connect_ex(("localhost", port)) == 0 - - def _port_is_ours(self, port: int, ports_to_check: Optional[Set[int]]) -> bool: - """True when `port` is one we should consider for reaping: either it is - in the explicit owned set, or (when no set is given) it falls in EITHER - Agent App pool (live 3100-3199 or shadow 3900-3999). Ownership is - structural — a port range and our own records — never a match on a - process command line.""" - if ports_to_check is not None: - return port in ports_to_check - from app.agent_app.instances import LIVE_RANGE, SHADOW_RANGE - - return (LIVE_RANGE[0] <= port <= LIVE_RANGE[1]) or ( - SHADOW_RANGE[0] <= port <= SHADOW_RANGE[1] - ) - - def _get_pids_on_ports( - self, ports_to_check: Optional[Set[int]] = None - ) -> Dict[int, str]: - """ - Get PIDs of processes listening on ports in the Agent App range. - Uses a single system call for efficiency. - - Args: - ports_to_check: Optional set of specific ports to check. - If None, checks all ports in the Agent App range. + """Check if a port is actually in use on the system. - Returns: - Dict mapping port numbers to PIDs + A timed-out probe is ambiguous: it is what a listener with a full + accept queue (a busy or wedged app) looks like, but on Windows a + closed loopback port times out too (SYN retries before the refusal). + Reading every timeout as "free" made the watchdog declare a live app + dead, so a timeout is settled by the OS listener table instead. """ - port_pids = {} - - if os.name == "nt": - # Windows: run netstat once and parse all results - try: - result = subprocess.run( - ["netstat", "-ano"], - capture_output=True, - text=True, - shell=True, - timeout=5, - creationflags=subprocess.CREATE_NO_WINDOW - if hasattr(subprocess, "CREATE_NO_WINDOW") - else 0, - ) - for line in result.stdout.split("\n"): - if "LISTENING" in line: - parts = line.split() - if len(parts) >= 5: - addr = parts[1] - pid = parts[-1] - if ":" in addr: - try: - port = int(addr.split(":")[-1]) - if self._port_is_ours(port, ports_to_check): - port_pids[port] = pid - except ValueError: - pass - except Exception as e: - logger.warning(f"[AGENT_APP] Failed to get ports via netstat: {e}") - else: - # Linux/Mac: use lsof - try: - result = subprocess.run( - ["lsof", "-i", "-P", "-n"], - capture_output=True, - text=True, - timeout=5, - ) - for line in result.stdout.split("\n"): - if "LISTEN" in line: - parts = line.split() - if len(parts) >= 2: - # PID is typically the second column - pid = parts[1] - # Find the port in the line - for part in parts: - if ":" in part: - try: - port = int(part.split(":")[-1]) - if self._port_is_ours(port, ports_to_check): - port_pids[port] = pid - break - except ValueError: - pass - except Exception as e: - logger.warning(f"[AGENT_APP] Failed to get ports via lsof: {e}") - - return port_pids - - def _kill_process_by_pid(self, pid: str) -> bool: - """ - Kill a process by its PID. - - Args: - pid: Process ID to kill - - Returns: - True if process was killed, False otherwise - """ - try: - if os.name == "nt": - subprocess.run( - ["taskkill", "/F", "/PID", pid], - capture_output=True, - shell=True, - creationflags=subprocess.CREATE_NO_WINDOW - if hasattr(subprocess, "CREATE_NO_WINDOW") - else 0, - ) - else: - subprocess.run(["kill", "-9", pid], capture_output=True) + with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as s: + s.settimeout(0.5) + rc = s.connect_ex(("127.0.0.1", port)) + if rc == 0: return True - except Exception as e: - logger.warning(f"[AGENT_APP] Failed to kill process {pid}: {e}") + if rc == errno.ECONNREFUSED: return False + return bool(listening_pids(port)) # ======================================================================== # Manifest-driven launch pipeline @@ -1260,8 +1165,20 @@ def _fail(step: str, errors: list) -> dict: except AgentAppRunnerUnavailable as e: return _fail("setup", [str(e)]) - # Clear any stale listener before binding the port. + # Clear any stale listener before binding the port. Only our own + # leftovers are killed; a program CraftBot did not start keeps the + # port, and booting anyway would "pass" health against THAT program. self._kill_process_on_port(port) + holder = self._foreign_listener(port) + if holder is not None: + return _fail( + "start", + [ + f"port {port} is in use by another program (pid {holder}) " + f"that CraftBot did not start, so it was left running. " + f"Close it, or move this app to a different port." + ], + ) # Renamed/deleted APPLIED migrations brick the boot with an error only # pocketbase.log ever sees — catch them here, before any process spawns. @@ -1343,7 +1260,12 @@ def _pb_log_since_boot(limit_lines: int = 30) -> str: except Exception as e: return _fail("start", [str(e)]) - if not await self.runner.wait_healthy(port): + # Healthy means OUR process answered: a 200 from whoever else holds + # the port (PocketBase exited on "address in use") is not a boot. + healthy = await self.runner.wait_healthy(port) + if healthy and (process.poll() is not None or process.pid not in listening_pids(port)): + healthy = False + if not healthy: self._terminate_process(process) # A dead health check with no cause starved the agent before — # the boot abort (bad migration, hook panic) is in pocketbase.log @@ -1456,10 +1378,7 @@ def _fail(step: str, errors: list) -> dict: except Exception: pass if self._is_port_in_use(port): - own_pid = str(os.getpid()) - holder = self._get_pids_on_ports({port}).get(port) - if holder is None or str(holder) != own_pid: - self._kill_process_on_port(port) + self._kill_process_on_port(port) # The app itself binds a fresh hidden internal port each launch. if project.internal_port: @@ -1603,6 +1522,12 @@ def _log_since_boot(limit_lines: int = 30) -> str: errors.append("app.log (this boot):\n" + boot_log) return _fail("health", errors) + # The start command runs under a shell; record the server itself too, + # so it stays recognisably ours even if that shell dies first. + get_ledger().adopt_listeners( + internal_port, ROLE_AGENT_APP, owner=project.id, label=f"server :{internal_port}" + ) + # A2App adapter in front of the healthy app: bind the project port, # then structurally self-check the surface (the identity probe is # the only reliable check — a status code never is). @@ -2254,6 +2179,12 @@ def _start_process( stderr=log_handle, shell=True, ) + get_ledger().register( + process.pid, + ROLE_AGENT_APP, + owner=project.id if project else "", + label=f"{command[:80]} :{port}" if port else command[:80], + ) return process def _create_frontend_log(project_path: Path) -> Path: @@ -2273,93 +2204,45 @@ def _read_log_tail(log_file: Path, chars: int = 1000) -> str: return "(could not read log)" def _terminate_process(self, process: subprocess.Popen) -> None: - """Terminate a subprocess, killing the entire process tree on Windows.""" + """Stop a subprocess WE hold the handle of, with its whole tree. + + shell=True puts a cmd.exe / sh between us and the server; stopping + only that shell strands the server still bound to its port (and, on + POSIX, reparented away from anything the ledger could trace back). + """ try: - if os.name == "nt": - # On Windows with shell=True, terminate() only kills cmd.exe, - # not the child python/uvicorn. Kill the whole tree via taskkill. - subprocess.run( - ["taskkill", "/T", "/F", "/PID", str(process.pid)], - capture_output=True, - shell=True, - creationflags=subprocess.CREATE_NO_WINDOW - if hasattr(subprocess, "CREATE_NO_WINDOW") - else 0, - ) - else: - process.terminate() + kill_tree(process.pid) process.wait(timeout=5) except (subprocess.TimeoutExpired, Exception): try: process.kill() except Exception: pass + get_ledger().forget(process.pid) - def _kill_process_on_port(self, port: int) -> bool: - """ - Kill any process listening on the specified port (Windows-specific). + @staticmethod + def _foreign_listener(port: int) -> Optional[int]: + """A pid listening on `port` that CraftBot did not start (and is not + CraftBot itself), or None.""" + ledger = get_ledger() + for pid in listening_pids(port): + if pid != os.getpid() and ledger.owner_of(pid) is None: + return pid + return None - Args: - port: The port to free + def _kill_process_on_port(self, port: int) -> bool: + """Free `port` by killing its listener — ONLY if CraftBot started it. - Returns: - True if a process was killed, False otherwise + The listener is matched on the exact port (never a substring of a + netstat line) and must be a process in the owned-process ledger, or a + descendant of one. A foreign service on the port is left alone; so is + CraftBot's own process (the in-process A2App proxy holds project + ports). Returns True if a process was killed. """ - if os.name != "nt": - # Linux/Mac: use lsof and kill - try: - result = subprocess.run( - ["lsof", "-ti", f":{port}"], capture_output=True, text=True - ) - if result.stdout.strip(): - pids = result.stdout.strip().split("\n") - for pid in pids: - subprocess.run(["kill", "-9", pid], capture_output=True) - logger.info(f"[AGENT_APP] Killed process(es) on port {port}") - return True - except Exception as e: - logger.warning( - f"[AGENT_APP] Failed to kill process on port {port}: {e}" - ) - return False - else: - # Windows: use netstat and taskkill - try: - no_window = ( - subprocess.CREATE_NO_WINDOW - if hasattr(subprocess, "CREATE_NO_WINDOW") - else 0 - ) - result = subprocess.run( - ["netstat", "-ano"], - capture_output=True, - text=True, - shell=True, - creationflags=no_window, - ) - killed = False - for line in result.stdout.split("\n"): - if f":{port}" in line and "LISTENING" in line: - parts = line.split() - if len(parts) >= 5: - pid = parts[-1] - # /T kills entire process tree (shell + child processes) - subprocess.run( - ["taskkill", "/T", "/F", "/PID", pid], - capture_output=True, - shell=True, - creationflags=no_window, - ) - logger.info( - f"[AGENT_APP] Killed process tree {pid} on port {port}" - ) - killed = True - if killed: - return True - except Exception as e: - logger.warning( - f"[AGENT_APP] Failed to kill process on port {port}: {e}" - ) + try: + return get_ledger().kill_port_listeners(port) + except Exception as e: + logger.warning(f"[AGENT_APP] Failed to free port {port}: {e}") return False def cleanup_on_startup(self) -> None: @@ -2367,41 +2250,28 @@ def cleanup_on_startup(self) -> None: Clean up orphan processes and folders on startup. This should be called after loading projects to: - 1. Kill any orphan Agent App server processes on tracked ports (frontend + backend) - 2. Delete project folders not tracked in the registry - 3. Reset all project statuses to 'stopped' - - Optimized to: - - Only check ports that are tracked in projects (not all 100 ports) - - Use a single netstat call to get all port info at once + 1. Kill leftover Agent App servers and tunnels from the previous run — + only the ones recorded in the owned-process ledger + 2. Log project folders not tracked in the registry + 3. Sweep dev-env leftover directories """ logger.info("[AGENT_APP] Running startup cleanup...") - # Nothing of ours is legitimately running yet: every app process died - # with the previous CraftBot. Reconcile against the OS by the ports WE - # OWN — every project's sticky LIVE port (persisted) plus every port a - # prior registry instance (live or shadow) claimed. Snapshotting and - # clearing the registry first releases the shadow reservations; the - # live ports were re-reserved from the project list at load. - # - # We kill only listeners sitting on ports WE own, by pid. Ownership is - # structural (our port ranges and our own records), so a foreign - # process on a port we never claimed is never touched, and there is no - # command-line/string matching to decide "is this ours". - prior = self.instances.reset() - owned_ports = {p.port for p in self.projects.values() if p.port} - owned_ports |= {i.port for i in prior} - - killed_count = 0 - if owned_ports: - for port, pid in self._get_pids_on_ports(owned_ports).items(): - if self._kill_process_by_pid(pid): - killed_count += 1 - logger.info( - f"[AGENT_APP] reclaimed owned port {port} (pid {pid})" - ) + # Nothing of ours is legitimately running yet: every app server and + # tunnel died with the previous CraftBot, or should have. Reap exactly + # the processes the ledger recorded starting — verified by pid AND + # creation time, so a recycled pid or a foreign service that happens + # to sit on one of our ports is never touched. Snapshotting and + # clearing the instance registry releases the shadow reservations; + # the live ports were re-reserved from the project list at load. + self.instances.reset() + try: + killed_count = get_ledger().reap() + except Exception as e: + killed_count = 0 + logger.warning(f"[AGENT_APP] owned-process reap failed: {e}") if killed_count > 0: - logger.info(f"[AGENT_APP] reclaimed {killed_count} leftover process(es)") + logger.info(f"[AGENT_APP] reaped {killed_count} leftover owned process(es)") # Log orphan project folders (do NOT delete — deleting them at boot has # destroyed real user projects; logging is the safe behavior). @@ -2412,8 +2282,8 @@ def cleanup_on_startup(self) -> None: ) # Sweep dev-env leftover directories. Process kills were handled above - # by owned-port reclaim across BOTH ranges (the shadow range included), - # so a deleted-project shadow can no longer leak an untracked process. + # by the ledger reap, which covers shadow boots too (the runner + # records every PocketBase it starts, live or shadow). try: reaped = self.lifecycle.reap_dirs() if reaped: @@ -4734,34 +4604,19 @@ async def start_tunnel( logger.info("[AGENT_APP] Stopping any existing tunnel...") await self.stop_tunnel(project_id) - # Only kill orphans on first tunnel start (no other tunnels active) - other_tunnels = any( - p.tunnel_process is not None and p.id != project_id + # Reap tunnels a previous CraftBot run left behind — by the exact + # pid+start-time we recorded when starting them, never by process + # name: a user's own cloudflared (e.g. a production tunnel) is not + # ours to kill. Tunnels other projects hold right now are skipped. + held = { + p.tunnel_process.pid for p in self.projects.values() - ) - if not other_tunnels: - logger.info( - "[AGENT_APP] No other tunnels active, cleaning orphan cloudflared processes..." - ) - try: - if os.name == "nt": - subprocess.run( - [ - "powershell", - "-Command", - "Stop-Process -Name cloudflared -Force -ErrorAction SilentlyContinue", - ], - capture_output=True, - timeout=5, - creationflags=subprocess.CREATE_NO_WINDOW - if hasattr(subprocess, "CREATE_NO_WINDOW") - else 0, - ) - else: - subprocess.run(["pkill", "-f", "cloudflared"], capture_output=True) - await asyncio.sleep(1) - except Exception: - pass + if p.tunnel_process is not None + } + ledger = get_ledger() + for entry in ledger.entries(ROLE_TUNNEL): + if entry.pid not in held: + ledger.kill_entry(entry) port = self._serving_port(project) if not port: @@ -4805,6 +4660,9 @@ async def start_tunnel( if os.name == "nt" and hasattr(subprocess, "CREATE_NO_WINDOW") else 0, ) + get_ledger().register( + proc.pid, ROLE_TUNNEL, owner=project.id, label=origin_url + ) logger.info(f"[AGENT_APP] cloudflared started, PID={proc.pid}, parsing URL...") url = await self._parse_cloudflare_url(proc, log_path, log_offset) logger.info(f"[AGENT_APP] cloudflared URL parse result: {url}") diff --git a/app/agent_app/runner.py b/app/agent_app/runner.py index 0dd276c8..7f9de53c 100644 --- a/app/agent_app/runner.py +++ b/app/agent_app/runner.py @@ -22,6 +22,7 @@ from app import node_runtime from app.node_runtime import MIN_NODE_MAJOR +from app.process_ledger import ROLE_AGENT_APP, get_ledger logger = logging.getLogger(__name__) @@ -505,6 +506,14 @@ async def start( stderr=subprocess.STDOUT, creationflags=subprocess.CREATE_NO_WINDOW if sys.platform == "win32" else 0, ) + # Recorded so cleanup can later kill exactly this process (and a + # restart can reap it) without guessing from the port. + get_ledger().register( + process.pid, + ROLE_AGENT_APP, + owner=project_dir.name, + label=f"pocketbase {app_env} :{port}", + ) logger.info( f"[AGENT_APP] started PocketBase pid={process.pid} port={port} " f"env={app_env}" diff --git a/app/process_ledger.py b/app/process_ledger.py new file mode 100644 index 00000000..10b5e27e --- /dev/null +++ b/app/process_ledger.py @@ -0,0 +1,433 @@ +"""Ledger of the processes CraftBot itself started — the ONLY ones it kills. + +Every process-cleanup path used to decide "is this ours?" from a name or a +port: `pkill -f cloudflared` / `Stop-Process -Name cloudflared` took down +every tunnel on the machine, and `f":{port}" in netstat_line` matched :3100 +against :31000 and killed whatever listened there. A foreign service that +happened to sit on "our" port was fair game too. + +The rule this module enforces: **a process is ours only if we recorded its +exact identity when we started it.** Identity is (pid, create_time) — the +creation time is what makes a persisted pid safe across a restart, because a +recycled pid belongs to a process with a different start time and is never +matched. Anything the ledger cannot vouch for is left alone and logged. + +A process counts as owned when it, or a live ancestor, is in the ledger — so +recording a shell (`shell=True` → cmd.exe/sh) also covers the server it +spawned. A server whose spawning shell later dies is covered by recording +the listener itself once the service is up (`adopt_listeners`). + +Each ledger file has ONE writer process (the launcher and the agent keep +separate files), so there is no cross-process write race; a lock covers +threads within the process. +""" + +from __future__ import annotations + +import json +import os +import subprocess +import threading +from dataclasses import asdict, dataclass +from pathlib import Path +from typing import Dict, List, Optional + +try: + from loguru import logger +except ImportError: # pragma: no cover - launcher may run before deps load + import logging + + logger = logging.getLogger(__name__) + +try: + import psutil +except ImportError: # pragma: no cover - declared in requirements.txt + psutil = None + +# psutil derives create_time from boot time + ticks on Linux, which can wobble +# by a fraction of a second between reads. A recycled pid would have to be +# started within this window of the original to be mistaken for it. +_CREATE_TIME_TOLERANCE = 1.0 + +ROLE_TUNNEL = "tunnel" +ROLE_AGENT_APP = "agent_app" +ROLE_LAUNCHER = "launcher" + + +@dataclass +class OwnedProcess: + pid: int + create_time: float + role: str + owner: str = "" # e.g. the project id the process serves + label: str = "" # human-readable, for logs only + + @classmethod + def from_dict(cls, d: Dict) -> "OwnedProcess": + return cls( + pid=int(d["pid"]), + create_time=float(d["create_time"]), + role=str(d.get("role", "")), + owner=str(d.get("owner", "")), + label=str(d.get("label", "")), + ) + + +def _create_time(pid: int) -> Optional[float]: + """Start time of a live pid, or None if it is gone/inaccessible.""" + if psutil is None or not pid or pid <= 0: + return None + try: + return psutil.Process(pid).create_time() + except Exception: + return None + + +def listening_pids(port: int) -> List[int]: + """Pids with a TCP socket LISTENING on exactly `port`. + + Exact integer comparison on the parsed local port — never a substring of + a netstat line. psutil first (locale-independent); netstat/lsof as the + fallback where psutil is missing or needs root (macOS). + """ + port = int(port) + if psutil is not None: + try: + return sorted( + { + c.pid + for c in psutil.net_connections(kind="tcp") + if c.pid + and c.status == psutil.CONN_LISTEN + and c.laddr + and c.laddr.port == port + } + ) + except Exception: + pass # AccessDenied on macOS without root — fall through + + pids = set() + try: + if os.name == "nt": + out = subprocess.run( + ["netstat", "-ano", "-p", "TCP"], + capture_output=True, + text=True, + timeout=10, + creationflags=getattr(subprocess, "CREATE_NO_WINDOW", 0), + ).stdout + for line in out.splitlines(): + # Proto Local Foreign State PID. A listener's foreign + # address is the wildcard (0.0.0.0:0 / [::]:0); matching on + # that keeps this independent of the localized state name. + parts = line.split() + if len(parts) < 5 or parts[0].upper() != "TCP": + continue + local, foreign, pid = parts[1], parts[2], parts[-1] + if not foreign.endswith(":0") or not pid.isdigit(): + continue + try: + if int(local.rsplit(":", 1)[1]) == port: + pids.add(int(pid)) + except (IndexError, ValueError): + continue + else: + out = subprocess.run( + ["lsof", "-nP", f"-iTCP:{port}", "-sTCP:LISTEN", "-t"], + capture_output=True, + text=True, + timeout=10, + ).stdout + pids = {int(p) for p in out.split() if p.isdigit()} + except Exception as e: + logger.warning(f"[PROCESS_LEDGER] could not list listeners on {port}: {e}") + pids.discard(0) + return sorted(pids) + + +def kill_tree(pid: int, grace: float = 5.0) -> None: + """Stop a process and every descendant: SIGTERM, then SIGKILL whatever is + left after `grace` seconds (Windows has no graceful signal for a windowless + tree, so it is `taskkill /T /F` as before). Callers must have verified + ownership first — this does not check.""" + if os.name == "nt": + # taskkill /T walks the tree even when psutil is unavailable. + subprocess.run( + ["taskkill", "/T", "/F", "/PID", str(pid)], + capture_output=True, + timeout=15, + creationflags=getattr(subprocess, "CREATE_NO_WINDOW", 0), + ) + return + if psutil is None: + import signal + + try: + os.kill(pid, signal.SIGTERM) + except OSError: + pass + return + try: + root = psutil.Process(pid) + procs = root.children(recursive=True) + [root] + except Exception: + return + for p in procs: + try: + p.terminate() + except Exception: + pass + try: + _, alive = psutil.wait_procs(procs, timeout=grace) + except Exception: + alive = procs + for p in alive: + try: + p.kill() + except Exception: + pass + + +class ProcessLedger: + """Persisted set of processes this CraftBot process started.""" + + def __init__(self, path: Path) -> None: + self._path = Path(path) + self._lock = threading.RLock() + self._by_pid: Dict[int, OwnedProcess] = {} + self._load() + + # ── recording ───────────────────────────────────────────────────────── + def register( + self, pid: Optional[int], role: str, owner: str = "", label: str = "" + ) -> Optional[OwnedProcess]: + """Record a process we just started. Returns None (and records + nothing) if its identity cannot be read — it then stays unkillable + by any cleanup path, which is the safe failure.""" + if not pid: + return None + ct = _create_time(int(pid)) + if ct is None: + logger.warning( + f"[PROCESS_LEDGER] cannot record pid {pid} ({role} {label}): " + f"identity unreadable{' (psutil missing)' if psutil is None else ''}" + ) + return None + entry = OwnedProcess(int(pid), ct, role, owner, label) + with self._lock: + self._by_pid[entry.pid] = entry + self._save() + return entry + + def adopt_listeners( + self, + port: int, + role: str, + owner: str = "", + label: str = "", + *, + include_self: bool = False, + ) -> List[OwnedProcess]: + """Record the listeners on `port` that are already ours — descendants + of a recorded process — as entries of their own. Covers a server whose + spawning shell later dies, which the ancestor walk in `owner_of` could + then no longer connect to its entry. + + A listener with no recorded ancestor is NOT adopted: a foreign service + that already held the port would otherwise pass our readiness check + and be killed by the next run. `include_self` lets the frozen launcher + (which serves its ports in-process) record its own pid. + """ + adopted = [] + for pid in listening_pids(port): + if pid == os.getpid(): + if not include_self: + continue # never record CraftBot's own process as killable + elif self.owner_of(pid) is None: + continue + entry = self.register(pid, role, owner, label or f"listener :{port}") + if entry: + adopted.append(entry) + return adopted + + def forget(self, pid: Optional[int]) -> None: + if not pid: + return + with self._lock: + if self._by_pid.pop(int(pid), None) is not None: + self._save() + + def forget_owner(self, owner: str, role: Optional[str] = None) -> None: + with self._lock: + before = len(self._by_pid) + self._by_pid = { + pid: e + for pid, e in self._by_pid.items() + if not (e.owner == owner and (role is None or e.role == role)) + } + if len(self._by_pid) != before: + self._save() + + # ── identity ────────────────────────────────────────────────────────── + def _live_entry(self, pid: int) -> Optional[OwnedProcess]: + """The entry for `pid` if that pid is still the SAME process.""" + entry = self._by_pid.get(pid) + if entry is None: + return None + ct = _create_time(pid) + if ct is None or abs(ct - entry.create_time) > _CREATE_TIME_TOLERANCE: + return None + return entry + + def owner_of(self, pid: int) -> Optional[OwnedProcess]: + """The ledger entry that vouches for `pid`: the pid itself, or its + nearest live ancestor that we recorded. None means not ours.""" + with self._lock: + entry = self._live_entry(int(pid)) + if entry is not None: + return entry + if psutil is None: + return None + try: + parents = psutil.Process(int(pid)).parents() + except Exception: + return None + for parent in parents: + entry = self._live_entry(parent.pid) + if entry is not None: + return entry + return None + + # ── killing (owned processes only) ──────────────────────────────────── + def kill_entry(self, entry: OwnedProcess) -> bool: + """Kill a recorded process tree if it is still the same process. + The entry is dropped either way (a dead/recycled pid is stale).""" + with self._lock: + live = self._live_entry(entry.pid) + self._by_pid.pop(entry.pid, None) + self._save() + if live is None: + return False + if live.pid == os.getpid(): + return False # a frozen launcher records itself; never suicide + kill_tree(live.pid) + logger.info( + f"[PROCESS_LEDGER] killed owned {live.role} pid {live.pid} " + f"({live.owner or '-'} {live.label})" + ) + return True + + def kill_pid(self, pid: Optional[int]) -> bool: + """Kill a pid we hold only as a number (e.g. from a persisted record) + — only if the ledger recorded exactly that process.""" + if not pid: + return False + with self._lock: + entry = self._by_pid.get(int(pid)) + if entry is None: + logger.warning( + f"[PROCESS_LEDGER] refusing to kill pid {pid}: not a process " + f"CraftBot recorded starting" + ) + return False + return self.kill_entry(entry) + + def kill_port_listeners(self, port: int) -> bool: + """Free `port` by killing ONLY listeners we own. A listener we cannot + vouch for is logged and left running. True if anything was killed.""" + killed = False + for pid in listening_pids(port): + if pid == os.getpid(): + continue # e.g. the in-process A2App proxy holding the port + entry = self.owner_of(pid) + if entry is None: + logger.warning( + f"[PROCESS_LEDGER] port {port} is held by pid {pid}, which " + f"CraftBot did not start — leaving it alone" + ) + continue + if entry.pid == os.getpid(): + continue # a child of THIS process; its owner stops it + # Kill from the recorded root so the shell that spawned the + # listener goes too. + killed = self.kill_entry(entry) or killed + return killed + + def reap(self, role: Optional[str] = None, owner: Optional[str] = None) -> int: + """Kill every still-alive recorded process (optionally filtered) — + leftovers from a previous run. Returns how many were killed.""" + with self._lock: + targets = [ + e + for e in self._by_pid.values() + if (role is None or e.role == role) + and (owner is None or e.owner == owner) + ] + return sum(1 for e in targets if self.kill_entry(e)) + + def entries(self, role: Optional[str] = None) -> List[OwnedProcess]: + with self._lock: + return [e for e in self._by_pid.values() if role is None or e.role == role] + + # ── persistence ─────────────────────────────────────────────────────── + def _load(self) -> None: + try: + raw = json.loads(self._path.read_text(encoding="utf-8")) + except Exception: + return + for d in raw.get("processes", []): + try: + entry = OwnedProcess.from_dict(d) + except Exception: + continue + self._by_pid[entry.pid] = entry + # Drop records whose process is gone (or whose pid was recycled) so + # the file only ever lists processes that could still be ours. + if psutil is not None: + stale = [pid for pid in self._by_pid if self._live_entry(pid) is None] + for pid in stale: + del self._by_pid[pid] + if stale: + self._save() + + def _save(self) -> None: + try: + self._path.parent.mkdir(parents=True, exist_ok=True) + payload = {"processes": [asdict(e) for e in self._by_pid.values()]} + tmp = self._path.with_suffix(self._path.suffix + ".tmp") + tmp.write_text(json.dumps(payload, indent=2) + "\n", encoding="utf-8") + os.replace(tmp, self._path) + except Exception as e: + logger.error(f"[PROCESS_LEDGER] could not persist {self._path}: {e}") + + +_ledgers: Dict[str, ProcessLedger] = {} +_ledgers_lock = threading.Lock() + + +def get_ledger(scope: str = "agent") -> ProcessLedger: + """The process-wide ledger for `scope`. One file per writer process: + "agent" (the CraftBot agent: Agent App servers, tunnels) and "launcher" + (run.py: frontend + agent backend).""" + with _ledgers_lock: + ledger = _ledgers.get(scope) + if ledger is None: + from app.config import AGENT_WORKSPACE_ROOT + + ledger = ProcessLedger( + Path(AGENT_WORKSPACE_ROOT) / f"owned_processes.{scope}.json" + ) + _ledgers[scope] = ledger + return ledger + + +__all__ = [ + "ROLE_TUNNEL", + "ROLE_AGENT_APP", + "ROLE_LAUNCHER", + "OwnedProcess", + "ProcessLedger", + "get_ledger", + "kill_tree", + "listening_pids", +] diff --git a/app/ui_layer/local_llm_setup.py b/app/ui_layer/local_llm_setup.py index 7c477fbe..3fb27595 100644 --- a/app/ui_layer/local_llm_setup.py +++ b/app/ui_layer/local_llm_setup.py @@ -11,7 +11,7 @@ import subprocess import urllib.error import urllib.request -from typing import Any, Callable, Dict +from typing import Any, Callable, Dict, Set, Tuple logger = logging.getLogger(__name__) @@ -260,6 +260,36 @@ def test_ollama_connection_sync(url: str) -> Dict[str, Any]: return {"success": False, "error": str(exc)} +_OLLAMA_TRAY_EXE = "ollama app.exe" + + +def _tray_app_snapshot() -> Set[Tuple[int, float]]: + """(pid, start time) of every running Ollama tray app.""" + try: + import psutil + except ImportError: + return set() + found = set() + for proc in psutil.process_iter(["name", "create_time"]): + name = (proc.info.get("name") or "").lower() + if name == _OLLAMA_TRAY_EXE: + found.add((proc.pid, proc.info.get("create_time") or 0.0)) + return found + + +def _stop_tray_apps_started_since(before: Set[Tuple[int, float]]) -> None: + """Close only the tray apps that appeared during our install — the one + the installer launched. Without psutil nothing can be told apart, so + nothing is closed (the tray app is harmless, just redundant).""" + from app.process_ledger import kill_tree + + for pid, _ in _tray_app_snapshot() - before: + try: + kill_tree(pid) + except Exception as exc: + logger.warning(f"Could not close Ollama tray app {pid}: {exc}") + + async def install_ollama(progress_callback: Callable) -> Dict[str, Any]: """Install Ollama for the current platform, streaming progress via callback.""" system = platform.system() @@ -268,6 +298,11 @@ async def install_ollama(progress_callback: Callable) -> Dict[str, Any]: if system == "Windows": # Try winget first await progress_callback("Checking for winget...") + # The installer auto-launches the Ollama tray app; we close THAT + # one afterwards. Snapshot first so a tray app the user already + # had running is never touched (this used to be + # `taskkill /IM "ollama app.exe"`, which killed every instance). + trays_before = await asyncio.to_thread(_tray_app_snapshot) try: proc = await asyncio.create_subprocess_exec( "winget", @@ -305,11 +340,7 @@ async def _stream_winget(stream: asyncio.StreamReader) -> None: # Verify actual install regardless of exit code — winget can return non-zero on success if (await asyncio.to_thread(get_ollama_status))["installed"]: - await asyncio.to_thread( - subprocess.run, - ["taskkill", "/F", "/IM", "ollama app.exe", "/T"], - capture_output=True, - ) + await asyncio.to_thread(_stop_tray_apps_started_since, trays_before) await progress_callback("Ollama installed successfully!") return {"success": True, "message": "Ollama installed via winget"} await progress_callback( @@ -363,11 +394,7 @@ async def _stream_ps(stream: asyncio.StreamReader) -> None: ) await run_proc.communicate() if (await asyncio.to_thread(get_ollama_status))["installed"]: - await asyncio.to_thread( - subprocess.run, - ["taskkill", "/F", "/IM", "ollama app.exe", "/T"], - capture_output=True, - ) + await asyncio.to_thread(_stop_tray_apps_started_since, trays_before) await progress_callback("Ollama installed successfully!") return {"success": True, "message": "Ollama installed"} return { diff --git a/craftbot.py b/craftbot.py index 17e79d0c..2ff81ea9 100644 --- a/craftbot.py +++ b/craftbot.py @@ -322,18 +322,65 @@ def _python_exe() -> str: return python -def _read_pid() -> Optional[int]: - """Read PID from the PID file. Returns None if file missing or invalid.""" +# A pid alone is not an identity: once CraftBot exits (crash, reboot) the OS +# reuses the number, and `stop` would kill whatever holds it now — with /T, +# its whole tree. The PID file therefore also records the process's start +# time, and a pid is only treated as CraftBot when both still match. +_START_TIME_TOLERANCE_S = 2.0 # macOS `ps lstart` has 1s resolution + + +def _read_pid_record() -> Optional[tuple]: + """(pid, recorded start time or None, file mtime), or None if no file. + Accepts the legacy bare-integer format (no start time).""" try: with open(PID_FILE) as f: - return int(f.read().strip()) - except (FileNotFoundError, ValueError): + raw = f.read().strip() + mtime = os.path.getmtime(PID_FILE) + except OSError: + return None + try: + if raw.startswith("{"): + import json + + data = json.loads(raw) + started = data.get("started") + return int(data["pid"]), (float(started) if started else None), mtime + return int(raw), None, mtime + except (ValueError, KeyError, TypeError): + return None + + +def _owned_pid() -> Optional[int]: + """The pid in the PID file IF that process is still the CraftBot we + started; otherwise None. A stale record (process gone, or pid reused by + another program) is removed so nothing ever acts on it.""" + record = _read_pid_record() + if record is None: return None + pid, started, mtime = record + actual = _helpers.process_start_time(pid) + if actual is None: + _remove_pid() + return None + if started is not None: + ours = abs(actual - started) <= _START_TIME_TOLERANCE_S + else: + # Legacy file (or start time unreadable at spawn): the file was + # written right after the spawn, so our process started before it. + # A reused pid belongs to a process started after ours died — i.e. + # after the file was written. + ours = actual <= mtime + _START_TIME_TOLERANCE_S + if not ours: + _remove_pid() + return None + return pid def _write_pid(pid: int) -> None: + import json + with open(PID_FILE, "w") as f: - f.write(str(pid)) + json.dump({"pid": pid, "started": _helpers.process_start_time(pid)}, f) def _remove_pid() -> None: @@ -414,23 +461,7 @@ def _wait_for_startup_exit( def _is_running(pid: int) -> bool: """Return True if a process with the given PID is currently alive.""" - if _PLATFORM == "win32": - try: - result = subprocess.run( - ["tasklist", "/FI", f"PID eq {pid}", "/NH"], - capture_output=True, - text=True, - timeout=5, - ) - return str(pid) in result.stdout - except Exception: - return False - else: - try: - os.kill(pid, 0) - return True - except (ProcessLookupError, PermissionError): - return False + return _helpers.process_start_time(pid) is not None def _stop_running_agent_if_alive(grace_s: float = 1.0) -> bool: @@ -439,8 +470,7 @@ def _stop_running_agent_if_alive(grace_s: float = 1.0) -> bool: install/repair so reinstall over a running agent doesn't fail with "Permission denied" on Windows. Returns True if a process was stopped. """ - pid = _read_pid() - if pid and _is_running(pid): + if _owned_pid(): cmd_stop() time.sleep(grace_s) return True @@ -538,8 +568,7 @@ def cmd_start(extra_args: List[str]) -> bool: Returns True once the service survives the early startup check; False when launch fails before CraftBot can be used. """ - pid = _read_pid() - if pid and _is_running(pid): + if _owned_pid(): cmd_stop() # service_mode=False — don't suppress the browser; we open it ourselves below @@ -636,14 +665,13 @@ def cmd_start(extra_args: List[str]) -> bool: def cmd_stop() -> None: """Stop the running CraftBot service.""" - pid = _read_pid() + had_record = _read_pid_record() is not None + pid = _owned_pid() # verified by start time; a reused pid is never ours if pid is None: - print("CraftBot does not appear to be running (no PID file found).") - return - - if not _is_running(pid): - print(f" {DIM}▸ PID {pid} not running — cleaning up stale PID file{RESET}") - _remove_pid() + if had_record: + print(f" {DIM}▸ CraftBot not running — cleaned up stale PID file{RESET}") + else: + print("CraftBot does not appear to be running (no PID file found).") return print(f" {ORANGE}▸{RESET} {WHITE}STOPPING CRAFTBOT{RESET} {DIM}PID {pid}{RESET}") @@ -681,9 +709,9 @@ def cmd_stop() -> None: def cmd_status() -> None: """Print whether CraftBot is currently running and whether auto-start is installed.""" W = max(50, len(LOG_FILE) + 12) - pid = _read_pid() + pid = _owned_pid() print(f"\n{ORANGE}╔{'═' * W}╗{RESET}") - if pid and _is_running(pid): + if pid: print( f"{ORANGE}║{RESET} {GREEN}▸ RUNNING{RESET} {DIM}PID {pid}{RESET}{' ' * (W - 17 - len(str(pid)))}{ORANGE}║{RESET}" ) @@ -691,8 +719,6 @@ def cmd_status() -> None: f"{ORANGE}║{RESET} {DIM}░░ LOG: {LOG_FILE[: W - 8]}{RESET}{' ' * max(0, W - 10 - len(LOG_FILE[: W - 8]))}{ORANGE}║{RESET}" ) else: - if pid: - _remove_pid() print( f"{ORANGE}║{RESET} {RED}▸ NOT RUNNING{RESET}{' ' * (W - 15)}{ORANGE}║{RESET}" ) diff --git a/environment.yml b/environment.yml index 51fdc147..de2bc5b8 100644 --- a/environment.yml +++ b/environment.yml @@ -8,6 +8,7 @@ dependencies: - requests=2.32.5 - pyyaml=6.0.3 - loguru=0.7.3 + - psutil - nest-asyncio=1.6.0 - pymongo=4.16.0 - tzlocal=5.3.1 diff --git a/installer/helpers.py b/installer/helpers.py index 6dfbe8e5..1e55b9ec 100644 --- a/installer/helpers.py +++ b/installer/helpers.py @@ -10,12 +10,20 @@ sys.platform — replaces the if win/elif darwin/else trinity that appears in `_full_install_frozen`, `cmd_uninstall`, `cmd_install`, `cmd_repair`, `_remove_desktop_shortcut`, and `_is_installed`. + +`process_start_time()` reads when a pid's process started, stdlib only (the +frozen installer has no psutil). A pid alone is not an identity: after the +process exits the OS hands the number to something else, so anything that +kills a pid it saved earlier must also check the start time. """ from __future__ import annotations +import os +import subprocess import sys -from typing import TypeVar +import time +from typing import Optional, TypeVar _PLATFORM = sys.platform T = TypeVar("T") @@ -53,3 +61,81 @@ def dispatch_per_platform(*, win: T, mac: T, linux: T) -> T: if _PLATFORM == "darwin": return mac return linux + + +def process_start_time(pid: int) -> Optional[float]: + """Epoch seconds at which `pid` started, or None if no such live process + (or it can't be inspected). Precision: sub-second on Windows/Linux, one + second on macOS — compare with a tolerance.""" + if not pid or pid <= 0: + return None + try: + return dispatch_per_platform( + win=_start_time_windows, mac=_start_time_macos, linux=_start_time_linux + )(int(pid)) + except Exception: + return None + + +def _start_time_windows(pid: int) -> Optional[float]: + import ctypes + from ctypes import wintypes + + PROCESS_QUERY_LIMITED_INFORMATION = 0x1000 + STILL_ACTIVE = 259 + k32 = ctypes.WinDLL("kernel32", use_last_error=True) + k32.OpenProcess.restype = wintypes.HANDLE + handle = k32.OpenProcess(PROCESS_QUERY_LIMITED_INFORMATION, False, pid) + if not handle: + return None + try: + code = wintypes.DWORD() + if not k32.GetExitCodeProcess(handle, ctypes.byref(code)): + return None + if code.value != STILL_ACTIVE: + return None # exited; the handle only keeps the pid reserved + created, exited, kernel, user = (wintypes.FILETIME() for _ in range(4)) + if not k32.GetProcessTimes( + handle, + ctypes.byref(created), + ctypes.byref(exited), + ctypes.byref(kernel), + ctypes.byref(user), + ): + return None + ticks = (created.dwHighDateTime << 32) | created.dwLowDateTime + return ticks / 1e7 - 11644473600 # FILETIME epoch 1601 → 1970 + finally: + k32.CloseHandle(handle) + + +def _start_time_linux(pid: int) -> Optional[float]: + try: + with open(f"/proc/{pid}/stat", encoding="utf-8") as f: + stat = f.read() + # comm (field 2) may contain spaces/parens; fields resume after ")". + fields = stat[stat.rindex(")") + 2 :].split() + if fields[0] == "Z": + return None # zombie: already exited + start_ticks = int(fields[19]) # field 22 overall + with open("/proc/stat", encoding="utf-8") as f: + btime = next(int(l.split()[1]) for l in f if l.startswith("btime")) + return btime + start_ticks / os.sysconf("SC_CLK_TCK") + except Exception: + return None + + +def _start_time_macos(pid: int) -> Optional[float]: + out = subprocess.run( + ["ps", "-o", "stat=,lstart=", "-p", str(pid)], + capture_output=True, + text=True, + timeout=5, + env={**os.environ, "LC_ALL": "C"}, + ).stdout.strip() + if not out: + return None + state, lstart = out.split(None, 1) + if state.startswith("Z"): + return None + return time.mktime(time.strptime(lstart.strip(), "%a %b %d %H:%M:%S %Y")) diff --git a/main.py b/main.py index 25e432bc..29189f87 100644 --- a/main.py +++ b/main.py @@ -7,7 +7,6 @@ import socket import signal import threading -import shutil # Needed for lsof check on Linux/macOS # --- CONFIGURATION --- # Path to the directory containing the docker-compose.yml file @@ -18,8 +17,6 @@ READY_HOST = "localhost" READY_PORT = 3001 MAX_WAIT_SECONDS = 60 -# Port to clean up at the very end -CLEANUP_PORT = 7861 # --------------------- @@ -67,155 +64,6 @@ def is_port_open(host: str, port: int, timeout: int = 1) -> bool: return False -def kill_process_on_port(port: int): - """Finds and kills any process listening on the specified TCP port (Cross-platform).""" - current_os = platform.system() - port_str = str(port) - print(f"[*] Checking for leftover processes on port {port}...") - - try: - if current_os == "Windows": - # SECURITY FIX: Use list-based subprocess call instead of shell=True - # This prevents command injection vulnerabilities - try: - # Use netstat without shell pipes - safer approach - output = subprocess.check_output( - ["netstat", "-ano"], text=True, stderr=subprocess.DEVNULL - ) - pids_to_kill = set() - for line in output.strip().split("\n"): - parts = line.strip().split() - # Format: PROTO LOCAL_ADDR FOREIGN_ADDR STATE PID - if len(parts) >= 5 and "LISTENING" in line and parts[-1].isdigit(): - pid = parts[-1] - try: - pid_int = int(pid) - if pid_int > 0: - pids_to_kill.add(pid) - except ValueError: - continue - - if not pids_to_kill: - print(f"[*] Port {port} is free.") - return - - for pid in pids_to_kill: - print( - f"[!] Found stale process (PID: {pid}) on port {port}. Killing it..." - ) - # SECURITY FIX: Use list-based call instead of f-string with shell=True - try: - subprocess.run( - ["taskkill", "/F", "/T", "/PID", pid], - stdout=subprocess.DEVNULL, - stderr=subprocess.DEVNULL, - timeout=5, - ) - except subprocess.TimeoutExpired: - print(f"[!] Timeout killing PID {pid}") - except Exception as e: - print(f"[!] Error killing PID {pid}: {e}") - print(f"[*] Port {port} cleared.") - time.sleep(0.5) - except subprocess.CalledProcessError: - print(f"[*] Port {port} is free.") - - else: # Linux/macOS - find_cmd = ["lsof", "-t", "-i", f"TCP:{port_str}"] - if shutil.which("lsof"): - try: - output = subprocess.check_output( - find_cmd, text=True, stderr=subprocess.DEVNULL - ) - pids = [ - p - for p in output.strip().split("\n") - if p.isdigit() and int(p) > 0 - ] - if not pids: - print(f"[*] Port {port} is free.") - return - for pid in pids: - print( - f"[!] Found stale process (PID: {pid}) on port {port}. Killing it..." - ) - subprocess.run( - ["kill", "-9", pid], - stdout=subprocess.DEVNULL, - stderr=subprocess.DEVNULL, - ) - print(f"[*] Port {port} cleared.") - time.sleep(0.5) - except subprocess.CalledProcessError: - print(f"[*] Port {port} is free.") - else: - print( - f"[!] Warning: 'lsof' not found. Cannot automatically clean port {port}." - ) - - except Exception as e: - print(f"[!] Warning: Failed to clean up port {port}: {e}") - - -def kill_process_on_port_quiet(port: int): - """Quietly kill any process listening on the specified TCP port.""" - current_os = platform.system() - port_str = str(port) - - try: - if current_os == "Windows": - # SECURITY FIX: Use list-based subprocess call instead of shell=True - try: - output = subprocess.check_output( - ["netstat", "-ano"], text=True, stderr=subprocess.DEVNULL - ) - pids_to_kill = set() - for line in output.strip().split("\n"): - parts = line.strip().split() - if len(parts) >= 5 and "LISTENING" in line and parts[-1].isdigit(): - try: - pid_int = int(parts[-1]) - if pid_int > 0: - pids_to_kill.add(parts[-1]) - except ValueError: - continue - - for pid in pids_to_kill: - try: - subprocess.run( - ["taskkill", "/F", "/T", "/PID", pid], - stdout=subprocess.DEVNULL, - stderr=subprocess.DEVNULL, - timeout=5, - ) - except (subprocess.TimeoutExpired, Exception): - pass - except subprocess.CalledProcessError: - pass - else: - find_cmd = ["lsof", "-t", "-i", f"TCP:{port_str}"] - if shutil.which("lsof"): - try: - output = subprocess.check_output( - find_cmd, text=True, stderr=subprocess.DEVNULL - ) - pids = [ - p - for p in output.strip().split("\n") - if p.isdigit() and int(p) > 0 - ] - for pid in pids: - subprocess.run( - ["kill", "-9", pid], - stdout=subprocess.DEVNULL, - stderr=subprocess.DEVNULL, - ) - except subprocess.CalledProcessError: - pass - except Exception: - pass - - # --- MAIN LOGIC --- @@ -329,9 +177,9 @@ def main(): run_command(["docker", "compose", "down"], cwd=VM_DIR, check=False) except Exception as e: print(f"[!] Warning: Error during docker shutdown: {e}") - - # 2. Clean up ports - kill_process_on_port(CLEANUP_PORT) + # No port sweep afterwards: `compose down` stops exactly the + # containers we started. Killing whatever else listens on the + # port would hit processes CraftBot never launched. else: print("[*] Skipping Docker cleanup (not started in CLI mode).") diff --git a/requirements.txt b/requirements.txt index c844e9f0..8010a92b 100644 --- a/requirements.txt +++ b/requirements.txt @@ -1,6 +1,7 @@ requests pyyaml loguru +psutil nest-asyncio pymongo tzlocal diff --git a/run.py b/run.py index 9e72cbaf..2dc250b9 100644 --- a/run.py +++ b/run.py @@ -349,59 +349,68 @@ def cleanup_background_processes(): atexit.register(cleanup_background_processes) -def _kill_stale_port_process(port: int) -> bool: - """Kill any process listening on the given port (stale leftovers from previous runs). +def _launcher_ledger(): + """Processes this launcher started (frontend, agent backend), recorded by + pid + start time so a later run can free our ports without touching + anything else. None if the ledger cannot load — then nothing is killed.""" + try: + from app.process_ledger import get_ledger - Returns True if a stale process was found and killed. - """ - if sys.platform != "win32": - try: - result = subprocess.run( - ["lsof", "-ti", f":{port}"], - capture_output=True, - text=True, - timeout=5, - ) - for pid_str in result.stdout.strip().split(): - pid = int(pid_str) - if pid != os.getpid(): - subprocess.run(["kill", "-9", str(pid)], timeout=5) - return True - except Exception: - pass - return False + return get_ledger("launcher") + except Exception: + return None - # Windows: parse netstat to find the PID, then taskkill it - try: - result = subprocess.run( - ["netstat", "-ano"], - capture_output=True, - text=True, - timeout=10, + +def _record_launched(process, label: str) -> None: + ledger = _launcher_ledger() + pid = getattr(process, "pid", None) + if ledger is not None and pid: + from app.process_ledger import ROLE_LAUNCHER + + ledger.register(pid, ROLE_LAUNCHER, label=label) + + +def _adopt_port_listeners(*ports: int) -> None: + """Record the servers on our ports once they are up — only ones descended + from a process we launched, plus this launcher itself (it serves the + static frontend in-process, and in frozen mode the agent too). That is what lets the NEXT run + free a port a crashed run left bound, without ever adopting a foreign + service that happened to answer on it.""" + ledger = _launcher_ledger() + if ledger is None: + return + from app.process_ledger import ROLE_LAUNCHER + + for port in ports: + ledger.adopt_listeners( + port, + ROLE_LAUNCHER, + label=f"listener :{port}", + # This process serves the static frontend in-process (and, when + # frozen, the agent too), so it is itself a port holder. + include_self=True, ) - for line in result.stdout.splitlines(): - # Match LISTENING lines for our port on any address - if f":{port}" in line and "LISTENING" in line: - parts = line.split() - pid = int(parts[-1]) - if pid and pid != os.getpid(): - subprocess.run( - ["taskkill", "/PID", str(pid), "/F"], - capture_output=True, - timeout=10, - ) - return True - except Exception: - pass - return False def _free_ports(*ports: int) -> None: - """Kill stale processes on the given ports before startup.""" + """Free our ports of leftovers from a previous run — only processes that + run recorded starting. Anything else on the port is reported, not killed.""" + ledger = _launcher_ledger() + if ledger is None: + return + from app.process_ledger import listening_pids + for port in ports: - if _kill_stale_port_process(port): + if ledger.kill_port_listeners(port): # Give the OS a moment to release the socket time.sleep(0.5) + foreign = [p for p in listening_pids(port) if p != os.getpid()] + if foreign: + print( + f"Warning: port {port} is in use by PID {', '.join(map(str, foreign))}, " + f"which CraftBot did not start. Close it or pick another port " + f"(--frontend-port / --backend-port)." + ) def _launch_static_frontend(silent: bool = False) -> Optional[subprocess.Popen]: @@ -807,6 +816,7 @@ def launch_frontend(silent: bool = False) -> Optional[subprocess.Popen]: ) process = subprocess.Popen(cmd, **popen_kwargs) _background_processes.append(process) + _record_launched(process, "frontend dev server") return process except FileNotFoundError: if not silent: @@ -1100,6 +1110,7 @@ def kill(self): stderr=sys.stderr, ) _background_processes.append(process) + _record_launched(process, "agent backend") return process except Exception as e: if not silent: @@ -1487,6 +1498,11 @@ def launch_agent(env_name: Optional[str], conda_base: Optional[str], use_conda: pass time.sleep(0.5) + _adopt_port_listeners( + *([FRONTEND_PORT] if frontend_ready else []), + *([BACKEND_PORT] if backend_ready else []), + ) + # Small delay to ensure agent's stdout is flushed before we print # The agent prints steps 3-8, and we want them to appear before the ready banner time.sleep(0.3) diff --git a/tests/test_agent_app_undeployed.py b/tests/test_agent_app_undeployed.py index 7a2c47d5..97348196 100644 --- a/tests/test_agent_app_undeployed.py +++ b/tests/test_agent_app_undeployed.py @@ -10,6 +10,7 @@ import asyncio import subprocess +import sys import types from pathlib import Path @@ -96,7 +97,10 @@ def test_an_edit_made_outside_the_action_layer_is_still_seen(self, agent, projec target = Path(project.path) / "frontend" / "src" / "app" / "App.tsx" target.write_text("const a = 1\n") self._shipped(project) - subprocess.run(["sed", "-i", "s/1/2/", str(target)], check=True) + # A separate process rewrites it, as run_shell would. Python rather + # than `sed -i`, which isn't portable (fails on Windows temp dirs). + edit = "import sys,pathlib; p=pathlib.Path(sys.argv[1]); p.write_text(p.read_text().replace('1','2'))" + subprocess.run([sys.executable, "-c", edit, str(target)], check=True) assert agent._unshipped_fingerprint(project) is not None def test_shipping_again_clears_it(self, agent, project): diff --git a/tests/test_pid_identity.py b/tests/test_pid_identity.py new file mode 100644 index 00000000..a354a219 --- /dev/null +++ b/tests/test_pid_identity.py @@ -0,0 +1,139 @@ +"""Saved pids are only acted on while they still name the SAME process. + +A pid alone is not an identity: after a crash or reboot the OS reuses the +number, and a stale PID file / marker would make CraftBot kill (with /T, a +whole tree) some unrelated program. Each check here pairs a real process with +a record that either matches it or looks like a recycled pid. +""" + +import json +import os +import subprocess +import sys +import time + +import pytest + +import craftbot +from agent_core.core.impl.action import cancellation +from app.ui_layer import local_llm_setup +from installer.helpers import process_start_time + + +@pytest.fixture +def sleeper(): + proc = subprocess.Popen([sys.executable, "-c", "import time; time.sleep(120)"]) + yield proc + if proc.poll() is None: + proc.kill() + proc.wait(timeout=10) + + +def _alive(proc) -> bool: + try: + proc.wait(timeout=3) + return False + except subprocess.TimeoutExpired: + return True + + +# ── installer.helpers.process_start_time ────────────────────────────────── + + +def test_start_time_of_live_and_dead_processes(sleeper): + started = process_start_time(sleeper.pid) + assert started is not None and abs(started - time.time()) < 60 + sleeper.kill() + sleeper.wait(timeout=10) + assert process_start_time(sleeper.pid) is None + + +# ── craftbot.py PID file ────────────────────────────────────────────────── + + +@pytest.fixture +def pid_file(tmp_path, monkeypatch): + path = tmp_path / "craftbot.pid" + monkeypatch.setattr(craftbot, "PID_FILE", str(path)) + return path + + +def test_pid_file_written_by_start_is_recognised(pid_file, sleeper): + craftbot._write_pid(sleeper.pid) + assert json.loads(pid_file.read_text())["started"] is not None + assert craftbot._owned_pid() == sleeper.pid + + +def test_reused_pid_is_not_stopped(pid_file, sleeper): + # Same pid, but recorded as starting an hour earlier: a recycled pid. + pid_file.write_text( + json.dumps({"pid": sleeper.pid, "started": process_start_time(sleeper.pid) - 3600}) + ) + craftbot.cmd_stop() + assert _alive(sleeper), "craftbot stop killed a process it did not start" + assert not pid_file.exists(), "the stale record should be cleaned up" + + +def test_legacy_pid_file_older_than_the_process_is_not_ours(pid_file, sleeper): + # A bare-integer file from before this change, written long before this + # pid's current process started — i.e. the number was reused since. + pid_file.write_text(str(sleeper.pid)) + old = time.time() - 3600 + os.utime(pid_file, (old, old)) + assert craftbot._owned_pid() is None + assert _alive(sleeper) + + +def test_legacy_pid_file_written_after_spawn_is_ours(pid_file, sleeper): + time.sleep(0.2) + pid_file.write_text(str(sleeper.pid)) + assert craftbot._owned_pid() == sleeper.pid + + +def test_stop_kills_the_process_we_started(pid_file, sleeper): + craftbot._write_pid(sleeper.pid) + craftbot.cmd_stop() + assert not _alive(sleeper) + + +# ── cancellation markers ────────────────────────────────────────────────── + + +@pytest.fixture +def session(monkeypatch, tmp_path): + monkeypatch.setattr(cancellation, "_marker_dir", lambda sid: tmp_path / sid) + return "sess-1" + + +def test_marked_child_is_killed_on_stop(session, sleeper): + cancellation.mark_subprocess(session, sleeper.pid) + assert cancellation.kill_session_processes(session) == 1 + assert not _alive(sleeper) + + +def test_stale_marker_for_reused_pid_is_ignored(session, sleeper, tmp_path): + marker = tmp_path / session / f"{sleeper.pid}.pid" + marker.parent.mkdir(parents=True) + marker.write_text( + json.dumps({"pid": sleeper.pid, "started": process_start_time(sleeper.pid) - 3600}) + ) + assert cancellation.kill_session_processes(session) == 0 + assert _alive(sleeper), "a stale marker killed an unrelated process" + assert not marker.exists() + + +# ── Ollama tray app ─────────────────────────────────────────────────────── + + +def test_only_the_tray_app_our_install_launched_is_closed(monkeypatch): + users_tray = (111, 1000.0) + installer_tray = (222, 2000.0) + monkeypatch.setattr( + local_llm_setup, "_tray_app_snapshot", lambda: {users_tray, installer_tray} + ) + killed = [] + import app.process_ledger + + monkeypatch.setattr(app.process_ledger, "kill_tree", killed.append) + local_llm_setup._stop_tray_apps_started_since({users_tray}) + assert killed == [222] diff --git a/tests/test_process_ledger.py b/tests/test_process_ledger.py new file mode 100644 index 00000000..d2c7ebc2 --- /dev/null +++ b/tests/test_process_ledger.py @@ -0,0 +1,192 @@ +"""Process cleanup kills only what CraftBot started (app/process_ledger.py). + +Pins the two over-broad kills this replaced: + * `pkill -f cloudflared` / `Stop-Process -Name cloudflared` took down every + cloudflared on the machine, including tunnels CraftBot never started. + * `f":{port}" in netstat_line` matched :3100 against :31000, and any + listener on "our" port was killed whether or not we launched it. + +Real processes, real sockets — no mocks of the OS layer. +""" + +import socket +import subprocess +import sys +import time + +import psutil +import pytest + +from app.process_ledger import ProcessLedger, ROLE_TUNNEL, listening_pids + +# A child that listens on the port given as argv[1] and sleeps. +_LISTENER = ( + "import socket,sys,time\n" + "s=socket.socket(); s.bind(('127.0.0.1', int(sys.argv[1]))); s.listen()\n" + "time.sleep(120)\n" +) + + +def _free_port() -> int: + with socket.socket() as s: + s.bind(("127.0.0.1", 0)) + return s.getsockname()[1] + + +def _spawn_listener(port: int) -> subprocess.Popen: + proc = subprocess.Popen([sys.executable, "-c", _LISTENER, str(port)]) + deadline = time.time() + 15 + while time.time() < deadline: + if proc.pid in listening_pids(port): + return proc + time.sleep(0.1) + proc.kill() + raise RuntimeError(f"listener on {port} never came up") + + +def _spawn_sleeper() -> subprocess.Popen: + return subprocess.Popen([sys.executable, "-c", "import time; time.sleep(120)"]) + + +@pytest.fixture +def procs(): + started = [] + yield started + for p in started: + if p.poll() is None: + p.kill() + p.wait(timeout=10) + + +def _gone(proc: subprocess.Popen) -> bool: + try: + proc.wait(timeout=10) + return True + except subprocess.TimeoutExpired: + return False + + +def test_listening_pids_matches_exact_port_only(procs): + port = _free_port() + proc = _spawn_listener(port) + procs.append(proc) + assert proc.pid in listening_pids(port) + # A port whose digits contain ours (and one that is its prefix) must + # not match — the old substring check would have hit both. + for other in (port * 10 % 65536 or 1, int(str(port)[:-1] or 1)): + if other != port: + assert proc.pid not in listening_pids(other) + + +def test_foreign_listener_on_our_port_is_left_alone(tmp_path, procs): + port = _free_port() + foreign = _spawn_listener(port) # started, but never recorded + procs.append(foreign) + ledger = ProcessLedger(tmp_path / "ledger.json") + + assert ledger.kill_port_listeners(port) is False + assert foreign.poll() is None, "a process CraftBot did not start was killed" + + +def test_owned_listener_is_killed(tmp_path, procs): + port = _free_port() + ours = _spawn_listener(port) + procs.append(ours) + ledger = ProcessLedger(tmp_path / "ledger.json") + ledger.register(ours.pid, "agent_app", owner="p1") + + assert ledger.kill_port_listeners(port) is True + assert _gone(ours) + assert ledger.entries() == [] + + +def test_only_recorded_tunnel_is_reaped(tmp_path, procs): + ours, theirs = _spawn_sleeper(), _spawn_sleeper() + procs += [ours, theirs] + path = tmp_path / "ledger.json" + ProcessLedger(path).register(ours.pid, ROLE_TUNNEL, owner="p1") + + # A fresh ledger (= after a CraftBot restart) reaps from the file. + assert ProcessLedger(path).reap(role=ROLE_TUNNEL) == 1 + assert _gone(ours) + assert theirs.poll() is None, "an unrelated tunnel was killed" + + +def test_recycled_pid_is_never_killed(tmp_path, procs): + victim = _spawn_sleeper() + procs.append(victim) + path = tmp_path / "ledger.json" + ledger = ProcessLedger(path) + entry = ledger.register(victim.pid, ROLE_TUNNEL) + # Same pid, different start time: what a recycled pid looks like. + entry.create_time -= 3600 + ledger._save() + + assert ProcessLedger(path).kill_pid(victim.pid) is False + assert victim.poll() is None, "a recycled pid was treated as ours" + + +def test_descendant_of_owned_shell_is_owned(tmp_path, procs): + port = _free_port() + # shell-style parent: a python that spawns the real listener and waits. + parent = subprocess.Popen( + [ + sys.executable, + "-c", + "import subprocess,sys; subprocess.run([sys.executable,'-c'," + f"{_LISTENER!r},'{port}'])", + ] + ) + procs.append(parent) + deadline = time.time() + 15 + while time.time() < deadline and not listening_pids(port): + time.sleep(0.1) + listener_pid = listening_pids(port)[0] + assert listener_pid != parent.pid + + ledger = ProcessLedger(tmp_path / "ledger.json") + ledger.register(parent.pid, "agent_app") + assert ledger.kill_port_listeners(port) is True + assert _gone(parent) + assert not psutil.pid_exists(listener_pid) or listener_pid not in listening_pids(port) + + +def test_adopt_skips_a_foreign_listener(tmp_path, procs): + # A service that already held the port answers our readiness check; it + # must not become "ours" (the next run would kill it). + port = _free_port() + foreign = _spawn_listener(port) + procs.append(foreign) + ledger = ProcessLedger(tmp_path / "ledger.json") + assert ledger.adopt_listeners(port, "launcher") == [] + assert ledger.kill_port_listeners(port) is False + assert foreign.poll() is None + + +def test_adopted_server_survives_its_shell_dying(tmp_path, procs): + port = _free_port() + parent = subprocess.Popen( + [ + sys.executable, + "-c", + "import subprocess,sys; subprocess.Popen([sys.executable,'-c'," + f"{_LISTENER!r},'{port}']); import time; time.sleep(120)", + ] + ) + procs.append(parent) + deadline = time.time() + 15 + while time.time() < deadline and not listening_pids(port): + time.sleep(0.1) + server_pid = listening_pids(port)[0] + + ledger = ProcessLedger(tmp_path / "ledger.json") + ledger.register(parent.pid, "agent_app") + assert [e.pid for e in ledger.adopt_listeners(port, "agent_app")] == [server_pid] + + parent.kill() # the shell dies; the server is orphaned + parent.wait(timeout=10) + assert ledger.kill_port_listeners(port) is True + deadline = time.time() + 10 + while time.time() < deadline and psutil.pid_exists(server_pid): + time.sleep(0.1) + assert not psutil.pid_exists(server_pid)