"""Real-server scheduler integration (T12). Proves the actual lifecycle (``build_scheduler_service`` + start/stop over the server's own stores and runtime) driving REAL workflows through the server's own runtime: constant completion, ``wf.std.runtime_error`` failure, a double-interrupt restart/resume chain, capacity saturation with late settlement, execution-death abandonment without replay (subprocess), graceful vs forced shutdown, manual-run coexistence, and paused/deleted schedule semantics. Schedules are created directly through ``FileScheduleStore`` (the T13 API does not exist yet). All synchronization is gate- or store-state-based with ``asyncio.timeout``-bounded polling; the only sleeps are tiny poll intervals (plus one short quiescence window for the no-replay negative assertions). Fresh ``tmp_path`` per test; every service is stopped in a ``finally`` block. """ from __future__ import annotations import asyncio import os import subprocess import sys import threading from collections.abc import Callable from datetime import UTC, datetime, timedelta from pathlib import Path from typing import Any, cast import pytest from pydantic import BaseModel from wf_api.models import RawWorkflowPlan from wf_artifacts import FileRunStore from wf_authoring import NodeSpec from wf_core import END from wf_platform import ( CapabilityBuckets, CapabilitySource, SourcePermissions, SourcePolicy, SourceVisibility, ) from wf_scheduling.lifecycle import SchedulerServiceConfig, SchedulerStartupError from wf_scheduling.models import Schedule from wf_scheduling.resume_gate import ( ScheduledCapacityBusyError, ScheduledResumeShutdownError, ) from wf_scheduling.store import FileScheduleStore, ScheduleExistsError from wf_server import WorkflowServer, build_local_static_workflow_server from wf_server.scheduling import build_scheduler_service GATE_FILE_ENV = "T12_SCHED_GATE_FILE" EFFECT_FILE_ENV = "T12_SCHED_EFFECT_FILE" _gate_open = threading.Event() class _GateInput(BaseModel): """Input model for the test gate node (wf.std.constant-shaped).""" value: Any = "hello" class _GateOutput(BaseModel): """Output model for the test gate node.""" value: Any = "hello" async def _gate_wait(payload: _GateInput, ctx: Any) -> _GateOutput: """Block until the test opens the gate, recording one visible effect. In-process tests open the gate via ``_gate_open``; subprocess children (which cannot share memory) open it via the ``T12_SCHED_GATE_FILE`` path instead. When ``T12_SCHED_EFFECT_FILE`` is set, exactly one line is appended before blocking, giving the death tests an externally visible dispatch effect. Cancellation is never swallowed, so shutdown drain observes the hang truthfully. """ effect = os.environ.get(EFFECT_FILE_ENV) if effect: with open(effect, "a", encoding="utf-8") as handle: handle.write("scheduled-effect\n") gate_file = os.environ.get(GATE_FILE_ENV) while True: if _gate_open.is_set(): break if gate_file and Path(gate_file).exists(): break await asyncio.sleep(0.02) return _GateOutput(value=payload.value) def _gate_sources() -> dict[str, CapabilitySource]: """Return the minimal custom source owning the ``test.gate.wait`` node.""" spec = NodeSpec( name="test.gate.wait", input_model=_GateInput, output_model=_GateOutput, outcomes=("ok",), fn=_gate_wait, description="Test-only gate: blocks until opened, then echoes.", is_async=True, ) return { "test.gate": CapabilitySource( id="test.gate", kind="python", capabilities=CapabilityBuckets(node_specs={"test.gate.wait": spec}), visibility=SourceVisibility(planner=True, client=True), permissions=SourcePermissions(safe_for_workflow=True), policy=SourcePolicy(platform=True, binding_required=False), description="Test-only gating source.", ) } def _constant_plan(name: str, value: str = "hello") -> RawWorkflowPlan: return RawWorkflowPlan.model_validate( { "name": name, "input_schema": {"type": "object", "properties": {}}, "state_schema": { "type": "object", "properties": {"result": {"type": "string"}}, }, "output_schema": { "type": "object", "properties": {"result": {"type": "string"}}, "required": ["result"], }, "outcomes": ["ok"], "start": "constant", "nodes": [ { "id": "constant", "type": "node", "node": "wf.std.constant", "input": [ { "value": value, "target": {"root": "local", "parts": ["value"]}, } ], "output": [ { "source": {"root": "local", "parts": ["value"]}, "target": {"root": "state", "parts": ["result"]}, } ], } ], "edges": [{"from": "constant", "outcome": "ok", "to": END}], "output": [ { "path": {"root": "state", "parts": ["result"]}, "target": {"root": "local", "parts": ["result"]}, } ], } ) def _failing_plan(name: str) -> RawWorkflowPlan: """A workflow whose real execution ends failed (raising node).""" return RawWorkflowPlan.model_validate( { "name": name, "input_schema": {"type": "object", "properties": {}}, "state_schema": {"type": "object", "properties": {}}, "output_schema": {"type": "object", "properties": {}}, "outcomes": ["ok"], "start": "fail", "nodes": [ { "id": "fail", "type": "node", "node": "wf.std.runtime_error", "input": [ { "value": "boom", "target": {"root": "local", "parts": ["message"]}, } ], "output": [], } ], "edges": [{"from": "fail", "outcome": "ok", "to": END}], "output": [], } ) def _double_interrupt_plan(name: str) -> RawWorkflowPlan: """Two sequential approval interrupts (ask1 -> ask2 -> end).""" return RawWorkflowPlan.model_validate( { "name": name, "input_schema": {"type": "object", "properties": {}}, "state_schema": {"type": "object", "properties": {}}, "output_schema": {"type": "object", "properties": {}}, "outcomes": ["submitted"], "start": "ask1", "nodes": [ {"id": "ask1", "type": "interrupt", "kind": "approval"}, {"id": "ask2", "type": "interrupt", "kind": "approval"}, {"id": "end_submitted", "type": "end", "outcome": "submitted"}, ], "edges": [ {"from": "ask1", "outcome": "submitted", "to": "ask2"}, {"from": "ask2", "outcome": "submitted", "to": "end_submitted"}, ], "output": [], } ) def _single_interrupt_plan(name: str) -> RawWorkflowPlan: """One approval interrupt for manual resume-interaction coverage.""" return RawWorkflowPlan.model_validate( { "name": name, "input_schema": {"type": "object", "properties": {}}, "state_schema": {"type": "object", "properties": {}}, "output_schema": {"type": "object", "properties": {}}, "outcomes": ["submitted"], "start": "ask", "nodes": [ {"id": "ask", "type": "interrupt", "kind": "approval"}, {"id": "end_submitted", "type": "end", "outcome": "submitted"}, ], "edges": [{"from": "ask", "outcome": "submitted", "to": "end_submitted"}], "output": [], } ) def _gate_plan(name: str) -> RawWorkflowPlan: """Constant-shaped workflow whose node blocks on the test gate.""" return RawWorkflowPlan.model_validate( { "name": name, "input_schema": {"type": "object", "properties": {}}, "state_schema": { "type": "object", "properties": {"result": {"type": "string"}}, }, "output_schema": { "type": "object", "properties": {"result": {"type": "string"}}, "required": ["result"], }, "outcomes": ["ok"], "start": "gate", "nodes": [ { "id": "gate", "type": "node", "node": "test.gate.wait", "input": [ { "value": "hello", "target": {"root": "local", "parts": ["value"]}, } ], "output": [ { "source": {"root": "local", "parts": ["value"]}, "target": {"root": "state", "parts": ["result"]}, } ], } ], "edges": [{"from": "gate", "outcome": "ok", "to": END}], "output": [ { "path": {"root": "state", "parts": ["result"]}, "target": {"root": "local", "parts": ["result"]}, } ], } ) def _one_shot( schedule_id: str, deployment_id: str, at: datetime, **overrides: Any ) -> Schedule: base: dict[str, Any] = { "id": schedule_id, "deployment_id": deployment_id, "trigger": {"kind": "oneshot", "at": at.isoformat()}, "input_bindings": [], "created_at": at.isoformat(), "updated_at": at.isoformat(), } base.update(overrides) return Schedule.model_validate(base) async def _seed( server: WorkflowServer, artifact_id: str, deployment_id: str, plan: RawWorkflowPlan, outcomes: list[str], ) -> None: """Create one executable artifact + deployment through the server API.""" await server.api.create_artifact_from_plan( artifact_id=artifact_id, version=1, title=artifact_id, plan=plan, outcomes=outcomes, source_bindings={}, ) await server.api.save_deployment( { "id": deployment_id, "artifact_id": artifact_id, "artifact_version": 1, "bindings": {}, } ) def _scheduler( server: WorkflowServer, *, capacity: int = 4, poll_interval_s: float = 0.02, drain_grace_s: float = 5.0, auto_tick: bool = True, ): return build_scheduler_service( server, SchedulerServiceConfig( poll_interval_s=poll_interval_s, capacity=capacity, drain_grace_s=drain_grace_s, auto_tick=auto_tick, ), ) async def _wait_for(condition: Callable[[], bool], timeout: float = 20.0) -> None: async with asyncio.timeout(timeout): while not condition(): await asyncio.sleep(0.02) def _entries(root: Path, schedule_id: str, kind: str) -> list[dict[str, Any]]: page = FileScheduleStore(root).list_occurrences(schedule_id, limit=100) rows = cast(list[dict[str, Any]], page["occurrences"]) return [row for row in rows if row["kind"] == kind] async def _kinds(root: Path, schedule_id: str) -> list[str]: page = FileScheduleStore(root).list_occurrences(schedule_id, limit=100) rows = cast(list[dict[str, Any]], page["occurrences"]) return [row["kind"] for row in rows] def _run_ids(root: Path) -> list[str]: return [record.id for record in FileRunStore(root).list_runs()] async def _start_with_retry(service: Any, timeout: float = 15.0) -> None: """Start a service, tolerating the post-kill lock-release beat.""" async with asyncio.timeout(timeout): while True: try: await service.start() return except SchedulerStartupError: await asyncio.sleep(0.2) _CHILD_SCRIPT = """\ import asyncio import sys sys.path.insert(0, {repo_root!r}) from pathlib import Path from tests.wf_server.test_scheduler_integration import _gate_sources from wf_scheduling.lifecycle import SchedulerServiceConfig from wf_server import build_local_static_workflow_server from wf_server.scheduling import build_scheduler_service async def _main() -> None: root = Path(sys.argv[1]) server = build_local_static_workflow_server(root, extra_sources=_gate_sources()) service = build_scheduler_service( server, SchedulerServiceConfig( poll_interval_s=0.02, capacity=4, drain_grace_s=30.0, auto_tick=True ), ) await service.start() print("child-ready", flush=True) await asyncio.sleep(120) asyncio.run(_main()) """ _REPO_ROOT = Path(__file__).resolve().parents[2] def _child_script() -> str: return _CHILD_SCRIPT.format(repo_root=str(_REPO_ROOT)) async def _wait_for_child_executing( proc: subprocess.Popen[str], root: Path, effect_file: Path, timeout: float = 30.0 ) -> str: """Wait until the child dispatched with a visible effect. The durable executing mark comes first (per the dispatch contract the mark persists before the executor runs), then the effect line proves the executor actually entered the node: killing only afterwards closes the mark-then-dispatch race deterministically. """ runs_store = FileRunStore(root) run_id: str | None = None async with asyncio.timeout(timeout): while run_id is None: if proc.poll() is not None: output = proc.stdout.read() if proc.stdout is not None else "" raise AssertionError( f"child exited early with rc={proc.returncode}: {output}" ) runs = runs_store.list_runs() if runs and runs_store.is_executing(runs[0].id): run_id = runs[0].id else: await asyncio.sleep(0.05) while True: if proc.poll() is not None: output = proc.stdout.read() if proc.stdout is not None else "" raise AssertionError(f"child exited before the effect landed: {output}") if effect_file.exists() and effect_file.read_text( encoding="utf-8" ).splitlines() == ["scheduled-effect"]: return run_id await asyncio.sleep(0.05) def _kill_child(proc: subprocess.Popen[str]) -> None: """Windows-safe forced termination (Popen.kill, no signals).""" if proc.poll() is None: proc.kill() try: proc.communicate(timeout=30) except subprocess.TimeoutExpired: proc.kill() proc.communicate(timeout=30) async def test_scheduled_completion(tmp_path: Path) -> None: root = tmp_path / "store" server = build_local_static_workflow_server(root) await _seed( server, "sched_const", "sched_const.default", _constant_plan("sched_const"), ["ok"], ) FileScheduleStore(root).create_schedule( _one_shot("once", "sched_const.default", datetime.now(UTC)) ) service = _scheduler(server) try: await service.start() await _wait_for(lambda: len(_run_ids(root)) == 1) run_id = _run_ids(root)[0] await _wait_for( lambda: FileRunStore(root).get_run(run_id).status.value == "completed" ) summary = await server.api.inspect_run(run_id=run_id) assert summary["status"] == "completed" assert summary["output"] is not None assert summary["output"]["result"] == "hello" completed = _entries(root, "once", "completed") assert len(completed) == 1 assert completed[0]["run_id"] == run_id finally: report = await service.stop() assert report.settled >= 1 @pytest.mark.parametrize("mutation", ["pause", "update", "delete", "shutdown"]) def test_admin_mutation_does_not_deadlock_dispatch_registration( tmp_path: Path, mutation: str ) -> None: """A store mutation must not block registration on the event loop.""" child = f""" import asyncio import sys import threading import time from pathlib import Path sys.path.insert(0, {str(_REPO_ROOT)!r}) from datetime import UTC, datetime, timedelta from tests.wf_server.test_scheduler_integration import ( _constant_plan, _one_shot, _scheduler, _seed, ) import wf_scheduling.lifecycle as lifecycle from wf_scheduling.ownership import SchedulerOwnership from wf_server import build_local_static_workflow_server async def main() -> None: root = Path(sys.argv[1]) server = build_local_static_workflow_server(root, schedules=True) await _seed( server, "constant", "constant.default", _constant_plan("constant"), ["ok"] ) due = _one_shot("once", "constant.default", datetime.now(UTC)) await server.api.schedules.create_schedule( schedule_id=due.id, deployment_id=due.deployment_id, trigger=due.trigger.model_dump(mode="json"), ) # API creation starts the watermark at creation time; move it behind the # already-created one-shot so the deterministic poll observes it as due. server.api.schedules._schedule_store().save_consumed( "once", due.trigger.at - timedelta(seconds=1) ) service = _scheduler(server, capacity=1, auto_tick=False) await service.start() registration_entered = threading.Event() release_registration = threading.Event() original = asyncio.run_coroutine_threadsafe def blocked_registration(coro, loop): registration_entered.set() if not release_registration.wait(5): raise AssertionError("registration release was not signalled") return original(coro, loop) lifecycle.asyncio.run_coroutine_threadsafe = blocked_registration poll_task = asyncio.create_task( service.poll_once(due.trigger.at + timedelta(seconds=1)) ) if not await asyncio.to_thread(registration_entered.wait, 5): raise AssertionError("poll did not reach the registration gate") if sys.argv[2] in ("pause", "shutdown"): admin_task = asyncio.create_task( server.api.schedules.pause_schedule(schedule_id="once") ) elif sys.argv[2] == "update": admin_task = asyncio.create_task( server.api.schedules.update_schedule( schedule_id="once", expected_revision=1, overlap="parallel" ) ) else: admin_task = asyncio.create_task( server.api.schedules.delete_schedule(schedule_id="once") ) threading.Thread( target=lambda: (time.sleep(0.2), release_registration.set()), daemon=True, ).start() try: if sys.argv[2] == "shutdown": shutdown_task = asyncio.create_task(service.stop()) result, paused, report = await asyncio.wait_for( asyncio.gather(poll_task, admin_task, shutdown_task), 4 ) assert report is not None assert service.running is False probe = SchedulerOwnership(root, owner="shutdown-probe").acquire() probe.release() else: result, paused = await asyncio.wait_for( asyncio.gather(poll_task, admin_task), 4 ) assert result["once"].startswith("admit:run-") run_id = result["once"].split(":", 1)[1] admission = server.stores.run_store.get_admission(run_id) assert admission.schedule_revision == 1 if sys.argv[2] in ("pause", "shutdown"): assert paused["paused"] is True elif sys.argv[2] == "update": assert paused["revision"] == 2 else: assert paused["deleted"] is True print("admin-registration-complete", flush=True) finally: release_registration.set() await service.stop() asyncio.run(main()) """ script = tmp_path / "admin_registration_deadlock.py" script.write_text(child, encoding="utf-8") process = subprocess.Popen( [sys.executable, str(script), str(tmp_path / "store"), mutation], cwd=_REPO_ROOT, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, text=True, ) try: output, _ = process.communicate(timeout=30) except subprocess.TimeoutExpired: process.kill() output, _ = process.communicate(timeout=5) pytest.fail(f"admin/registration deadlocked; child output: {output}") assert process.returncode == 0, output assert "admin-registration-complete" in output async def test_scheduled_failure(tmp_path: Path) -> None: root = tmp_path / "store" server = build_local_static_workflow_server(root) await _seed( server, "sched_fail", "sched_fail.default", _failing_plan("sched_fail"), ["ok"] ) FileScheduleStore(root).create_schedule( _one_shot("once", "sched_fail.default", datetime.now(UTC)) ) service = _scheduler(server) try: await service.start() await _wait_for(lambda: len(_run_ids(root)) == 1) run_id = _run_ids(root)[0] await _wait_for( lambda: FileRunStore(root).get_run(run_id).status.value == "failed" ) record = FileRunStore(root).get_run(run_id) assert record.status.value == "failed" assert record.resume_readiness.value == "not_applicable" summary = await server.api.inspect_run(run_id=run_id) assert summary["status"] == "failed" assert summary["resume_readiness"] == "not_applicable" failed = _entries(root, "once", "failed") assert len(failed) == 1 assert failed[0]["run_id"] == run_id finally: await service.stop() async def test_interrupt_restart_resume_reinterrupt(tmp_path: Path) -> None: root = tmp_path / "store" server = build_local_static_workflow_server(root) await _seed( server, "ask_twice", "ask_twice.default", _double_interrupt_plan("ask_twice"), ["submitted"], ) FileScheduleStore(root).create_schedule( _one_shot("asker", "ask_twice.default", datetime.now(UTC)) ) first = _scheduler(server) try: await first.start() await _wait_for(lambda: len(_run_ids(root)) == 1) run_id = _run_ids(root)[0] await _wait_for( lambda: FileRunStore(root).get_run(run_id).status.value == "interrupted" ) assert FileRunStore(root).get_run(run_id).resume_readiness.value == "ready" assert ( FileRunStore(root).get_run(run_id).latest_checkpoint_id == f"{run_id}.000001" ) await _wait_for(lambda: len(_entries(root, "asker", "interrupted")) == 1) assert len(_entries(root, "asker", "interrupted")) == 1 finally: await first.stop() # Real restart boundary: fresh server + service instances, same root. server_b = build_local_static_workflow_server(root) second = _scheduler(server_b) try: await second.start() record = FileRunStore(root).get_run(run_id) assert record.status.value == "interrupted" assert record.resume_readiness.value == "ready" resumed = await server_b.api.resume_run( run_id=run_id, resume_payload={}, resume_outcome="submitted" ) assert resumed["status"] == "interrupted" assert resumed["resume_readiness"] == "ready" assert ( FileRunStore(root).get_run(run_id).latest_checkpoint_id == f"{run_id}.000002" ) finally: await second.stop() # Startup recovery reconciles the resume-produced second interruption. server_c = build_local_static_workflow_server(root) third = _scheduler(server_c) try: await third.start() interrupted = _entries(root, "asker", "interrupted") assert len(interrupted) == 2 assert {row["checkpoint_id"] for row in interrupted} == { f"{run_id}.000001", f"{run_id}.000002", } finished = await server_c.api.resume_run( run_id=run_id, resume_payload={}, resume_outcome="submitted" ) assert finished["status"] == "completed" assert ( FileRunStore(root).get_run(run_id).latest_checkpoint_id == f"{run_id}.000003" ) finally: await third.stop() server_d = build_local_static_workflow_server(root) fourth = _scheduler(server_d) try: await fourth.start() assert FileRunStore(root).get_run(run_id).status.value == "completed" completed = _entries(root, "asker", "completed") assert len(completed) == 1 assert completed[0]["checkpoint_id"] == f"{run_id}.000003" finally: await fourth.stop() async def test_capacity_saturation_late_settlement( tmp_path: Path, monkeypatch: pytest.MonkeyPatch ) -> None: _gate_open.clear() monkeypatch.delenv(GATE_FILE_ENV, raising=False) monkeypatch.delenv(EFFECT_FILE_ENV, raising=False) root = tmp_path / "store" server = build_local_static_workflow_server(root, extra_sources=_gate_sources()) await _seed(server, "gated", "gated.default", _gate_plan("gated"), ["ok"]) await _seed( server, "plain", "plain.default", _constant_plan("plain", "hello"), ["ok"] ) base = datetime.now(UTC) store = FileScheduleStore(root) store.create_schedule(_one_shot("a", "gated.default", base)) store.create_schedule(_one_shot("b", "plain.default", base)) # Operator-driven ticks keep the live-execution bound deterministic: # admission of B only happens on an explicit poll after A has settled. service = _scheduler(server, capacity=1, auto_tick=False) peak = 0 try: await service.start() first = await service.poll_once(base + timedelta(seconds=1)) assert first["a"].startswith("admit:run-") assert first["b"] == "admit:held-undecided" run_a = first["a"].split(":", 1)[1] await _wait_for(lambda: FileRunStore(root).is_executing(run_a)) assert FileRunStore(root).get_run(run_a).status.value == "admitted" assert service.live_executions == 1 peak = max(peak, service.live_executions) # Capacity is saturated: no second admission while A executes. await service.poll_once(base + timedelta(seconds=2)) assert _run_ids(root) == [run_a] assert service.live_executions == 1 _gate_open.set() await _wait_for( lambda: FileRunStore(root).get_run(run_a).status.value == "completed" ) await _wait_for(lambda: service.live_executions == 0) peak = max(peak, service.live_executions) nxt = await service.poll_once(base + timedelta(seconds=3)) assert nxt["b"].startswith("admit:run-") run_b = nxt["b"].split(":", 1)[1] assert service.live_executions <= 1 peak = max(peak, service.live_executions) await _wait_for( lambda: FileRunStore(root).get_run(run_b).status.value == "completed" ) await _wait_for(lambda: service.live_executions == 0) assert FileRunStore(root).get_run(run_a).status.value == "completed" assert len(_entries(root, "a", "completed")) == 1 assert len(_entries(root, "b", "completed")) == 1 finally: _gate_open.set() await service.stop() assert peak <= 1 @pytest.mark.slow async def test_restart_after_execution_death_no_replay_subprocess( tmp_path: Path, ) -> None: """SUBPROCESS death test: kill mid-execution, restart abandons, no replay.""" root = tmp_path / "store" effect_file = tmp_path / "effects.log" child_env = {**os.environ, EFFECT_FILE_ENV: str(effect_file)} child_env.pop(GATE_FILE_ENV, None) server = build_local_static_workflow_server(root, extra_sources=_gate_sources()) await _seed( server, "death_gate", "death_gate.default", _gate_plan("death_gate"), ["ok"] ) FileScheduleStore(root).create_schedule( _one_shot("doomed", "death_gate.default", datetime.now(UTC)) ) script = tmp_path / "child_death.py" script.write_text(_child_script(), encoding="utf-8") proc = subprocess.Popen( [sys.executable, str(script), str(root)], stdout=subprocess.PIPE, stderr=subprocess.STDOUT, text=True, env=child_env, ) try: run_id = await _wait_for_child_executing(proc, root, effect_file, timeout=30.0) finally: _kill_child(proc) # Real restart boundary in-process: recovery must abandon without replay. server_b = build_local_static_workflow_server(root, extra_sources=_gate_sources()) service = _scheduler(server_b) try: await _start_with_retry(service, timeout=15.0) record = FileRunStore(root).get_run(run_id) assert record.status.value == "failed" assert record.resume_readiness.value == "not_applicable" assert any( item.code == "schedule-recovery" and "external effects may already" in item.message for item in record.diagnostics ) assert len(_entries(root, "doomed", "failed")) == 1 assert len(_entries(root, "doomed", "completed")) == 0 assert effect_file.exists() assert effect_file.read_text(encoding="utf-8").splitlines() == [ "scheduled-effect" ] # Quiescence window (not synchronization): ticks keep running, so a # replay would append a second line here. await asyncio.sleep(1.0) assert effect_file.read_text(encoding="utf-8").splitlines() == [ "scheduled-effect" ] assert FileRunStore(root).get_run(run_id).status.value == "failed" finally: await service.stop() async def test_graceful_shutdown_cancel_leaves_executing_for_recovery( tmp_path: Path, monkeypatch: pytest.MonkeyPatch ) -> None: _gate_open.clear() monkeypatch.delenv(GATE_FILE_ENV, raising=False) monkeypatch.delenv(EFFECT_FILE_ENV, raising=False) root = tmp_path / "store" server = build_local_static_workflow_server(root, extra_sources=_gate_sources()) await _seed(server, "hang", "hang.default", _gate_plan("hang"), ["ok"]) FileScheduleStore(root).create_schedule( _one_shot("hanging", "hang.default", datetime.now(UTC)) ) service = _scheduler(server, drain_grace_s=0.05) try: await service.start() await _wait_for(lambda: len(_run_ids(root)) == 1) run_id = _run_ids(root)[0] await _wait_for(lambda: FileRunStore(root).is_executing(run_id)) assert FileRunStore(root).get_run(run_id).status.value == "admitted" report = await service.stop() assert report.cancelled >= 1 assert report.settled == 0 # Truthful handling: the run keeps its executing mark for recovery. assert FileRunStore(root).get_run(run_id).status.value == "admitted" assert FileRunStore(root).is_executing(run_id) is True finally: _gate_open.set() await service.stop() # Restart abandons the cancelled execution without replay. _gate_open.clear() server_b = build_local_static_workflow_server(root, extra_sources=_gate_sources()) restart = _scheduler(server_b) try: await restart.start() record = FileRunStore(root).get_run(run_id) assert record.status.value == "failed" assert record.resume_readiness.value == "not_applicable" assert any( item.code == "schedule-recovery" and "external effects may already" in item.message for item in record.diagnostics ) inspected = await server_b.api.inspect_run(run_id=run_id) assert inspected["status"] == "failed" assert inspected["resume_readiness"] == "not_applicable" assert inspected["error"] is not None assert "external effects may already" in inspected["error"] assert len(_entries(root, "hanging", "failed")) == 1 finally: _gate_open.set() await restart.stop() @pytest.mark.slow async def test_forced_termination_abandons_without_replay_subprocess( tmp_path: Path, ) -> None: """SUBPROCESS forced test with a releasable gate: kill, restart, release.""" root = tmp_path / "store" effect_file = tmp_path / "effects.log" gate_file = tmp_path / "gate.open" child_env = { **os.environ, EFFECT_FILE_ENV: str(effect_file), GATE_FILE_ENV: str(gate_file), } server = build_local_static_workflow_server(root, extra_sources=_gate_sources()) await _seed( server, "force_gate", "force_gate.default", _gate_plan("force_gate"), ["ok"] ) FileScheduleStore(root).create_schedule( _one_shot("forced", "force_gate.default", datetime.now(UTC)) ) script = tmp_path / "child_forced.py" script.write_text(_child_script(), encoding="utf-8") proc = subprocess.Popen( [sys.executable, str(script), str(root)], stdout=subprocess.PIPE, stderr=subprocess.STDOUT, text=True, env=child_env, ) try: run_id = await _wait_for_child_executing(proc, root, effect_file, timeout=30.0) finally: _kill_child(proc) server_b = build_local_static_workflow_server(root, extra_sources=_gate_sources()) service = _scheduler(server_b) try: await _start_with_retry(service, timeout=15.0) assert FileRunStore(root).get_run(run_id).status.value == "failed" assert len(_entries(root, "forced", "failed")) == 1 assert effect_file.read_text(encoding="utf-8").splitlines() == [ "scheduled-effect" ] # Release the gate after death: a replay would now finish and append. gate_file.write_text("open", encoding="utf-8") await asyncio.sleep(1.0) assert effect_file.read_text(encoding="utf-8").splitlines() == [ "scheduled-effect" ] assert FileRunStore(root).get_run(run_id).status.value == "failed" assert _run_ids(root) == [run_id] finally: await service.stop() async def test_manual_run_resume_interaction(tmp_path: Path) -> None: root = tmp_path / "store" server = build_local_static_workflow_server(root) await _seed( server, "sched_const", "sched_const.default", _constant_plan("sched_const"), ["ok"], ) await _seed( server, "manual_ask", "manual_ask.default", _single_interrupt_plan("manual_ask"), ["submitted"], ) FileScheduleStore(root).create_schedule( _one_shot("ticker", "sched_const.default", datetime.now(UTC)) ) service = _scheduler(server) try: await service.start() await _wait_for(lambda: len(_run_ids(root)) == 1) scheduled_id = _run_ids(root)[0] await _wait_for( lambda: FileRunStore(root).get_run(scheduled_id).status.value == "completed" ) # A manual run completes concurrently with its own run id. manual = await server.api.run_deployment( deployment_id="sched_const.default", workflow_input={} ) assert manual["status"] == "completed" manual_id = manual["run_id"] assert isinstance(manual_id, str) assert manual_id != scheduled_id assert sorted(_run_ids(root)) == sorted([scheduled_id, manual_id]) # Manual interrupt + resume works while the scheduler stays active. asked = await server.api.run_deployment( deployment_id="manual_ask.default", workflow_input={} ) assert asked["status"] == "interrupted" ask_id = asked["run_id"] assert isinstance(ask_id, str) assert service.running is True resumed = await server.api.resume_run( run_id=ask_id, resume_payload={}, resume_outcome="submitted" ) assert resumed["status"] == "completed" assert service.running is True # The scheduled overlap slot is unaffected by manual runs: a second # one-shot admits on its own occurrence after the manual work. FileScheduleStore(root).create_schedule( _one_shot("ticker2", "sched_const.default", datetime.now(UTC)) ) await _wait_for(lambda: len(_run_ids(root)) == 4) await _wait_for( lambda: all( FileRunStore(root).get_run(rid).status.value == "completed" for rid in _run_ids(root) ) ) assert len(_entries(root, "ticker", "completed")) == 1 assert len(_entries(root, "ticker2", "completed")) == 1 finally: report = await service.stop() assert report.settled >= 2 async def test_paused_deleted_run_completion( tmp_path: Path, monkeypatch: pytest.MonkeyPatch ) -> None: _gate_open.clear() monkeypatch.delenv(GATE_FILE_ENV, raising=False) monkeypatch.delenv(EFFECT_FILE_ENV, raising=False) root = tmp_path / "store" server = build_local_static_workflow_server(root, extra_sources=_gate_sources()) await _seed(server, "gated", "gated.default", _gate_plan("gated"), ["ok"]) await _seed( server, "ask_twice", "ask_twice.default", _double_interrupt_plan("ask_twice"), ["submitted"], ) store = FileScheduleStore(root) store.create_schedule(_one_shot("pausable", "gated.default", datetime.now(UTC))) service = _scheduler(server) try: await service.start() # Pause mid-execution: the admitted run still completes on release. await _wait_for(lambda: len(_run_ids(root)) == 1) gated_id = _run_ids(root)[0] await _wait_for(lambda: FileRunStore(root).is_executing(gated_id)) paused = store.get_schedule("pausable") paused.paused = True store.save_schedule(paused) _gate_open.set() await _wait_for( lambda: FileRunStore(root).get_run(gated_id).status.value == "completed" ) await _wait_for(lambda: len(_entries(root, "pausable", "completed")) == 1) assert len(_entries(root, "pausable", "completed")) == 1 # A paused double-interrupt schedule still resumes through the API. _gate_open.clear() store.create_schedule( _one_shot("pausable_ask", "ask_twice.default", datetime.now(UTC)) ) await _wait_for(lambda: len(_run_ids(root)) == 2) ask_id = [rid for rid in _run_ids(root) if rid != gated_id][0] await _wait_for( lambda: FileRunStore(root).get_run(ask_id).status.value == "interrupted" ) pausing = store.get_schedule("pausable_ask") pausing.paused = True store.save_schedule(pausing) first = await server.api.resume_run( run_id=ask_id, resume_payload={}, resume_outcome="submitted" ) assert first["status"] == "interrupted" second = await server.api.resume_run( run_id=ask_id, resume_payload={}, resume_outcome="submitted" ) assert second["status"] == "completed" finally: _gate_open.set() await service.stop() # Restart recovery reconciles the resume-produced completion onto the # still-paused schedule. server_b = build_local_static_workflow_server(root, extra_sources=_gate_sources()) recovery_service = _scheduler(server_b) try: await recovery_service.start() assert len(_entries(root, "pausable_ask", "interrupted")) >= 1 assert len(_entries(root, "pausable_ask", "completed")) == 1 finally: await recovery_service.stop() # Delete with an active run: completion still updates retained history, # and the id is never reusable. _gate_open.clear() server_c = build_local_static_workflow_server(root, extra_sources=_gate_sources()) deleter = _scheduler(server_c) try: await deleter.start() FileScheduleStore(root).create_schedule( _one_shot("doomed", "gated.default", datetime.now(UTC)) ) await _wait_for(lambda: len(_run_ids(root)) == 3) doomed_id = [rid for rid in _run_ids(root) if rid not in {gated_id, ask_id}][0] await _wait_for(lambda: FileRunStore(root).is_executing(doomed_id)) gone = FileScheduleStore(root).get_schedule("doomed") gone.deleted = True FileScheduleStore(root).save_schedule(gone) _gate_open.set() await _wait_for( lambda: FileRunStore(root).get_run(doomed_id).status.value == "completed" ) await _wait_for(lambda: len(_entries(root, "doomed", "completed")) == 1) assert len(_entries(root, "doomed", "completed")) == 1 with pytest.raises(ScheduleExistsError): FileScheduleStore(root).create_schedule( _one_shot("doomed", "gated.default", datetime.now(UTC)) ) finally: _gate_open.set() await deleter.stop() def _interrupt_then_fail_plan(name: str) -> RawWorkflowPlan: """Approval interrupt whose resumed continuation raises (resume failure).""" return RawWorkflowPlan.model_validate( { "name": name, "input_schema": {"type": "object", "properties": {}}, "state_schema": {"type": "object", "properties": {}}, "output_schema": {"type": "object", "properties": {}}, "outcomes": ["submitted"], "start": "ask", "nodes": [ {"id": "ask", "type": "interrupt", "kind": "approval"}, { "id": "fail", "type": "node", "node": "wf.std.runtime_error", "input": [ { "value": "boom on resume", "target": {"root": "local", "parts": ["message"]}, } ], "output": [], }, {"id": "end_submitted", "type": "end", "outcome": "submitted"}, ], "edges": [ {"from": "ask", "outcome": "submitted", "to": "fail"}, {"from": "fail", "outcome": "ok", "to": "end_submitted"}, ], "output": [], } ) def _interrupt_then_gate_plan(name: str) -> RawWorkflowPlan: """Approval interrupt whose resumed continuation blocks on the test gate.""" return RawWorkflowPlan.model_validate( { "name": name, "input_schema": {"type": "object", "properties": {}}, "state_schema": { "type": "object", "properties": {"result": {"type": "string"}}, }, "output_schema": { "type": "object", "properties": {"result": {"type": "string"}}, }, "outcomes": ["submitted"], "start": "ask", "nodes": [ {"id": "ask", "type": "interrupt", "kind": "approval"}, { "id": "gate", "type": "node", "node": "test.gate.wait", "input": [ { "value": "hello", "target": {"root": "local", "parts": ["value"]}, } ], "output": [ { "source": {"root": "local", "parts": ["value"]}, "target": {"root": "state", "parts": ["result"]}, } ], }, {"id": "end_submitted", "type": "end", "outcome": "submitted"}, ], "edges": [ {"from": "ask", "outcome": "submitted", "to": "gate"}, {"from": "gate", "outcome": "ok", "to": "end_submitted"}, ], "output": [], } ) async def test_scheduled_resume_rejected_when_capacity_saturated( tmp_path: Path, ) -> None: """A scheduled resume acquires the shared execution slot (B2). With the only slot occupied by a gated scheduled run, resuming the interrupted scheduled run is rejected BEFORE any dispatch: no ACTIVE attempt, no executing mark, both runs untouched. Freeing the slot lets the same resume through, releasing the slot afterwards. """ _gate_open.clear() root = tmp_path / "store" server = build_local_static_workflow_server(root, extra_sources=_gate_sources()) await _seed( server, "ask1", "ask1.default", _single_interrupt_plan("ask1"), ["submitted"] ) await _seed(server, "gated", "gated.default", _gate_plan("gated"), ["ok"]) store = FileScheduleStore(root) store.create_schedule(_one_shot("asker", "ask1.default", datetime.now(UTC))) service = _scheduler(server, capacity=1) try: await service.start() await _wait_for(lambda: len(_run_ids(root)) == 1) ask_id = _run_ids(root)[0] await _wait_for( lambda: FileRunStore(root).get_run(ask_id).status.value == "interrupted" ) store.create_schedule(_one_shot("blocker", "gated.default", datetime.now(UTC))) await _wait_for(lambda: len(_run_ids(root)) == 2) gate_id = [rid for rid in _run_ids(root) if rid != ask_id][0] await _wait_for(lambda: FileRunStore(root).is_executing(gate_id)) with pytest.raises(ScheduledCapacityBusyError): await server.api.resume_run( run_id=ask_id, resume_payload={}, resume_outcome="submitted" ) assert FileRunStore(root).get_resume_attempt(ask_id) is None assert not FileRunStore(root).is_executing(ask_id) assert FileRunStore(root).get_run(ask_id).status.value == "interrupted" assert FileRunStore(root).is_executing(gate_id) assert len(_entries(root, "asker", "completed")) == 0 _gate_open.set() await _wait_for( lambda: FileRunStore(root).get_run(gate_id).status.value == "completed" ) resumed = await server.api.resume_run( run_id=ask_id, resume_payload={}, resume_outcome="submitted" ) assert resumed["status"] == "completed" attempt = FileRunStore(root).get_resume_attempt(ask_id) assert attempt is not None and attempt.state == "DONE" assert not FileRunStore(root).is_executing(ask_id) finally: _gate_open.set() await service.stop() async def test_scheduled_resume_releases_slot_on_reinterruption( tmp_path: Path, ) -> None: """A scheduled resume that interrupts again frees its slot (B2). The re-interruption stays durably resumable (attempt DONE, new checkpoint) and a second resume of the same run completes. """ root = tmp_path / "store" server = build_local_static_workflow_server(root) await _seed( server, "ask_twice", "ask_twice.default", _double_interrupt_plan("ask_twice"), ["submitted"], ) store = FileScheduleStore(root) store.create_schedule(_one_shot("asker", "ask_twice.default", datetime.now(UTC))) service = _scheduler(server, capacity=1) try: await service.start() await _wait_for(lambda: len(_run_ids(root)) == 1) ask_id = _run_ids(root)[0] await _wait_for( lambda: FileRunStore(root).get_run(ask_id).status.value == "interrupted" ) first = await server.api.resume_run( run_id=ask_id, resume_payload={}, resume_outcome="submitted" ) assert first["status"] == "interrupted" attempt = FileRunStore(root).get_resume_attempt(ask_id) assert attempt is not None and attempt.state == "DONE" assert not FileRunStore(root).is_executing(ask_id) second = await server.api.resume_run( run_id=ask_id, resume_payload={}, resume_outcome="submitted" ) assert second["status"] == "completed" assert not FileRunStore(root).is_executing(ask_id) finally: await service.stop() async def test_scheduled_resume_releases_slot_on_failure(tmp_path: Path) -> None: """A scheduled resume that fails releases its slot (B2). The failed stopped result is durably persisted with its attempt cleared, and the freed slot admits new scheduled work. """ root = tmp_path / "store" server = build_local_static_workflow_server(root) await _seed( server, "ask_fail", "ask_fail.default", _interrupt_then_fail_plan("ask_fail"), ["submitted"], ) await _seed( server, "const", "const.default", _constant_plan("const"), ["ok"], ) store = FileScheduleStore(root) store.create_schedule(_one_shot("asker", "ask_fail.default", datetime.now(UTC))) service = _scheduler(server, capacity=1) try: await service.start() await _wait_for(lambda: len(_run_ids(root)) == 1) ask_id = _run_ids(root)[0] await _wait_for( lambda: FileRunStore(root).get_run(ask_id).status.value == "interrupted" ) resumed = await server.api.resume_run( run_id=ask_id, resume_payload={}, resume_outcome="submitted" ) assert resumed["status"] == "failed" attempt = FileRunStore(root).get_resume_attempt(ask_id) assert attempt is not None and attempt.state == "DONE" assert not FileRunStore(root).is_executing(ask_id) # The freed slot admits new scheduled work. store.create_schedule(_one_shot("after", "const.default", datetime.now(UTC))) await _wait_for(lambda: len(_run_ids(root)) == 2) after_id = [rid for rid in _run_ids(root) if rid != ask_id][0] await _wait_for( lambda: FileRunStore(root).get_run(after_id).status.value == "completed" ) finally: await service.stop() async def test_service_stop_drains_inflight_scheduled_resume(tmp_path: Path) -> None: """Shutdown waits for an in-flight scheduled resume (B2 drain).""" _gate_open.clear() root = tmp_path / "store" server = build_local_static_workflow_server(root, extra_sources=_gate_sources()) await _seed( server, "ask_gate", "ask_gate.default", _interrupt_then_gate_plan("ask_gate"), ["submitted"], ) store = FileScheduleStore(root) store.create_schedule(_one_shot("asker", "ask_gate.default", datetime.now(UTC))) service = _scheduler(server, capacity=2, drain_grace_s=10.0) try: await service.start() await _wait_for(lambda: len(_run_ids(root)) == 1) ask_id = _run_ids(root)[0] await _wait_for( lambda: FileRunStore(root).get_run(ask_id).status.value == "interrupted" ) resume_task = asyncio.create_task( server.api.resume_run( run_id=ask_id, resume_payload={}, resume_outcome="submitted" ) ) await _wait_for(lambda: ask_id in service._live_resumes) await _wait_for(lambda: FileRunStore(root).is_executing(ask_id)) stop_task = asyncio.create_task(service.stop()) await asyncio.sleep(0.3) # The gate is still closed: the drain must still be waiting. assert not stop_task.done() assert not resume_task.done() _gate_open.set() report = await asyncio.wait_for(stop_task, timeout=20.0) resumed = await asyncio.wait_for(resume_task, timeout=20.0) assert resumed["status"] == "completed" assert not FileRunStore(root).is_executing(ask_id) assert report.cancelled == 0 finally: _gate_open.set() await service.stop() async def test_scheduled_resume_completes_within_grace_with_history( tmp_path: Path, ) -> None: """A resume finishing inside grace stops cleanly with history (shutdown). Operator-driven service (no background ticks): the in-flight resume completes during the drain, stop reports no cancellation, and the resumed completion lands in occurrence history exactly once. """ _gate_open.clear() root = tmp_path / "store" server = build_local_static_workflow_server(root, extra_sources=_gate_sources()) await _seed( server, "ask_gate", "ask_gate.default", _interrupt_then_gate_plan("ask_gate"), ["submitted"], ) store = FileScheduleStore(root) store.create_schedule(_one_shot("asker", "ask_gate.default", datetime.now(UTC))) service = _scheduler(server, capacity=1, drain_grace_s=10.0, auto_tick=False) try: await service.start() async with asyncio.timeout(20): while not _run_ids(root): await service.poll_once(datetime.now(UTC)) await asyncio.sleep(0.02) ask_id = _run_ids(root)[0] await _wait_for( lambda: FileRunStore(root).get_run(ask_id).status.value == "interrupted" ) resume_task = asyncio.create_task( server.api.resume_run( run_id=ask_id, resume_payload={}, resume_outcome="submitted" ) ) await _wait_for(lambda: ask_id in service._live_resumes) stop_task = asyncio.create_task(service.stop()) await asyncio.sleep(0.3) # The gate is still closed: the drain must still be waiting. assert not stop_task.done() assert not resume_task.done() _gate_open.set() resumed = await asyncio.wait_for(resume_task, timeout=20.0) report = await asyncio.wait_for(stop_task, timeout=20.0) assert resumed["status"] == "completed" assert report.cancelled == 0 attempt = FileRunStore(root).get_resume_attempt(ask_id) assert attempt is not None and attempt.state == "DONE" assert not FileRunStore(root).is_executing(ask_id) assert await _kinds(root, "asker") == ["admitted", "interrupted", "completed"] finally: _gate_open.set() await service.stop() async def test_shutdown_timeout_fences_inflight_resume_and_restart_keeps_decision( tmp_path: Path, ) -> None: """A resume outliving grace cannot overwrite the next owner (shutdown). Stop cancels and joins the gated resume before releasing ownership: the task is done-cancelled at stop return, durable state keeps the exact crash shape (interrupted, executing mark, ACTIVE attempt), and after restart recovery fails the run closed the old task cannot flip the decision back to completed or clear the ACTIVE attempt. """ _gate_open.clear() root = tmp_path / "store" server = build_local_static_workflow_server(root, extra_sources=_gate_sources()) await _seed( server, "ask_gate", "ask_gate.default", _interrupt_then_gate_plan("ask_gate"), ["submitted"], ) store = FileScheduleStore(root) store.create_schedule(_one_shot("asker", "ask_gate.default", datetime.now(UTC))) service = _scheduler(server, capacity=1, drain_grace_s=0.05, auto_tick=False) try: await service.start() async with asyncio.timeout(20): while not _run_ids(root): await service.poll_once(datetime.now(UTC)) await asyncio.sleep(0.02) ask_id = _run_ids(root)[0] await _wait_for( lambda: FileRunStore(root).get_run(ask_id).status.value == "interrupted" ) interrupted_checkpoint = FileRunStore(root).get_run(ask_id).latest_checkpoint_id resume_task = asyncio.create_task( server.api.resume_run( run_id=ask_id, resume_payload={}, resume_outcome="submitted" ) ) await _wait_for(lambda: ask_id in service._live_resumes) await _wait_for(lambda: FileRunStore(root).is_executing(ask_id)) # The gate stays closed: the drain must fence, not wait forever. report = await service.stop() assert report.cancelled == 1 assert resume_task.done() assert resume_task.cancelled() # Cancellation truthfulness: crash-shaped durable state, no result. assert FileRunStore(root).get_run(ask_id).status.value == "interrupted" assert FileRunStore(root).is_executing(ask_id) attempt = FileRunStore(root).get_resume_attempt(ask_id) assert attempt is not None and attempt.state == "ACTIVE" assert await _kinds(root, "asker") == ["admitted", "interrupted"] revived = _scheduler(server, capacity=1, auto_tick=False) try: await revived.start() # Restart recovery fails the ambiguous resume closed. assert FileRunStore(root).get_run(ask_id).status.value == "failed" attempt2 = FileRunStore(root).get_resume_attempt(ask_id) assert attempt2 is not None and attempt2.state == "ACTIVE" assert not FileRunStore(root).is_executing(ask_id) assert ( FileRunStore(root).get_run(ask_id).latest_checkpoint_id == interrupted_checkpoint ) assert len(_entries(root, "asker", "failed")) == 1 # The fenced old task is already dead: opening the gate and # polling under the new owner changes nothing. _gate_open.set() await revived.poll_once(datetime.now(UTC)) assert FileRunStore(root).get_run(ask_id).status.value == "failed" attempt3 = FileRunStore(root).get_resume_attempt(ask_id) assert attempt3 is not None and attempt3.state == "ACTIVE" assert _entries(root, "asker", "completed") == [] assert len(_entries(root, "asker", "failed")) == 1 finally: _gate_open.set() await revived.stop() finally: _gate_open.set() await service.stop() async def test_caller_cancellation_reconciles_and_frees_capacity( tmp_path: Path, ) -> None: """Caller cancellation fails one ambiguous resume without fencing capacity.""" _gate_open.clear() root = tmp_path / "store" server = build_local_static_workflow_server(root, extra_sources=_gate_sources()) await _seed( server, "ask_gate", "ask_gate.default", _interrupt_then_gate_plan("ask_gate"), ["submitted"], ) await _seed(server, "const", "const.default", _constant_plan("const"), ["ok"]) store = FileScheduleStore(root) store.create_schedule(_one_shot("asker", "ask_gate.default", datetime.now(UTC))) service = _scheduler(server, capacity=1, auto_tick=False) try: await service.start() async with asyncio.timeout(20): while not _run_ids(root): await service.poll_once(datetime.now(UTC)) await asyncio.sleep(0.02) ask_id = _run_ids(root)[0] await _wait_for( lambda: FileRunStore(root).get_run(ask_id).status.value == "interrupted" ) resume_task = asyncio.create_task( server.api.resume_run( run_id=ask_id, resume_payload={}, resume_outcome="submitted" ) ) await _wait_for(lambda: ask_id in service._live_resumes) await _wait_for(lambda: FileRunStore(root).is_executing(ask_id)) await _wait_for( lambda: ( (attempt := FileRunStore(root).get_resume_attempt(ask_id)) is not None and attempt.state == "ACTIVE" ) ) resume_task.cancel() with pytest.raises(asyncio.CancelledError): await resume_task assert service.running assert not service._live_resumes cancelled = FileRunStore(root).get_run(ask_id) assert cancelled.status.value == "failed" assert not FileRunStore(root).is_executing(ask_id) attempt = FileRunStore(root).get_resume_attempt(ask_id) assert attempt is not None and attempt.state == "ACTIVE" assert len(_entries(root, "asker", "failed")) == 1 with pytest.raises(ValueError, match="ambiguous active resume attempt"): await server.api.resume_run( run_id=ask_id, resume_payload={}, resume_outcome="submitted" ) assert _entries(root, "asker", "completed") == [] store.create_schedule(_one_shot("after", "const.default", datetime.now(UTC))) async with asyncio.timeout(20): while True: await service.poll_once(datetime.now(UTC)) if len(_run_ids(root)) == 2: after_id = [rid for rid in _run_ids(root) if rid != ask_id][0] if FileRunStore(root).get_run(after_id).status.value == "completed": break await asyncio.sleep(0.02) assert await _kinds(root, "asker") == [ "admitted", "interrupted", "failed", ] assert await _kinds(root, "after") == ["admitted", "completed"] finally: _gate_open.set() await service.stop() async def test_scheduled_resume_requested_during_drain_rejected( tmp_path: Path, ) -> None: """New scheduled resumes are rejected once shutdown begins (shutdown). With a free slot available the rejection still fires (drain, not capacity): no dispatch, no ACTIVE attempt, no executing mark. Manual resumes bypass the gate and complete mid-drain. """ _gate_open.clear() root = tmp_path / "store" server = build_local_static_workflow_server(root, extra_sources=_gate_sources()) await _seed( server, "ask1", "ask1.default", _single_interrupt_plan("ask1"), ["submitted"] ) await _seed(server, "gated", "gated.default", _gate_plan("gated"), ["ok"]) await _seed( server, "manual_ask", "manual_ask.default", _single_interrupt_plan("manual_ask"), ["submitted"], ) store = FileScheduleStore(root) store.create_schedule(_one_shot("blocker", "gated.default", datetime.now(UTC))) store.create_schedule(_one_shot("asker", "ask1.default", datetime.now(UTC))) service = _scheduler(server, capacity=2, drain_grace_s=10.0) try: await service.start() await _wait_for(lambda: len(_run_ids(root)) == 2) ids = { FileRunStore(root).get_admission(rid).schedule_id or "": rid for rid in _run_ids(root) } gate_id = ids["blocker"] ask_id = ids["asker"] await _wait_for( lambda: FileRunStore(root).get_run(ask_id).status.value == "interrupted" ) await _wait_for(lambda: FileRunStore(root).is_executing(gate_id)) asked = await server.api.run_deployment( deployment_id="manual_ask.default", workflow_input={} ) assert asked["status"] == "interrupted" manual_id = asked["run_id"] assert isinstance(manual_id, str) # Hold the drain on the gated execution; one slot stays free, so # only the shutdown check can reject the resume below. stop_task = asyncio.create_task(service.stop()) await _wait_for(lambda: service._stopping) assert not stop_task.done() with pytest.raises(ScheduledResumeShutdownError): await server.api.resume_run( run_id=ask_id, resume_payload={}, resume_outcome="submitted" ) assert FileRunStore(root).get_resume_attempt(ask_id) is None assert not FileRunStore(root).is_executing(ask_id) assert FileRunStore(root).get_run(ask_id).status.value == "interrupted" manual_resumed = await server.api.resume_run( run_id=manual_id, resume_payload={}, resume_outcome="submitted" ) assert manual_resumed["status"] == "completed" _gate_open.set() report = await asyncio.wait_for(stop_task, timeout=20.0) assert report.cancelled == 0 assert FileRunStore(root).get_run(gate_id).status.value == "completed" finally: _gate_open.set() await service.stop() async def test_shutdown_timeout_spares_healthy_sibling(tmp_path: Path) -> None: """A fenced resume leaves a healthy sibling's outcome/history intact.""" _gate_open.clear() root = tmp_path / "store" server = build_local_static_workflow_server(root, extra_sources=_gate_sources()) await _seed( server, "ask_gate", "ask_gate.default", _interrupt_then_gate_plan("ask_gate"), ["submitted"], ) await _seed(server, "const", "const.default", _constant_plan("const"), ["ok"]) store = FileScheduleStore(root) store.create_schedule(_one_shot("asker", "ask_gate.default", datetime.now(UTC))) store.create_schedule(_one_shot("sib", "const.default", datetime.now(UTC))) service = _scheduler(server, capacity=1, drain_grace_s=0.05, auto_tick=False) try: await service.start() async with asyncio.timeout(20): while len(_run_ids(root)) < 2: await service.poll_once(datetime.now(UTC)) await asyncio.sleep(0.02) ids = { FileRunStore(root).get_admission(rid).schedule_id or "": rid for rid in _run_ids(root) } ask_id = ids["asker"] sib_id = ids["sib"] await _wait_for( lambda: FileRunStore(root).get_run(ask_id).status.value == "interrupted" ) await _wait_for( lambda: FileRunStore(root).get_run(sib_id).status.value == "completed" ) resume_task = asyncio.create_task( server.api.resume_run( run_id=ask_id, resume_payload={}, resume_outcome="submitted" ) ) await _wait_for(lambda: ask_id in service._live_resumes) report = await service.stop() assert report.cancelled == 1 assert resume_task.cancelled() assert FileRunStore(root).get_run(sib_id).status.value == "completed" assert await _kinds(root, "sib") == ["admitted", "completed"] revived = _scheduler(server, capacity=1, auto_tick=False) try: await revived.start() assert FileRunStore(root).get_run(ask_id).status.value == "failed" assert FileRunStore(root).get_run(sib_id).status.value == "completed" assert await _kinds(root, "sib") == ["admitted", "completed"] assert _entries(root, "asker", "completed") == [] finally: _gate_open.set() await revived.stop() finally: _gate_open.set() await service.stop() async def test_scheduled_resume_records_completion_live(tmp_path: Path) -> None: """A resumed completion lands in occurrence history live (B3). No restart recovery is involved: after the API resume plus ordinary polls, history reads admitted/interrupted/completed, and repeated polling plus a restart never duplicate the entries. """ root = tmp_path / "store" server = build_local_static_workflow_server(root) await _seed( server, "ask1", "ask1.default", _single_interrupt_plan("ask1"), ["submitted"] ) store = FileScheduleStore(root) store.create_schedule(_one_shot("asker", "ask1.default", datetime.now(UTC))) service = _scheduler(server, capacity=2) try: await service.start() await _wait_for(lambda: len(_run_ids(root)) == 1) ask_id = _run_ids(root)[0] await _wait_for( lambda: FileRunStore(root).get_run(ask_id).status.value == "interrupted" ) resumed = await server.api.resume_run( run_id=ask_id, resume_payload={}, resume_outcome="submitted" ) assert resumed["status"] == "completed" assert await _kinds(root, "asker") == ["admitted", "interrupted", "completed"] completed = _entries(root, "asker", "completed") assert len(completed) == 1 assert completed[0]["run_id"] == ask_id assert ( completed[0]["checkpoint_id"] == FileRunStore(root).get_run(ask_id).latest_checkpoint_id ) assert completed[0]["revision"] == 1 await service.poll_once(datetime.now(UTC)) await service.poll_once(datetime.now(UTC)) assert await _kinds(root, "asker") == ["admitted", "interrupted", "completed"] finally: await service.stop() server_b = build_local_static_workflow_server(root) revived = _scheduler(server_b) try: await revived.start() await revived.poll_once(datetime.now(UTC)) assert await _kinds(root, "asker") == ["admitted", "interrupted", "completed"] finally: await revived.stop() async def test_scheduled_resume_records_failure_and_reinterruption_live( tmp_path: Path, ) -> None: """Resumed failures and re-interruptions each record live (B3).""" root = tmp_path / "store" server = build_local_static_workflow_server(root) await _seed( server, "ask_twice", "ask_twice.default", _double_interrupt_plan("ask_twice"), ["submitted"], ) await _seed( server, "ask_fail", "ask_fail.default", _interrupt_then_fail_plan("ask_fail"), ["submitted"], ) store = FileScheduleStore(root) store.create_schedule(_one_shot("twicer", "ask_twice.default", datetime.now(UTC))) store.create_schedule(_one_shot("failer", "ask_fail.default", datetime.now(UTC))) service = _scheduler(server, capacity=2) try: await service.start() await _wait_for(lambda: len(_run_ids(root)) == 2) twice_id = fail_id = "" for rid in _run_ids(root): admission = FileRunStore(root).get_admission(rid) if admission.schedule_id == "twicer": twice_id = rid else: fail_id = rid assert twice_id and fail_id await _wait_for( lambda: FileRunStore(root).get_run(twice_id).status.value == "interrupted" ) await _wait_for( lambda: FileRunStore(root).get_run(fail_id).status.value == "interrupted" ) first = await server.api.resume_run( run_id=twice_id, resume_payload={}, resume_outcome="submitted" ) assert first["status"] == "interrupted" second = await server.api.resume_run( run_id=twice_id, resume_payload={}, resume_outcome="submitted" ) assert second["status"] == "completed" failed = await server.api.resume_run( run_id=fail_id, resume_payload={}, resume_outcome="submitted" ) assert failed["status"] == "failed" assert await _kinds(root, "twicer") == [ "admitted", "interrupted", "interrupted", "completed", ] twicer_page = FileScheduleStore(root).list_occurrences("twicer", limit=100) twicer_ckpts = [ row["checkpoint_id"] for row in cast(list[dict[str, Any]], twicer_page["occurrences"]) if row["kind"] in ("interrupted", "completed") ] assert len(set(twicer_ckpts)) == 3 assert await _kinds(root, "failer") == ["admitted", "interrupted", "failed"] fail_rows = _entries(root, "failer", "failed") assert len(fail_rows) == 1 assert fail_rows[0]["run_id"] == fail_id finally: await service.stop() async def test_paused_and_deleted_schedules_keep_resume_history( tmp_path: Path, ) -> None: """Pausing/deleting never suppresses retained resume history (B3).""" root = tmp_path / "store" server = build_local_static_workflow_server(root) await _seed( server, "ask1", "ask1.default", _single_interrupt_plan("ask1"), ["submitted"] ) store = FileScheduleStore(root) store.create_schedule(_one_shot("pausable", "ask1.default", datetime.now(UTC))) store.create_schedule(_one_shot("doomed", "ask1.default", datetime.now(UTC))) service = _scheduler(server, capacity=2) try: await service.start() await _wait_for(lambda: len(_run_ids(root)) == 2) ids = { FileRunStore(root).get_admission(rid).schedule_id or "": rid for rid in _run_ids(root) } await _wait_for( lambda: all( FileRunStore(root).get_run(rid).status.value == "interrupted" for rid in ids.values() ) ) paused = store.get_schedule("pausable") paused.paused = True store.save_schedule(paused) gone = store.get_schedule("doomed") gone.deleted = True store.save_schedule(gone) for rid in ids.values(): resumed = await server.api.resume_run( run_id=rid, resume_payload={}, resume_outcome="submitted" ) assert resumed["status"] == "completed" assert await _kinds(root, "pausable") == [ "admitted", "interrupted", "completed", ] assert await _kinds(root, "doomed") == ["admitted", "interrupted", "completed"] finally: await service.stop()