mirror of
https://github.com/razzant/ouroboros.git
synced 2026-10-03 04:07:04 +00:00
fix: preserve executor probe uncertainty
Co-authored-by: Ouroboros <311266734+ouroboros-agent@users.noreply.github.com>
This commit is contained in:
parent
dd10bd51c5
commit
a849c9a6b5
4 changed files with 387 additions and 27 deletions
|
|
@ -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:
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue