fix: preserve executor probe uncertainty

Co-authored-by: Ouroboros <311266734+ouroboros-agent@users.noreply.github.com>
This commit is contained in:
Ouroboros 2026-08-21 23:57:51 +03:00
parent dd10bd51c5
commit a849c9a6b5
4 changed files with 387 additions and 27 deletions

View file

@ -387,6 +387,7 @@ def _docker_exec_pidfile_stop_shell(pidfile: str) -> str:
"kill -TERM -$pid 2>/dev/null || kill -TERM $pid 2>/dev/null || true; "
"sleep 0.5; "
"kill -KILL -$pid 2>/dev/null || kill -KILL $pid 2>/dev/null || true; "
"if kill -0 -$pid 2>/dev/null || kill -0 $pid 2>/dev/null; then exit 1; fi; "
f"rm -f {quoted_pidfile}"
)
@ -412,7 +413,15 @@ def _dispatch_docker_record_cleanup(record: dict[str, Any]) -> bool:
errors="replace",
timeout=5,
)
return proc.returncode == 0
if proc.returncode != 0:
return False
# The bounded shell dispatch is only an attempt. Service records carry
# a backend PID, so panic/``wait=False`` cleanup performs the same
# explicit terminal probe before allowing the durable record to go.
backend_pid = str(record.get("backend_pid") or "").strip()
if backend_pid:
return _docker_pid_state(container, backend_pid) == "exited"
return True
except Exception:
return False
@ -491,7 +500,6 @@ def _register_service_process(drive_root: pathlib.Path | None, record: _Executor
"container_name": record.executor.container_name,
"backend_pid": record.backend_pid,
"backend_log_path": record.backend_log_path,
"readiness": dict(record.readiness),
"backend_cwd": record.backend_cwd,
"host_cwd": str(record.host_cwd),
"cwd_root": record.cwd_root,
@ -651,17 +659,17 @@ def _kill_docker_record(record: dict[str, Any], *, wait: bool = True) -> bool:
return False
if proc.returncode != 0:
return False
if _docker_pid_is_running(container, backend_pid):
if _docker_pid_state(container, backend_pid) != "exited":
return False
return True
def _docker_pid_is_running(container_name: str, backend_pid: str) -> bool:
"""Return true only when Docker's kill-0 probe confirms a live backend."""
def _docker_pid_state(container_name: str, backend_pid: str) -> str:
"""Probe Docker and preserve ``unknown`` when the probe is inconclusive."""
pid = str(backend_pid or "").strip()
if not pid:
return False
return "unknown"
try:
bootstrap_process_path()
proc = subprocess.run(
@ -679,8 +687,18 @@ def _docker_pid_is_running(container_name: str, backend_pid: str) -> bool:
timeout=5,
)
except Exception:
return True
return proc.returncode == 0 and "running" in (proc.stdout or "")
return "unknown"
if proc.returncode != 0:
return "unknown"
raw_output = proc.stdout or ""
if isinstance(raw_output, bytes):
raw_output = raw_output.decode("utf-8", errors="replace")
output = str(raw_output).strip()
if output == "running":
return "running"
if output == "exited":
return "exited"
return "unknown"
def kill_all_foreground(drive_root: pathlib.Path | None = None, *, wait: bool = True) -> list[dict[str, Any]]:
@ -884,6 +902,8 @@ def stop_service(ctx: Any, name: str) -> dict[str, Any] | None:
)
except Exception as exc:
payload = _service_payload(record)
if record.executor.kind == "docker_exec":
payload["cleanup_dispatched"] = False
payload["stop_failed"] = True
payload["stop_error"] = f"{type(exc).__name__}: {exc}"
return payload
@ -895,10 +915,14 @@ def stop_service(ctx: Any, name: str) -> dict[str, Any] | None:
# A successful shell dispatch is not custody settlement by itself.
# Confirm the backend PID is no longer observable before dropping the
# in-memory handle and durable process record.
if _service_state(record) == "running":
probe_state = _safe_service_state(record)
if probe_state != "exited":
payload = _service_payload(record)
payload["stop_failed"] = True
payload["stop_error"] = "docker service stop returned success but kill-0 still reports the service running"
payload["stop_error"] = (
"docker service stop returned success but kill-0 confirmation is "
f"{probe_state}"
)
return payload
with _STATE_LOCK:
_SERVICES.pop(key, None)
@ -911,7 +935,7 @@ def stop_service(ctx: Any, name: str) -> dict[str, Any] | None:
def stop_task_services(ctx: Any) -> list[dict[str, Any]]:
task_id = str(getattr(ctx, "task_id", "") or "manual")
kept = [
_service_payload(record, state=_service_state(record), note="keep_alive")
_service_payload(record, state=_safe_service_state(record), note="keep_alive")
for record in _services_snapshot()
if record.task_id == task_id
and bool(getattr(record, "keep_alive", False))
@ -965,8 +989,9 @@ def kill_all_services(
text=True,
timeout=10 if wait else 5,
)
terminal = proc.returncode == 0 and _service_state(record) != "running"
payload = _service_payload(record, state="stopped" if terminal else _service_state(record))
probe_state = _safe_service_state(record) if proc.returncode == 0 else "unknown"
terminal = proc.returncode == 0 and probe_state == "exited"
payload = _service_payload(record, state="stopped" if terminal else probe_state)
payload["cleanup_dispatched"] = terminal
if terminal:
with _STATE_LOCK:
@ -974,10 +999,16 @@ def kill_all_services(
_forget_process(record.durable_record_path)
else:
payload["stop_failed"] = True
payload["stop_error"] = proc.stderr.strip() or proc.stdout.strip() or "docker service stop failed"
payload["stop_error"] = (
proc.stderr.strip()
or proc.stdout.strip()
or f"docker service stop confirmation is {probe_state}"
)
stopped.append(payload)
except Exception as exc:
payload = _service_payload(record)
if record.executor.kind == "docker_exec":
payload["cleanup_dispatched"] = False
payload["stop_failed"] = True
payload["stop_error"] = f"{type(exc).__name__}: {exc}"
stopped.append(payload)
@ -1100,11 +1131,17 @@ def _docker_service_stop_shell(backend_pid: str) -> str:
def _service_payload(record: _ExecutorService, *, state: str | None = None, note: str = "") -> dict[str, Any]:
actual_state = state or _service_state(record)
actual_state = state if state is not None else _safe_service_state(record)
if actual_state == "running":
_refresh_executor_service_readiness(record)
# Readiness scanning and the state probe are separate operations. A
# process can settle during the scan, so re-poll before serializing a
# running+ready claim (for both local and Docker backends).
actual_state = _safe_service_state(record)
else:
record.ready = False
if actual_state != "running":
record.ready = False
payload = {
"service_id": record.service_id,
"name": record.name,
@ -1139,15 +1176,43 @@ def _service_payload(record: _ExecutorService, *, state: str | None = None, note
def _service_state(record: _ExecutorService) -> str:
if record.executor.kind == "local":
proc = record.local_proc
return "running" if proc is not None and proc.poll() is None else "exited"
proc = subprocess.run(
["docker", "exec", record.executor.container_name, "sh", "-lc", f"kill -0 {shlex.quote(record.backend_pid)} 2>/dev/null && echo running || echo exited"],
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
text=True,
timeout=10,
)
return "running" if "running" in (proc.stdout or "") else "exited"
if proc is None:
return "exited"
try:
return "running" if proc.poll() is None else "exited"
except Exception:
return "unknown"
try:
proc = subprocess.run(
["docker", "exec", record.executor.container_name, "sh", "-lc", f"kill -0 {shlex.quote(record.backend_pid)} 2>/dev/null && echo running || echo exited"],
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
text=True,
timeout=10,
)
except Exception:
return "unknown"
if proc.returncode != 0:
return "unknown"
raw_output = proc.stdout or ""
if isinstance(raw_output, bytes):
raw_output = raw_output.decode("utf-8", errors="replace")
output = str(raw_output).strip()
if output == "running":
return "running"
if output == "exited":
return "exited"
return "unknown"
def _safe_service_state(record: _ExecutorService) -> str:
"""Never let an inconclusive state probe discard cleanup custody."""
try:
state = _service_state(record)
except Exception:
return "unknown"
return state if state in {"running", "exited", "unknown"} else "unknown"
def _read_service_tail(record: _ExecutorService, chars: int) -> str:

View file

@ -280,6 +280,98 @@ def test_executor_local_status_refreshes_late_readiness_after_timeout(tmp_path):
assert payload["ready"] is True
def test_executor_payload_repolls_local_terminal_after_readiness_refresh(tmp_path):
import ouroboros.workspace_executor as workspace_executor
class RacingProcess:
def __init__(self):
self.calls = 0
def poll(self):
self.calls += 1
return None if self.calls == 1 else 0
log_path = tmp_path / "race.log"
log_path.write_bytes(b"READY\n")
record = SimpleNamespace(
service_id="task:race",
name="race",
task_id="task",
executor=SimpleNamespace(executor_id="local", kind="local", network="host"),
backend_pid="1234",
backend_cwd="/workspace",
host_cwd=tmp_path,
cwd_root="active_workspace",
cwd_base=str(tmp_path),
cwd_source="active_workspace",
skill_name="",
cmd=["service"],
outputs=[],
keep_alive=False,
backend_log_path=str(log_path),
readiness={"log_contains": "READY"},
started_at=time.time(),
ready=False,
ready_observed_at="",
local_proc=RacingProcess(),
readiness_log_offset=0,
readiness_log_carry=b"",
readiness_log_identity=None,
)
payload = workspace_executor._service_payload(record)
assert payload["state"] == "exited"
assert payload["ready"] is False
assert payload["ready_observed_at"]
def test_executor_payload_repolls_docker_terminal_after_readiness_refresh(monkeypatch, tmp_path):
import ouroboros.workspace_executor as workspace_executor
record = SimpleNamespace(
service_id="task:docker-race",
name="docker-race",
task_id="task",
executor=SimpleNamespace(executor_id="docker", kind="docker_exec", network="none", container_name="svc"),
backend_pid="1234",
backend_cwd="/workspace",
host_cwd=tmp_path,
cwd_root="active_workspace",
cwd_base=str(tmp_path),
cwd_source="active_workspace",
skill_name="",
cmd=["service"],
outputs=[],
keep_alive=False,
backend_log_path="/tmp/service.log",
readiness={"log_contains": "READY"},
started_at=time.time(),
ready=False,
ready_observed_at="",
readiness_log_offset=0,
readiness_log_carry=b"",
readiness_log_identity=(7, 42),
)
state_calls = 0
def fake_run(cmd, **kwargs):
nonlocal state_calls
shell = str(cmd[-1])
if "stat -c" in shell:
return subprocess.CompletedProcess(cmd, 0, stdout=b"7:42:5\nREADY", stderr=b"")
state_calls += 1
result = "running\n" if state_calls == 1 else "exited\n"
return subprocess.CompletedProcess(cmd, 0, stdout=result, stderr="")
monkeypatch.setattr(workspace_executor.subprocess, "run", fake_run)
payload = workspace_executor._service_payload(record)
assert payload["state"] == "exited"
assert payload["ready"] is False
assert state_calls == 2
def test_executor_docker_readiness_uses_incremental_cursor_and_carry(monkeypatch):
import ouroboros.workspace_executor as workspace_executor

View file

@ -775,6 +775,7 @@ def test_executor_service_status_and_durable_record_redact_secret_like_args(tmp_
assert secret not in status_raw
assert secret not in durable_text
assert '"readiness"' not in durable_text
assert "***REDACTED***" in status_raw
assert "***REDACTED***" in durable_text
@ -1082,6 +1083,170 @@ def test_docker_executor_stop_success_without_terminal_kill_preserves_handle(tmp
assert workspace_executor.service_status(ctx, "svc") is not None
def test_docker_executor_stop_unknown_probe_preserves_handle(tmp_path, monkeypatch):
import ouroboros.workspace_executor as workspace_executor
workspace = tmp_path / "workspace"
workspace.mkdir()
data = tmp_path / "data"
data.mkdir()
ctx = ToolContext(
repo_dir=tmp_path / "repo",
drive_root=data,
workspace_root=workspace,
workspace_mode="external",
task_id="docker-stop-unknown",
executor_ref={
"type": "docker_exec",
"id": "pb-container",
"container_name": "pb-container",
"network": "none",
"workspace_host_path": str(workspace),
"workspace_backend_path": "/workspace",
},
)
workspace_executor._SERVICES.clear()
def fake_run(cmd, **kwargs):
if cmd[:3] == ["docker", "inspect", "-f"]:
return subprocess.CompletedProcess(cmd, 0, stdout="none\n", stderr="")
if cmd[:2] == ["docker", "exec"] and "nohup" in str(cmd[-1]):
return subprocess.CompletedProcess(cmd, 0, stdout="12345\n", stderr="")
if cmd[:2] == ["docker", "exec"] and "kill -TERM" in str(cmd[-1]):
return subprocess.CompletedProcess(cmd, 0, stdout="", stderr="")
if cmd[:2] == ["docker", "exec"] and "kill -0" in str(cmd[-1]):
return subprocess.CompletedProcess(cmd, 7, stdout="", stderr="daemon unavailable")
raise AssertionError(cmd)
monkeypatch.setattr(workspace_executor.subprocess, "run", fake_run)
workspace_executor.start_service(
ctx,
name="svc",
cmd=["sleep", "30"],
host_cwd=workspace,
cwd_root="active_workspace",
readiness={},
outputs=[],
before_outputs={},
)
failed = workspace_executor.stop_service(ctx, "svc")
assert failed and failed["stop_failed"] is True
assert "unknown" in failed["stop_error"]
assert failed["state"] == "unknown"
assert workspace_executor.service_status(ctx, "svc") is not None
def test_docker_executor_stop_state_exception_preserves_handle(tmp_path, monkeypatch):
import ouroboros.workspace_executor as workspace_executor
workspace = tmp_path / "workspace"
workspace.mkdir()
data = tmp_path / "data"
data.mkdir()
ctx = ToolContext(
repo_dir=tmp_path / "repo",
drive_root=data,
workspace_root=workspace,
workspace_mode="external",
task_id="docker-stop-exception",
executor_ref={
"type": "docker_exec",
"id": "pb-container",
"container_name": "pb-container",
"network": "none",
"workspace_host_path": str(workspace),
"workspace_backend_path": "/workspace",
},
)
workspace_executor._SERVICES.clear()
def fake_run(cmd, **kwargs):
if cmd[:3] == ["docker", "inspect", "-f"]:
return subprocess.CompletedProcess(cmd, 0, stdout="none\n", stderr="")
if cmd[:2] == ["docker", "exec"] and "nohup" in str(cmd[-1]):
return subprocess.CompletedProcess(cmd, 0, stdout="12345\n", stderr="")
if cmd[:2] == ["docker", "exec"] and "kill -TERM" in str(cmd[-1]):
return subprocess.CompletedProcess(cmd, 0, stdout="", stderr="")
raise AssertionError(cmd)
monkeypatch.setattr(workspace_executor.subprocess, "run", fake_run)
monkeypatch.setattr(workspace_executor, "_service_state", lambda _record: (_ for _ in ()).throw(RuntimeError("probe boom")))
workspace_executor.start_service(
ctx,
name="svc",
cmd=["sleep", "30"],
host_cwd=workspace,
cwd_root="active_workspace",
readiness={},
outputs=[],
before_outputs={},
)
failed = workspace_executor.stop_service(ctx, "svc")
assert failed and failed["stop_failed"] is True
assert failed["state"] == "unknown"
assert workspace_executor.service_status(ctx, "svc") is not None
def test_docker_executor_global_cleanup_unknown_state_keeps_handle(tmp_path, monkeypatch):
import ouroboros.workspace_executor as workspace_executor
workspace = tmp_path / "workspace"
workspace.mkdir()
data = tmp_path / "data"
data.mkdir()
ctx = ToolContext(
repo_dir=tmp_path / "repo",
drive_root=data,
workspace_root=workspace,
workspace_mode="external",
task_id="docker-cleanup-unknown",
executor_ref={
"type": "docker_exec",
"id": "pb-container",
"container_name": "pb-container",
"network": "none",
"workspace_host_path": str(workspace),
"workspace_backend_path": "/workspace",
},
)
workspace_executor._SERVICES.clear()
def fake_run(cmd, **kwargs):
if cmd[:3] == ["docker", "inspect", "-f"]:
return subprocess.CompletedProcess(cmd, 0, stdout="none\n", stderr="")
if cmd[:2] == ["docker", "exec"] and "nohup" in str(cmd[-1]):
return subprocess.CompletedProcess(cmd, 0, stdout="12345\n", stderr="")
if cmd[:2] == ["docker", "exec"] and "kill -TERM" in str(cmd[-1]):
return subprocess.CompletedProcess(cmd, 0, stdout="", stderr="")
if cmd[:2] == ["docker", "exec"] and "kill -0" in str(cmd[-1]):
return subprocess.CompletedProcess(cmd, 7, stdout="", stderr="daemon unavailable")
raise AssertionError(cmd)
monkeypatch.setattr(workspace_executor.subprocess, "run", fake_run)
workspace_executor.start_service(
ctx,
name="svc",
cmd=["sleep", "30"],
host_cwd=workspace,
cwd_root="active_workspace",
readiness={},
outputs=[],
before_outputs={},
)
result = workspace_executor.kill_all_services(data)
current = next(item for item in result if item.get("name") == "svc")
assert current["cleanup_dispatched"] is False
assert current["stop_failed"] is True
assert current["state"] == "unknown"
assert workspace_executor.service_status(ctx, "svc") is not None
def test_docker_durable_cleanup_keeps_record_until_kill_zero_terminal(tmp_path, monkeypatch):
import ouroboros.workspace_executor as workspace_executor
@ -1118,6 +1283,42 @@ def test_docker_durable_cleanup_keeps_record_until_kill_zero_terminal(tmp_path,
assert path.exists()
def test_docker_durable_cleanup_keeps_record_on_unknown_kill_zero(tmp_path, monkeypatch):
import ouroboros.workspace_executor as workspace_executor
data = tmp_path / "data"
data.mkdir()
workspace_executor._SERVICES.clear()
path = workspace_executor._register_process(
data,
{
"record_type": "service",
"executor_type": "docker_exec",
"executor_id": "pb-container",
"container_name": "pb-container",
"backend_pid": "12345",
"service_id": "task:durable-unknown",
"task_id": "task",
"name": "durable-unknown",
},
)
assert path is not None
def fake_run(cmd, **kwargs):
if cmd[:2] == ["docker", "exec"] and "kill -TERM" in str(cmd[-1]):
return subprocess.CompletedProcess(cmd, 0, stdout="", stderr="")
if cmd[:2] == ["docker", "exec"] and "kill -0" in str(cmd[-1]):
return subprocess.CompletedProcess(cmd, 7, stdout="", stderr="daemon unavailable")
raise AssertionError(cmd)
monkeypatch.setattr(workspace_executor.subprocess, "run", fake_run)
result = workspace_executor._kill_durable_service_records(data)
assert result[0]["state"] == "cleanup_pending"
assert result[0]["cleanup_dispatched"] is False
assert path.exists()
def test_docker_executor_service_shell_uses_process_group_stop():
from ouroboros.workspace_executor import _docker_service_start_shell, _docker_service_stop_shell

View file

@ -79,6 +79,8 @@ def test_executor_panic_cleanup_wait_false_uses_bounded_docker_stop(tmp_path, mo
def fake_docker_wait(cmd, **kwargs):
docker_run_calls.append([str(part) for part in cmd])
if "kill -0" in str(cmd[-1]):
return subprocess.CompletedProcess(cmd, 0, stdout="exited\n", stderr="")
return subprocess.CompletedProcess(cmd, 0, stdout="", stderr="")
class FakePopen:
@ -91,9 +93,9 @@ def test_executor_panic_cleanup_wait_false_uses_bounded_docker_stop(tmp_path, mo
killed_foreground = workspace_executor.kill_all_foreground(data, wait=False)
killed_services = workspace_executor.kill_all_services(data, wait=False)
# The in-memory service now performs an explicit kill-0 confirmation after
# the stop shell, in addition to the three cleanup dispatches.
assert len(docker_run_calls) == 4
# Both in-memory and durable services perform an explicit kill-0
# confirmation after the stop shell, in addition to the dispatches.
assert len(docker_run_calls) == 5
assert all(call[:2] == ["docker", "exec"] for call in docker_run_calls)
assert any(item.get("executor_type") == "docker_exec" for item in killed_foreground)
assert any(item.get("state") == "stopped" for item in killed_services)