diff --git a/docs/README.md b/docs/README.md index 319a0ce0..e7deef2a 100644 --- a/docs/README.md +++ b/docs/README.md @@ -39,6 +39,8 @@ implementation plans are kept for context, not as active instructions. deployments, dependency compatibility, and interrupt limitations. - [`durable_run_operations.md`](durable_run_operations.md): `run_deployment`, `inspect_run`, bounded trace reads, and `resume_run` behavior. +- [`deployment_scheduling.md`](deployment_scheduling.md): opt-in schedule + administration, occurrence inspection, and server lifecycle behavior. - [`workflow_drafts.md`](workflow_drafts.md): LLM/human draft authoring format above raw workflow plans. diff --git a/docs/current_roadmap.md b/docs/current_roadmap.md index e0eab0c5..4c69d4b0 100644 --- a/docs/current_roadmap.md +++ b/docs/current_roadmap.md @@ -149,6 +149,16 @@ The active sequence can assume these foundations: - Python client reconstruction of capabilities, artifacts, deployments, and runs through the API +The active sequence can also assume deployment scheduling: an opt-in +same-server scheduler starts ordinary deployment runs without a connected +client (one-shot and recurring cron, durable admission, coalesced +missed-start recovery, bounded parallel runs, occurrence inspection, and +API/Python-client administration). Scheduling is disabled by default and +enabled per server. The current contract is +[`deployment scheduling`](superpowers/specs/2026-09-08-deployment-scheduling-design.md); +operator usage is under +[`deployment scheduling operations`](deployment_scheduling.md). + The current foreach return contract is [`foreach back-edge design`](superpowers/specs/2026-09-04-foreach-back-edge-design.md). The current context contract is diff --git a/docs/deployment_scheduling.md b/docs/deployment_scheduling.md new file mode 100644 index 00000000..06de1972 --- /dev/null +++ b/docs/deployment_scheduling.md @@ -0,0 +1,198 @@ +# Deployment Scheduling Operations + +Schedules start ordinary deployment runs without a connected client. The +scheduler is opt-in, lives in the workflow server, and introduces no +workflow node type. It is not the core runtime's frame scheduler. + +Current contract: +[`deployment scheduling spec`](superpowers/specs/2026-09-08-deployment-scheduling-design.md). + +## Enablement + +Scheduling is disabled by default. Enable it for a local/static server +with the config section or the CLI flag (MCP-backed servers reject it): + +```json +{"server": {"scheduler": {"enabled": true}}} +``` + +```bash +wf-rpc-server --store-root .wf_store --enable-scheduler +``` + +Tuning (`server.scheduler`): `poll_interval_s` (default 1.0), +`max_concurrent_runs` (default 4, the execution-slot bound), +`drain_grace_s` (default 30.0). Schedule data lives at +`/schedules` next to run data; one lock file at +`/scheduler.lock` proves exclusive ownership. A second +scheduler over the same stores is rejected; shut the first down before +starting another. + +## Mental model: schedule vs occurrence vs run + +- A **schedule** is a durable definition: deployment, trigger + (one-shot or cron), input bindings, overlap/misfire policies, and a + revision. Edits bump the revision and affect only future admissions. +- An **occurrence** is one resolved calendar instant, identified by + `(schedule_id, resolved UTC instant)`. An occurrence is immutable once + admitted and is never replayed. +- A **run** is the execution of one admitted occurrence, with a + store-backed `run-000001`-style id, a pinned input/artifact snapshot, + and a stopped checkpoint when it stops. + +## Triggers and time zones + +One-shot timestamps must include an offset. Cron uses five-field Unix +expressions (`0` and `7` both mean Sunday; day-of-month/day-of-week +match with OR) plus an explicit IANA time zone, default UTC. +Occurrence instants persist as UTC; the definition keeps the zone name. +`croniter` owns calendar resolution including DST gaps and folds; there +is no custom calendar filtering. Reject invalid zones and naive times. + +## Overlap + +Overlap is per schedule, not per deployment. With default +`overlap="skip"`, an unfinished scheduled run — including an +interrupted run awaiting input — blocks another occurrence from that +schedule. Manual runs and other schedules never participate. +`overlap="parallel"` admits independent runs up to a required positive +`max_active_runs`; admitted, running, and interrupted runs all count, +and lowering the limit blocks new admission without cancelling existing +runs. A known completed or failed run releases overlap. + +## Misfire + +Default `misfire="skip"` drops missed starts (past the 60-second +configurable lateness allowance). `misfire="latest"` retains at most +one latest unadmitted candidate per schedule and runs it as soon as +overlap and capacity allow; it never replays a burst. Creation never +backfills time before the revision: the consumed watermark starts at +creation. + +## Pause + +Pausing stops future admission, not an active run. Resume selects the +next future occurrence; paused times are not replayed. Pause is not +downtime: the paused span is consumed, so downtime catch-up never +resurrects it. + +## Deletion + +Deleting a schedule stops future admission. Existing runs and +occurrence history survive, including their schedule identity; active +runs continue and may still complete or resume against the retained +history. Schedule ids are never reused. + +## Failure and restart + +An abandoned in-flight execution (the process died mid-run) becomes +failed without retrying its occurrence; the failure discloses that +external effects may already have occurred. A durably interrupted run +is not abandoned: it stays resumable and keeps occupying its +schedule's overlap slot across restarts. Every resume marks a +store-backed attempt first, so recovery can tell a fresh result from a +stale checkpoint and fail the ambiguous case closed instead of +retrying it. Corrupt or contradictory records fail closed with +diagnostics and block the schedule rather than clearing overlap. + +On shutdown the server stops admission first and drains active tasks +within the grace period; anything still running keeps its executing +mark, and the next startup recovery abandons it truthfully. + +## Occurrence inspection + +`list_schedule_occurrences` pages stored history (`pending`, +`coalesced`/`superseded`, `skipped-overlap`, `skipped-misfire`, +`preflight-rejected`, `admitted`/`running`, `interrupted`, +`completed`, `failed`, `exhausted`) oldest-first with a `next_cursor`. +A currently held (unadmitted) candidate is synthesized as a `pending` +row at the top of the first page; once admitted, the durable +`admitted` entry replaces it. While a candidate is held, the first +page may carry one row more than `limit`. + +## Administration surface + +Python API (`server.api.schedules`), JSON-RPC (`workflow.schedules.*`), +and the Python client (`App` schedule methods) share these operations: + +- `create_schedule` (`workflow.schedules.create`): caller-chosen id; + trigger, deployment, binding, and sample-schema checks. +- `get_schedule` (`workflow.schedules.get`): full definition payload. +- `list_schedules` (`workflow.schedules.list`): deleted excluded + unless asked. +- `update_schedule` (`workflow.schedules.update`): `expected_revision` + required; provided fields only; `None` means unpatched. +- `pause_schedule` (`workflow.schedules.pause`): idempotent; clears + candidates, consumes the span. +- `resume_schedule` (`workflow.schedules.resume`): idempotent; resumes + from the next future instant. +- `delete_schedule` (`workflow.schedules.delete`): soft delete; runs + and history survive. +- `list_schedule_occurrences` + (`workflow.schedules.occurrences.list`): cursor pages, `limit` 1–100, + live pending synthesis. + +`inspect_run` also reads admitted (checkpoint-less) runs: status +`admitted` with no fabricated trace, output, or checkpoint. + +## Hypothetically used as follows + +EXECUTABLE example (runs in CI as +`tests/examples/test_scheduled_deployment_example.py`; run it with +`uv run pytest -q tests/examples/test_scheduled_deployment_example.py`): + +```python +server = build_local_static_workflow_server(root, schedules=True) +await server.api.create_artifact_from_plan( + artifact_id="scheduled_hello", version=1, ..., +) +await server.api.save_deployment({...}) +due = datetime.now(UTC) + timedelta(seconds=0.5) +await server.api.schedules.create_schedule( + schedule_id="hello-once", + deployment_id="scheduled_hello.default", + trigger={"kind": "oneshot", "at": due.isoformat()}, +) +service = build_scheduler_service(server, SchedulerServiceConfig(...)) +await service.start() +# ... the service admits the occurrence and completes the run ... +inspected = await server.api.inspect_run(run_id=run_id) +assert inspected["output"]["result"] == "hello on a schedule" +page = await server.api.schedules.list_schedule_occurrences( + schedule_id="hello-once" +) +await service.stop() +``` + +ILLUSTRATIVE example (not executed; shows a cron week with an operator +pause — same calls, longer horizons): + +```python +# Monday: an hourly report, latest-wins catch-up, at most two at once. +await schedules.create_schedule( + schedule_id="hourly-report", + deployment_id="report.default", + trigger={"kind": "cron", "expression": "0 * * * *", "timezone": "UTC"}, + misfire="latest", + overlap="parallel", + max_active_runs=2, +) +# Friday: pause for maintenance; Monday: resume from the next hour. +await schedules.pause_schedule(schedule_id="hourly-report") +await schedules.resume_schedule(schedule_id="hourly-report") +# A bad edit is rejected without touching the running definition: +await schedules.update_schedule( + schedule_id="hourly-report", expected_revision=1, ... +) +``` + +## Known limitations + +- One composition's stores nested inside another live composition's + store subtree (without sharing its identical roots) is unsupported + operator error; shared-store cross layouts are rejected outright. +- Manual runs and resumes bypass scheduler capacity by design; + capacity governs scheduled dispatch only. +- A set `max_steps` budget cannot be cleared back to unset through + update (recreate the schedule for an unbounded budget). +- MCP-backed servers reject scheduler enablement for now. diff --git a/docs/superpowers/plans/2026-09-09-deployment-scheduling-implementation-plan.md b/docs/historical/superpowers/plans/2026-09-09-deployment-scheduling-implementation-plan.md similarity index 93% rename from docs/superpowers/plans/2026-09-09-deployment-scheduling-implementation-plan.md rename to docs/historical/superpowers/plans/2026-09-09-deployment-scheduling-implementation-plan.md index 318c0e16..271b8063 100644 --- a/docs/superpowers/plans/2026-09-09-deployment-scheduling-implementation-plan.md +++ b/docs/historical/superpowers/plans/2026-09-09-deployment-scheduling-implementation-plan.md @@ -16,6 +16,48 @@ Scheduling policies remain settled. Calendar behavior follows the subsequent approved simplification: croniter owns DST resolution. Findings below record the missing implementation seams and the completed verification work. +## Completion status (2026-09-09; plan retired to historical/) + +All gates and tasks completed on branch `opencode/sched-verify-plan`: + +- R0, R1, R2, R3: passed during phased implementation. +- R4 (fault-injection review): passed in three waves — wave 2 bound + ownership to store composition and stabilized recovery failure; wave 3 + added one canonical lock identity plus checkpoint coherence and + ordering authority; wave 4 closed the shared-store overlap hole + (sibling-distinct identity, overlapping layouts rejected). +- T01–T11: implemented per phase (calendar, expressions, admission, + scheduler core, resume safety, recovery, ownership). +- T12: opt-in server lifecycle with bounded real execution, drain, and + subprocess-tested death paths; independent review passed. +- T13: administration API + JSON-RPC surface + Python client; an + independent review gated it on torn-admin ordering and admin/poll + staleness, both fixed (crash-safe admin ordering, creation watermark, + per-schedule poll freshness) and re-reviewed to pass. +- T14: this completion record; spec updated in place; roadmap and + user docs updated; disposable probes deleted after production-test + equivalence was verified item by item (calendar Part A mirrored in + `test_calendar_adapter.py`; APScheduler rejected-candidate evidence + preserved in the design spec; expression pins mirrored or superseded + by T03/T04 plus core tests; all 31 state-model behaviors mapped to + production tests, adding overlap-independence and store-backed + run-identity pins where no equivalent existed). + +Probe retirement map: `probes/deployment_scheduling_verify/` deleted. +`test_calendar_probe.py` Part A is subsumed by +`tests/scheduling/test_calendar_adapter.py`; Part B (APScheduler gap +phantom + fold replay strict xfails) is preserved as narrative evidence +in the design spec, not as runnable tests (it needs an isolated env +with `apscheduler` installed). `test_expression_contract_probe.py` is +subsumed by `tests/scheduling/test_schedule_expressions.py`, +`tests/core/test_input_sources.py`, and core strict-JSON tests, except +the pre-T03 no-seam observation (deliberately superseded) and one +unowned `InputBinding` micro-pin of untouched core code (noted, not +ported). `test_schedule_state_model.py` is subsumed by +`tests/scheduling/` (matrix, coalescing, slots, pause/edit/delete, +restart, downtime, fairness, fault injection, ownership, corrupt +handling) plus the two pins added at retirement. + ## Gate 1 — Spec audit against actual code Three parallel audit sweeps (expression bindings, deployment invocation, diff --git a/docs/superpowers/specs/2026-09-08-deployment-scheduling-design.md b/docs/superpowers/specs/2026-09-08-deployment-scheduling-design.md index 2f08c5d2..dccdfd7f 100644 --- a/docs/superpowers/specs/2026-09-08-deployment-scheduling-design.md +++ b/docs/superpowers/specs/2026-09-08-deployment-scheduling-design.md @@ -1,6 +1,9 @@ # Deployment Scheduling -Status: draft for review; not implemented. +Status: implemented (slices T01–T14, reviews R0–R4 passed); this document +remains the current contract. The implementation plan that built it is +archived at +`docs/historical/superpowers/plans/2026-09-09-deployment-scheduling-implementation-plan.md`. ## Purpose and scope @@ -247,6 +250,16 @@ anything else under an active attempt is stale and fails closed without retry. Recovery must distinguish both cases from an interruption that was merely waiting across server restart. +Implemented checkpoint authority: a checkpoint decides a summary only when +its content is coherent — same run identity, checkpoint id of the form +`{run_id}.{sequence:06d}`, runtime state decodable, and decoded stopped +status equal to the outer reason. A durable failed decision is superseded +only by a present, coherent, strictly newer checkpoint with matching +attempt provenance; a missing or older referenced checkpoint keeps the +decision (noted, history preserved) and never rolls a failed run back to +interrupted. A genuinely newer coherent result still repairs a torn +summary, and repeated recovery is silent and stable. + The process must own the store exclusively before recovery. Enforce scheduler ownership with a held cross-process lock, not a stale PID file or a lease that can expire while the old owner still runs. Unsupported locking must reject @@ -254,6 +267,17 @@ scheduler startup. Other processes mutating/resuming the same store remain unsupported. This does not upgrade the rest of the file stores to multi-writer safety or claim exactly-once external effects. +Implemented lock identity: one store composition has exactly one lock file, +at the deepest composition root containing every store root as itself or a +direct child (`canonical_lock_root`). Identical, sibling, and nested roots +covering the same store files contend on that one file; ancestor locks prove +nothing and are rejected, and cross pairs reusing one protected store with a +different partner have no lock identity at all and are rejected before any +write. The acquired identity is frozen at acquisition, so mutating the +handle's root afterwards cannot redirect authority. The server layout points +both stores at the composition root itself (run data at `/runs`, +schedules at `/schedules`, one lock at `/scheduler.lock`). + ## Lifecycle, administration, and resource bounds Expose create/get/list/update/pause/resume/delete and paginated occurrence @@ -290,6 +314,43 @@ grace period. Record cancellation/failure when possible; abrupt termination uses startup recovery. Paused/deleted schedule definitions must not prevent run completion or resume from updating retained occurrence history. +Implemented service and administration surface. The opt-in same-server +scheduler (`SchedulerService`, enabled by `server.scheduler.enabled` or +`wf-rpc-server --enable-scheduler`; local/static servers only) acquires +canonical ownership, recovers without executing, then ticks calendar +polling without blocking on long workflows: each dispatch spawns exactly +one bounded execution task behind the async-completion seam, and the +scheduler's own capacity gate is the execution-slot bound (an executing +run keeps its slot until it stops). Shutdown stops admission, drains +within `drain_grace_s`, leaves unfinished work under its executing mark +for startup recovery to abandon truthfully, and releases ownership last. +A failed startup releases the lock and raises. + +Administration (`WorkflowApi` schedules methods, `workflow.schedules.*` +RPC, Python client): `create/get/list/update/pause/resume/delete_schedule` +plus paginated `list_schedule_occurrences`. Creation validates the +trigger, the deployment, the binding shapes, and a sample-occurrence +resolution against the pinned root schema, and starts the consumed +watermark at creation (no pre-creation backfill). Updates are +revision-checked; pause/resume/delete mirror the poll-loop transitions; +all mutating admin ops clear the old revision's unadmitted work and +advance the watermark BEFORE the revision bump or flag flip lands, so a +crash can only leave the op unapplied (retryable), never a new revision +that backfills. Occurrence pages carry the stored history plus a live +held-candidate `pending` synthesis on the first page. Manual runs and +resumes bypass scheduler capacity by design (unchanged API behavior); +scheduler capacity governs scheduled dispatch only, and a scheduled +interrupted run resumed manually reconciles its terminal history through +recovery. + +Known limitations: pointing one composition's stores inside another live +composition's store subtree (without sharing its identical roots) is +unsupported operator error; `max_steps: None` means "unpatched" on +update (a set budget cannot be cleared back to unset); a first +occurrence page may carry one row more than `limit` while a candidate is +held; calendar iteration within a tick may use the tick-start source +(trigger edits take effect on the next tick). + ## Verification gates Use injected clocks and controlled executors, not real-time sleeps: diff --git a/docs/wf_api_architecture.md b/docs/wf_api_architecture.md index 48590ca9..269a4e75 100644 --- a/docs/wf_api_architecture.md +++ b/docs/wf_api_architecture.md @@ -233,6 +233,7 @@ WorkflowApiSurface WorkflowArtifactSurface WorkflowDeploymentSurface WorkflowRunSurface + WorkflowScheduleSurface ``` The domain services below are the process-local implementation pieces, not the diff --git a/docs/wf_cli.md b/docs/wf_cli.md index 244eac38..98a95fa0 100644 --- a/docs/wf_cli.md +++ b/docs/wf_cli.md @@ -109,6 +109,17 @@ small and avoids dumping arbitrary MCP resource payloads. registry. `--store-root` is for the local/static server path and cannot be combined with `--mcp-config`. +Opt in to deployment scheduling on a local/static server (MCP-backed +servers reject it for now): + +```bash +wf-rpc-server --store-root .wf_store --enable-scheduler +``` + +or set `server.scheduler.enabled` (plus optional `poll_interval_s`, +`max_concurrent_runs`, `drain_grace_s`) in the neutral config. See +[`deployment scheduling operations`](deployment_scheduling.md). + `admin registry` shows desired persisted source entries. It is separate from workflow artifacts and deployments, so it can be empty even when the server has runtime sources and saved workflows. diff --git a/probes/deployment_scheduling_verify/README.md b/probes/deployment_scheduling_verify/README.md deleted file mode 100644 index bf185a8d..00000000 --- a/probes/deployment_scheduling_verify/README.md +++ /dev/null @@ -1,40 +0,0 @@ -# DISPOSABLE VERIFICATION PROBE — NOT PRODUCTION CODE - -This directory holds throwaway verification probes for the -deployment-scheduling slice -(`docs/superpowers/specs/2026-09-08-deployment-scheduling-design.md`). -They are evidence-gathering scripts, not a production implementation: - -- `test_calendar_probe.py` — calendar-library probe (16 passed, 2 strict - xfailed: Part A pins the croniter contract, Part B keeps APScheduler - rejected-candidate evidence). Requires `croniter==6.2.4`, `tzdata` - (`apscheduler==3.11.0` for Part B only) on Python 3.14. - Run it from an isolated project (NOT the repo env, which does not - depend on croniter; the file `importorskip`s itself elsewhere): - - ```powershell - $probe = "C:\Users\Admin\AppData\Local\Temp\opencode\sched-cal-probe" - $file = "\probes\deployment_scheduling_verify\test_calendar_probe.py" - uv run --project $probe python -m pytest $file -q -p no:cacheprovider ` - -o addopts="" - ``` - -- `test_schedule_state_model.py` — pure-stdlib reference model of the - scheduling state machine (25 tests: overlap x misfire, coalescing, - slots, pause, faults, ownership). Runs in the repo env: - - ```powershell - uv run pytest -q probes/deployment_scheduling_verify/test_schedule_state_model.py - ``` - -- `test_expression_contract_probe.py` — pins the current input-expression - contract for the Phase 1 implementer (8 tests). Runs in the repo env: - - ```powershell - uv run pytest -q probes/deployment_scheduling_verify/test_expression_contract_probe.py - ``` - -Do not import these probes from `src/`. Do not copy their logic into -production without going through the sequenced implementation plan at -`docs/superpowers/plans/2026-09-09-deployment-scheduling-implementation-plan.md`. -Delete this directory once the slice is implemented. diff --git a/probes/deployment_scheduling_verify/test_calendar_probe.py b/probes/deployment_scheduling_verify/test_calendar_probe.py deleted file mode 100644 index b1919708..00000000 --- a/probes/deployment_scheduling_verify/test_calendar_probe.py +++ /dev/null @@ -1,413 +0,0 @@ -# DISPOSABLE CALENDAR-LIBRARY PROBE — NOT PRODUCTION CODE. -# See README.md in this directory. Requires croniter==6.2.4 and tzdata -# (plus apscheduler==3.11.0 for the Part B rejected-candidate evidence) -# on Python 3.14 (isolated env, not repo env). -"""Pin the croniter 6.2.4 calendar contract for deployment scheduling. - -Policy: croniter owns calendar calculation, including DST resolution. -The adapter consumes it thinly — convert the query instant into the -schedule's named zone, ask for the next/previous occurrence, convert -the result to UTC — and applies no correction of its own. Each test -prints OBSERVED lines pinning 6.2.4 behavior. After any -calendar-dependency upgrade a failure IS the finding: re-probe, do not -hand-roll around it. - -Part A pins the chosen-library contract (shipping acceptance). -Part B preserves rejected-candidate evidence (NOT acceptance). - -Run: - uv run --project python -m pytest - -q -p no:cacheprovider -o addopts="" -""" - -from __future__ import annotations - -import time -from datetime import datetime, timedelta, timezone - -import pytest - -# This probe runs ONLY in the isolated calendar env (see README.md). -# Skip — do not error — when collected by a default repo test run. -pytest.importorskip("apscheduler", reason="isolated calendar-probe env only") -pytest.importorskip("croniter", reason="isolated calendar-probe env only") - -from apscheduler.triggers.cron import CronTrigger # noqa: E402 (Part B only) -from croniter import ( # noqa: E402 - CroniterBadCronError, - CroniterBadDateError, -) -from croniter import ( - croniter as Croniter, -) - -try: # noqa: E402 - from zoneinfo import ZoneInfo -except ImportError: # pragma: no cover - from backports.zoneinfo import ZoneInfo # type: ignore[no-redef] - -from importlib.metadata import version as _pkg_version # noqa: E402 - -import apscheduler # noqa: E402 - -CRONITER_VERSION = _pkg_version("croniter") - -print(f"croniter=={CRONITER_VERSION}") -print(f"apscheduler=={apscheduler.__version__} (rejected candidate, Part B only)") - -UTC = timezone.utc -HCMC = ZoneInfo("Asia/Ho_Chi_Minh") -NYC = ZoneInfo("America/New_York") - -# 2026 DST transitions (US): spring forward 2026-03-08 02:00 -> 03:00, -# fall back 2026-11-01 02:00 -> 01:00. - - -def utc(*args) -> datetime: - return datetime(*args, tzinfo=UTC) - - -# --------------------------------------------------------------------------- -# Part A — chosen-library contract (croniter; shipping acceptance) -# --------------------------------------------------------------------------- - - -def test_versions_pinned(): - print(f"OBSERVED croniter version={CRONITER_VERSION}") - print(f"OBSERVED apscheduler.__version__={apscheduler.__version__}") - # Exact pin: a calendar-dependency upgrade must fail here deliberately, - # forcing a re-probe before any pin update. - assert CRONITER_VERSION == "6.2.4" - assert apscheduler.__version__.startswith("3.") - - -def test_utc_daily_next_is_strictly_increasing_and_unique(): - it = Croniter("30 9 * * *", utc(2026, 9, 1, 0, 0, 0)) - seen: set[datetime] = set() - prev = None - for _ in range(10): - nxt = it.get_next(datetime) - assert nxt.tzinfo is not None - as_utc = nxt.astimezone(UTC) - assert as_utc not in seen, f"duplicate UTC instant {as_utc}" - seen.add(as_utc) - if prev is not None: - assert as_utc > prev, "not strictly increasing" - print(f"OBSERVED utc-daily next={as_utc.isoformat()}") - prev = as_utc - assert len(seen) == 10 - - -def test_hcmc_daily_converts_and_hourly_counts(): - # Asia/Ho_Chi_Minh is UTC+7 with no DST. 09:30 local == 02:30 UTC. - nxt = Croniter("30 9 * * *", datetime(2026, 9, 1, 0, 0, tzinfo=HCMC)).get_next( - datetime - ) - as_utc = nxt.astimezone(UTC) - print(f"OBSERVED hcmc next local={nxt.isoformat()} utc={as_utc.isoformat()}") - assert (as_utc.hour, as_utc.minute) == (2, 30) - assert as_utc.date().isoformat() == "2026-09-01" - # No-DST zone: hourly occurrence count over two 48-hour windows - # (March and November, starting off-tick) must be exactly 48 each. - for label, start in ( - ("march", utc(2026, 3, 7, 17, 30, 0)), - ("november", utc(2026, 10, 31, 17, 30, 0)), - ): - it = Croniter("0 * * * *", start.astimezone(HCMC)) - count = 0 - while True: - hit = it.get_next(datetime).astimezone(UTC) - if hit >= start + timedelta(hours=48): - break - count += 1 - assert count < 60 - print(f"OBSERVED hcmc hourly count window={label} count={count}") - assert count == 48, f"{label}: no-DST zone must yield exactly 48, got {count}" - - -def test_expression_forms_wildcard_step_list_range(): - """Wildcards, steps, lists, and ranges resolve through plain - get_next on an ordinary day (no DST involved).""" - cases = { - "* * * * *": (datetime(2026, 9, 8, 0, 0, tzinfo=UTC), "2026-09-08T00:01:00"), - "5/15 * * * *": (datetime(2026, 9, 8, 0, 0, tzinfo=UTC), "2026-09-08T00:05:00"), - "0,30 1-2 * * *": ( - datetime(2026, 9, 8, 0, 0, tzinfo=UTC), - "2026-09-08T01:00:00", - ), - "*/20 1-3 * * *": ( - datetime(2026, 9, 8, 1, 50, tzinfo=UTC), - "2026-09-08T02:00:00", - ), - } - for expr, (start, want) in cases.items(): - nxt = Croniter(expr, start).get_next(datetime) - print( - f"OBSERVED croniter {expr!r} from {start.isoformat()} -> {nxt.isoformat()}" - ) - assert nxt.astimezone(UTC).isoformat() == want + "+00:00" - - -def test_weekday_names_and_numbers_unix(): - """Pin croniter 6.2.4 weekday dialect (re-probe on upgrade). Unix - convention: numeric 0 AND 7 mean Sunday; 1/mon mean Monday.""" - monday = datetime(2026, 9, 7, 0, 0, tzinfo=UTC) # a Monday - cases = { - "0 12 * * 0": 6, # Sunday - "0 12 * * 7": 6, # Sunday (alias) - "0 12 * * 1": 0, # Monday - "0 12 * * mon": 0, - "0 12 * * sun": 6, - } - for expr, want_wd in cases.items(): - nxt = Croniter(expr, monday).get_next(datetime) - print( - f"OBSERVED croniter {expr!r} -> {nxt.date().isoformat()} weekday={nxt.weekday()}" - ) - assert nxt.weekday() == want_wd, f"{expr}: want weekday={want_wd}" - assert Croniter("0 12 * * 0", monday).get_next(datetime) == Croniter( - "0 12 * * sun", monday - ).get_next(datetime) - - -def test_dom_dow_day_or_true_is_selected(): - """Day-of-month/day-of-week uses croniter's standard day_or=True - (Unix OR). The alternative (AND) is pinned for reference only.""" - start = datetime(2026, 9, 1, 0, 0, tzinfo=UTC) - it_or = Croniter("0 12 13 * fri", start, day_or=True) - hits_or = [it_or.get_next(datetime).date().isoformat() for _ in range(4)] - print(f"OBSERVED croniter day_or=True hits={hits_or}") - # Fridays plus the 13th (a Sunday): classic Unix OR. - assert hits_or == ["2026-09-04", "2026-09-11", "2026-09-13", "2026-09-18"] - it_and = Croniter("0 12 13 * fri", start, day_or=False) - first_and = it_and.get_next(datetime).date() - print( - f"OBSERVED croniter day_or=False first={first_and.isoformat()} (not selected)" - ) - assert (first_and.day, first_and.weekday()) == (13, 4) - - -def test_zone_conversion_is_the_callers_job(): - """croniter iterates in the tz of the supplied datetime and performs - no conversion: a UTC start yields UTC results. The adapter converts - now into the schedule zone before querying and back to UTC after.""" - utc_start = Croniter( - "30 2 * * *", datetime(2026, 3, 7, 12, 0, tzinfo=UTC) - ).get_next(datetime) - print( - f"OBSERVED croniter utc-start gap next={utc_start.isoformat()} (no zone conversion)" - ) - assert utc_start.tzinfo is UTC - hcmc_start = Croniter( - "30 9 * * *", datetime(2026, 9, 1, 0, 0, tzinfo=HCMC) - ).get_next(datetime) - print(f"OBSERVED croniter hcmc-start next={hcmc_start.isoformat()}") - assert hcmc_start.utcoffset() == timedelta(hours=7) - - -def test_gap_nonexistent_wall_times_resolve_forward(): - """Observed 6.2.4 gap behavior: daily 02:30 on 2026-03-08 (a wall - time that never existed in America/New_York) resolves to 03:00-04:00 - the same day. That resolution IS the occurrence — there is no - separate validity concept and no day is skipped over.""" - first = Croniter("30 2 * * *", datetime(2026, 3, 7, 12, 0, tzinfo=NYC)).get_next( - datetime - ) - print(f"OBSERVED croniter gap-day daily-0230 resolves={first.isoformat()}") - assert first.isoformat() == "2026-03-08T03:00:00-04:00" - following = Croniter("30 2 * * *", first).get_next(datetime) - print(f"OBSERVED croniter gap following={following.isoformat()}") - assert following.isoformat() == "2026-03-09T02:30:00-04:00" - - -def test_gap_per_minute_stream_jumps_forward(): - """Per-minute iteration across the spring gap jumps 01:59 EST - straight to 03:00 EDT: no 02:xx wall times are emitted, UTC is - strictly increasing and unique.""" - it = Croniter("* * * * *", datetime(2026, 3, 8, 0, 0, tzinfo=NYC)) - seq = [it.get_next(datetime) for _ in range(300)] - us = [d.astimezone(UTC) for d in seq] - assert all(b > a for a, b in zip(us, us[1:])) - assert len(set(us)) == len(us) - gap_day = [d for d in seq if d.date().isoformat() == "2026-03-08"] - assert all(d.hour != 2 for d in gap_day), "no 02:xx wall times emitted" - assert "2026-03-08T01:59:00-05:00" in (d.isoformat() for d in gap_day) - assert "2026-03-08T03:00:00-04:00" in (d.isoformat() for d in gap_day) - print( - "OBSERVED per-minute gap jump 01:59-05:00 -> 03:00-04:00, " - f"{len(gap_day)} gap-day hits" - ) - - -def test_gap_backward_queries_return_library_resolution(): - """Backward queries resolve the same way: get_prev on the gap day - returns the forward-resolved 03:00-04:00, not a skipped-over day. - Latest-missed on a gap day is that resolved instant.""" - cases = { - "2026-03-08T12:00": "2026-03-08T03:00:00-04:00", - "2026-03-09T00:00": "2026-03-08T03:00:00-04:00", - "2026-03-07T12:00": "2026-03-07T02:30:00-05:00", - "2026-03-09T12:00": "2026-03-09T02:30:00-04:00", - } - for start_iso, want in cases.items(): - start = datetime.fromisoformat(start_iso).replace(tzinfo=NYC) - got = Croniter("30 2 * * *", start).get_prev(datetime) - print(f"OBSERVED croniter gap get_prev from {start_iso} -> {got.isoformat()}") - assert got.isoformat() == want - - -def test_fold_repeated_times_are_distinct_utc(): - """Both 01:30s on 2026-11-01 occur as distinct UTC instants - (05:30Z EDT, then 06:30Z EST).""" - it = Croniter("30 1 * * *", datetime(2026, 10, 31, 12, 0, tzinfo=NYC)) - first = it.get_next(datetime) - second = it.get_next(datetime) - print( - f"OBSERVED croniter fold first={first.isoformat()} second={second.isoformat()}" - ) - assert first.astimezone(UTC) == utc(2026, 11, 1, 5, 30) - assert second.astimezone(UTC) == utc(2026, 11, 1, 6, 30) - - -def test_fold_per_minute_unique_increasing(): - """Per-minute iteration across the fall fold emits both fold hours - with strictly increasing unique UTC instants.""" - it = Croniter("* * * * *", datetime(2026, 10, 31, 20, 0, tzinfo=NYC)) - seq = [it.get_next(datetime) for _ in range(560)] - us = [d.astimezone(UTC) for d in seq] - print(f"OBSERVED croniter fold-span count={len(us)} last={us[-1].isoformat()}") - assert all(b > a for a, b in zip(us, us[1:])), "must be strictly increasing" - assert len(set(us)) == len(us), "UTC identities must be unique" - iso = {d.isoformat() for d in us} - assert "2026-11-01T05:30:00+00:00" in iso, "first 01:30 (EDT) must occur" - assert "2026-11-01T06:30:00+00:00" in iso, "second 01:30 (EST) must occur" - - -def test_latest_missed_after_years_of_downtime(): - """Per-minute schedule, ~3 years of downtime: get_prev answers the - latest missed occurrence directly — no enumeration of missed years.""" - now = utc(2026, 9, 8, 12, 0, 0) - t0 = time.perf_counter() - latest = Croniter("* * * * *", now).get_prev(datetime) - elapsed = time.perf_counter() - t0 - print( - f"OBSERVED croniter get_prev({now.isoformat()})={latest.isoformat()} " - f"elapsed={elapsed:.4f}s" - ) - assert latest.tzinfo is not None, "get_prev must preserve tz-awareness" - assert elapsed < 5, "latest-missed lookup must be bounded" - assert latest <= now - assert (now - latest) < timedelta(minutes=2) - now_nyc = datetime(2026, 9, 8, 12, 0, tzinfo=NYC) - latest_nyc = Croniter("* * * * *", now_nyc).get_prev(datetime) - print( - f"OBSERVED croniter nyc get_prev now={now_nyc.isoformat()} " - f"latest={latest_nyc.isoformat()}" - ) - assert latest_nyc.tzinfo is not None - assert latest_nyc <= now_nyc - assert (now_nyc - latest_nyc) < timedelta(minutes=2) - assert latest_nyc.astimezone(UTC).isoformat() == "2026-09-08T15:59:00+00:00" - - -def test_iteration_is_exclusive_both_directions(): - """Both get_next and get_prev are exclusive of their start: a query - from exactly a due instant returns the neighboring occurrence, not - the instant itself. Adapter consequence (T01, not enforced here): - query forward from the last-consumed instant and compare catch-up - results against the same watermark, so a due occurrence is admitted - exactly once — neither missed by exclusivity nor admitted twice.""" - due = utc(2026, 9, 8, 13, 0, 0) - fwd = Croniter("0 13 * * *", due).get_next(datetime) - back = Croniter("0 13 * * *", due).get_prev(datetime) - print(f"OBSERVED exclusive get_next(due)={fwd.isoformat()}") - print(f"OBSERVED exclusive get_prev(due)={back.isoformat()}") - assert fwd.astimezone(UTC) == utc(2026, 9, 9, 13, 0, 0) - assert (back.year, back.month, back.day, back.hour, back.minute) == ( - 2026, - 9, - 7, - 13, - 0, - ) - just_before = Croniter("0 13 * * *", due - timedelta(seconds=1)).get_next(datetime) - assert just_before.astimezone(UTC) == due, "a tick just before due admits it" - - -def test_impossible_schedule_raises_promptly(): - """February 30th never occurs: the library raises its documented - exhaustion error promptly instead of searching forever. The adapter - maps this to exhausted — exhaustion must never look like progress.""" - start = time.perf_counter() - with pytest.raises(CroniterBadDateError): - Croniter("0 12 30 2 *", utc(2026, 1, 1)).get_next(datetime) - elapsed = time.perf_counter() - start - print( - f"OBSERVED impossible-schedule raised CroniterBadDateError elapsed={elapsed:.3f}s" - ) - assert elapsed < 5, f"search must be bounded, took {elapsed:.1f}s" - - -def test_bad_expressions_rejected_naive_passes_through(): - """Malformed expressions fail at construction with a documented - error. Naive datetimes are NOT rejected by the library — they pass - straight through — so the adapter rejects naive itself per spec.""" - with pytest.raises(CroniterBadCronError): - Croniter("nonsense", utc(2026, 1, 1)) - print("OBSERVED malformed expression rejected with CroniterBadCronError") - naive = Croniter("* * * * *", datetime(2026, 9, 8, 12, 0)).get_next(datetime) - print(f"OBSERVED naive start passes through tzinfo={naive.tzinfo}") - assert naive.tzinfo is None - - -# --------------------------------------------------------------------------- -# Part B — rejected-candidate evidence (APScheduler; NOT acceptance) -# --------------------------------------------------------------------------- - - -@pytest.mark.xfail( - strict=True, - reason="REJECTED CANDIDATE: APScheduler 3.11.0 fabricates a phantom " - "02:30-05:00 for the DST gap. Kept as evidence for the rejection; " - "the chosen library's gap behavior is pinned in Part A. Strict so " - "re-opening the candidacy fails loudly. See implementation plan.", -) -def test_rejected_candidate_gap_phantom(): - # 02:30 does not exist in America/New_York on 2026-03-08. The - # disqualifier is the phantom itself: whatever the library returns - # must not be a nonexistent wall time. - trig = CronTrigger(hour=2, minute=30, second=0, timezone="America/New_York") - nxt = trig.get_next_fire_time(None, utc(2026, 3, 7, 12, 0, 0)) - assert nxt is not None - local = nxt.astimezone(NYC) - print(f"OBSERVED apscheduler dst-gap next local={local.isoformat()}") - assert not ( - (local.year, local.month, local.day) == (2026, 3, 8) and local.hour == 2 - ), f"library must not fabricate a nonexistent wall time, got {local.isoformat()}" - - -@pytest.mark.xfail( - strict=True, - reason="REJECTED CANDIDATE: APScheduler 3.11.0 replays 05:01Z-06:00Z " - "(~59 past minutes) after the fall fold, duplicating UTC occurrence " - "identities. Kept as evidence for the rejection; the chosen " - "library's fold behavior is pinned in Part A. See implementation plan.", -) -def test_rejected_candidate_fold_replay(): - trig = CronTrigger(minute="*", second=0, timezone="America/New_York") - now = utc(2026, 11, 1, 0, 0, 0) - prev = None - seen: set[str] = set() - last: datetime | None = None - for _ in range(450): # spans past 07:00Z: covers BOTH fold hours - nxt = trig.get_next_fire_time(prev, now) - assert nxt is not None - as_utc = nxt.astimezone(UTC) - ident = f"sched-1|{as_utc.isoformat()}" - assert ident not in seen, f"duplicate occurrence identity {ident}" - seen.add(ident) - if last is not None: - assert as_utc > last, "occurrences must be strictly increasing" - last = as_utc - prev, now = nxt, nxt - print(f"OBSERVED fold-span count={len(seen)} unique, last={last.isoformat()}") diff --git a/probes/deployment_scheduling_verify/test_expression_contract_probe.py b/probes/deployment_scheduling_verify/test_expression_contract_probe.py deleted file mode 100644 index 7f41baa9..00000000 --- a/probes/deployment_scheduling_verify/test_expression_contract_probe.py +++ /dev/null @@ -1,138 +0,0 @@ -# DISPOSABLE EXPRESSION-CONTRACT PROBE — NOT PRODUCTION CODE. -# See README.md in this directory. Runs in the repo env: -# uv run pytest -q probes/deployment_scheduling_verify/test_expression_contract_probe.py -"""Pin the CURRENT input-expression contract that the scheduling slice must -reuse (spec: Input authoring and serialization). - -These tests document what exists today for the Phase 1 implementer: the -closed 4-kind union, budget enforcement point, strict-JSON literals, -target-conflict detection, closed GraphSourcePath roots, and the single -hardcoded graph-context evaluator. A schedule occurrence source must plug -into this traversal (T03/T04 of the implementation plan), not copy it. -""" - -from __future__ import annotations - -import pytest -from pydantic import ValidationError - -from wf_core.local_paths import has_overlapping_paths -from wf_core.models.input_bindings import ( - MAX_INPUT_EXPRESSION_DEPTH, - MAX_INPUT_EXPRESSION_NODES, - ArrayExpression, - InputExpressionBinding, - LiteralExpression, - ObjectExpression, - PathExpression, - validate_input_expression_limits, -) -from wf_core.models.json_values import validate_strict_json_value -from wf_core.paths import GraphSourcePath -from wf_core.runtime.input_bindings import ( - resolve_input_expression, - resolve_step_input_bindings, -) - - -def test_expression_union_is_closed_to_four_kinds(): - assert LiteralExpression(kind="literal", value=1).kind == "literal" - assert PathExpression(kind="path", path="input.a").kind == "path" - assert ArrayExpression(kind="array", items=[]).kind == "array" - assert ObjectExpression(kind="object", fields={}).kind == "object" - with pytest.raises(ValidationError): - InputExpressionBinding( - target="x", expression={"kind": "occurrence", "field": "scheduled_at"} - ) # type: ignore[dict-item] - print("OBSERVED occurrence kind rejected: union closed to 4 kinds") - - -def test_budget_constants_and_validator_entry_point(): - assert (MAX_INPUT_EXPRESSION_DEPTH, MAX_INPUT_EXPRESSION_NODES) == (64, 1024) - deep: dict = {"kind": "literal", "value": 0} - for _ in range(MAX_INPUT_EXPRESSION_DEPTH + 5): - deep = {"kind": "array", "items": [deep]} - with pytest.raises(ValueError, match="limit exceeded"): - validate_input_expression_limits(deep) # type: ignore[arg-type] - print("OBSERVED depth budget enforced by validate_input_expression_limits") - - -def test_strict_json_rejects_non_finite_and_non_string_keys(): - with pytest.raises(ValueError): - validate_strict_json_value(float("inf")) - with pytest.raises(ValueError): - validate_strict_json_value({1: "x"}) - assert validate_strict_json_value({"a": [1, None, "x"]}) == {"a": [1, None, "x"]} - print("OBSERVED strict-JSON validator rejects inf and non-string keys") - - -def test_target_conflicts_detected_on_local_paths(): - assert has_overlapping_paths(["a.b", "a.b.c"]) - assert not has_overlapping_paths(["a.b", "a.c"]) - print("OBSERVED overlapping local-path targets detected") - - -def test_graph_source_roots_closed_to_input_state_context(): - assert GraphSourcePath.parse("input.a").root == "input" - assert GraphSourcePath.parse("state.a").root == "state" - assert GraphSourcePath.parse("context.a").root == "context" - with pytest.raises(ValueError): - GraphSourcePath.parse("occurrence.scheduled_at") - print("OBSERVED occurrence root rejected: GraphSourcePath closed") - - -def test_runtime_resolver_composes_literal_object_array_and_paths(): - expr = ObjectExpression( - kind="object", - fields={ - "team": LiteralExpression(kind="literal", value="eng"), - "tags": ArrayExpression( - kind="array", - items=[LiteralExpression(kind="literal", value="a")], - ), - "req": PathExpression(kind="path", path="input.request_id"), - }, - ) - resolved = resolve_input_expression( - expr, - state={}, - workflow_input={"request_id": "r1"}, - context={}, - label="probe", - location="$", - ) - assert resolved == {"team": "eng", "tags": ["a"], "req": "r1"} - print(f"OBSERVED composed resolution -> {resolved}") - - -def test_runtime_resolver_takes_only_graph_context_mappings(): - # There is no source-resolver seam: the only injection point is the - # concrete state/workflow_input/context mappings (faking occurrence - # values through context is exactly what the spec forbids). - import inspect - - sig = inspect.signature(resolve_input_expression) - assert list(sig.parameters) == [ - "expression", - "state", - "workflow_input", - "context", - "label", - "location", - ], f"no resolver parameter exists: {list(sig.parameters)}" - assert "resolver" not in inspect.signature(resolve_step_input_bindings).parameters - print("OBSERVED resolver signatures are concrete graph mappings; no seam") - - -def test_node_input_binding_rejects_expression_kind_at_top_level(): - # StepInputBinding allows expressions, but plain InputBinding - # (deployment-level inputs) does not carry them — resolved data only. - from pydantic import TypeAdapter - - from wf_core.models.input_bindings import InputBinding - - with pytest.raises(ValidationError): - TypeAdapter(InputBinding).validate_python( - {"target": "x", "expression": {"kind": "literal", "value": 1}} - ) - print("OBSERVED top-level InputBinding carries no expressions") diff --git a/probes/deployment_scheduling_verify/test_schedule_state_model.py b/probes/deployment_scheduling_verify/test_schedule_state_model.py deleted file mode 100644 index f69f68ff..00000000 --- a/probes/deployment_scheduling_verify/test_schedule_state_model.py +++ /dev/null @@ -1,1803 +0,0 @@ -# DISPOSABLE SCHEDULING STATE-MODEL PROBE — NOT PRODUCTION CODE. -# See README.md in this directory. Pure stdlib; runs in the repo env: -# uv run pytest -q probes/deployment_scheduling_verify/test_schedule_state_model.py -"""Executable reference model of the deployment-scheduling state machine. - -Covers the spec's state gates with injected clocks and controlled -execution (no sleeps, no threads): -overlap=skip|parallel x misfire=skip|latest, latest-means-ONE-candidate, -supersession, terminal skips, slot accounting, pause-vs-downtime, -edit/delete/restart, capacity/fairness, fault boundaries, crash-during- -resume, exclusive ownership. - -The calendar is an abstract due-instant source here (next_after / -prev_before only — deliberately NO iter_between, to forbid unbounded -enumeration). Calendar math itself is probed in test_calendar_probe.py. -""" - -from __future__ import annotations - -from dataclasses import dataclass, field -from datetime import datetime, timedelta, timezone - -import pytest - -UTC = timezone.utc - - -def ts(y, mo, d, h=0, mi=0, s=0) -> datetime: - return datetime(y, mo, d, h, mi, s, tzinfo=UTC) - - -# -------------------------------------------------------------------------- -# Errors -# -------------------------------------------------------------------------- - - -class InjectedFault(Exception): - pass - - -class SecondOwnerError(Exception): - pass - - -class StartupRejected(Exception): - pass - - -class BlockedSchedule(Exception): - pass - - -class ExecutorCrashed(Exception): - pass - - -# -------------------------------------------------------------------------- -# Records -# -------------------------------------------------------------------------- - -RUNNING = "running" -INTERRUPTED = "interrupted" # durably waiting, resumable -COMPLETED = "completed" -FAILED = "failed" -ACTIVE_STATES = (RUNNING, INTERRUPTED) - - -@dataclass -class Schedule: - id: str - rev: int = 1 - enabled: bool = True - paused: bool = False - deleted: bool = False - overlap: str = "skip" # skip | parallel - misfire: str = "skip" # skip | latest - max_active: int = 1 - allowance_s: float = 60.0 - deployment_id: str = "dep-1" - created_dep_rev: int = 1 - exhausted: bool = False - blocked_reason: str | None = None - - -@dataclass -class Candidate: - sched_id: str - intended: datetime - rev: int - - -@dataclass -class Run: - id: str - sched_id: str - intended: datetime - rev: int - frozen_input: dict - state: str = RUNNING - dispatched_unknown: bool = False # dispatched, no stopped result yet - needs_dispatch: bool = False # admitted + view exists, never dispatched - attempt_id: int = 0 # resume attempt the run currently belongs to - result_attempt: int | None = None # attempt that produced the stopped result - fail_reason: str = "" - - -@dataclass -class Record: - kind: str # admitted|completed|failed|skipped-overlap|skipped-misfire| - # superseded|preflight-rejected|exhausted|interval-summary|interrupted - sched_id: str - intended: datetime | None = None - run_id: str | None = None - reason: str = "" - interval: tuple[datetime, datetime] | None = None - count: int = 0 - - -# -------------------------------------------------------------------------- -# Occurrence sources: next_after / prev_before ONLY (no enumeration seam) -# -------------------------------------------------------------------------- - - -class CountingMixin: - def __init__(self) -> None: - self.next_calls = 0 - self.prev_calls = 0 - - -class PeriodicSource(CountingMixin): - def __init__(self, period: timedelta, start: datetime) -> None: - super().__init__() - self.period = period - self.start = start - - def next_after(self, instant: datetime) -> datetime | None: - self.next_calls += 1 - if instant < self.start: - return self.start - n = (instant - self.start) // self.period + 1 - return self.start + n * self.period - - def prev_before(self, instant: datetime) -> datetime | None: - self.prev_calls += 1 - if instant <= self.start: - return None - n = (instant - self.start - timedelta(microseconds=1)) // self.period - return self.start + n * self.period - - -class OneShotSource(CountingMixin): - def __init__(self, at: datetime) -> None: - super().__init__() - self.at = at - - def next_after(self, instant: datetime) -> datetime | None: - self.next_calls += 1 - return self.at if instant < self.at else None - - def prev_before(self, instant: datetime) -> datetime | None: - self.prev_calls += 1 - return self.at if instant > self.at else None - - -# -------------------------------------------------------------------------- -# Store with fault injection -# -------------------------------------------------------------------------- - - -@dataclass -class MemStore: - lockable: bool = True - lock_holder: str | None = None - id_seq: int = 0 # store-backed run counter: survives scheduler restart - attempt_seq: int = 0 # store-backed resume-attempt counter: same reason - schedules: dict[str, Schedule] = field(default_factory=dict) - candidates: dict[str, Candidate | None] = field(default_factory=dict) - consumed: dict[str, datetime] = field(default_factory=dict) - admissions: dict[str, dict] = field(default_factory=dict) # run_id -> record - runs: dict[str, Run] = field(default_factory=dict) - resume_attempts: dict[str, str] = field(default_factory=dict) # run -> ACTIVE|DONE - history: list[Record] = field(default_factory=list) - deployments: dict[str, dict] = field(default_factory=dict) # id -> {rev, required} - faults: dict[str, str] = field(default_factory=dict) # op -> before|after - op_log: list[str] = field(default_factory=list) - poll_cursor: int = 0 - - def check_fault(self, op: str) -> None: - mode = self.faults.pop(op, None) - self.op_log.append(op) - if mode == "before": - raise InjectedFault(f"{op}:before") - if mode == "after": - self.op_log.append(f"{op}:after-pending") - - def after_ok(self, op: str) -> None: - # A test driver calls this after the op's effect to confirm the - # injected 'after' fault (crash between effect and next step). - if f"{op}:after-pending" in self.op_log: - self.op_log.remove(f"{op}:after-pending") - raise InjectedFault(f"{op}:after") - - -EPOCH = ts(2020, 1, 1) - - -# -------------------------------------------------------------------------- -# Scheduler reference model -# -------------------------------------------------------------------------- - -SCAN_CAP = 100 # max next_after calls per schedule per poll before jumping - - -class Scheduler: - def __init__( - self, - store: MemStore, - owner: str, - capacity: int, - sources: dict[str, CountingMixin], - outcomes: dict[str, str] | None = None, - ) -> None: - if not store.lockable: - raise StartupRejected("unsupported locking rejects scheduler startup") - if store.lock_holder is not None and store.lock_holder != owner: - raise SecondOwnerError(f"store owned by {store.lock_holder}") - store.lock_holder = owner - self.store = store - self.owner = owner - self.capacity = capacity - self.sources = sources - self.outcomes = outcomes or {} - - def close(self) -> None: - if self.store.lock_holder == self.owner: - self.store.lock_holder = None - - def __enter__(self) -> Scheduler: - return self - - def __exit__(self, *exc: object) -> None: - self.close() - - # -- helpers --------------------------------------------------------- - def _active(self, sched_id: str) -> list[Run]: - return [ - r - for r in self.store.runs.values() - if r.sched_id == sched_id and r.state in ACTIVE_STATES - ] - - def _task_load(self) -> int: - # Admitted-but-never-dispatched runs hold a schedule slot but no - # executing-task slot: nothing is running on their behalf. - return sum( - 1 - for r in self.store.runs.values() - if r.state == RUNNING and not r.needs_dispatch - ) - - def _record(self, **kw: object) -> None: - self.store.history.append(Record(**kw)) # type: ignore[arg-type] - - def _admit(self, sched: Schedule, intended: datetime, now: datetime) -> str | None: - """Ordered admission. Returns run_id, 'held', or None (terminal).""" - st = self.store - if sched.blocked_reason: - raise BlockedSchedule(sched.blocked_reason) - if not sched.enabled or sched.deleted or sched.paused: - return None - # 0. recheck + preflight against CURRENT deployment contract - dep = st.deployments.get(sched.deployment_id) - if dep is None: - self._record( - kind="preflight-rejected", - sched_id=sched.id, - intended=intended, - reason="deployment-deleted", - ) - st.consumed[sched.id] = intended - return None - frozen = { - "team": "eng", - "report_time": intended.isoformat(), - "sched": sched.id, - "dep_rev": dep["rev"], - } - missing = [k for k in dep["required"] if k not in frozen] - if missing: - self._record( - kind="preflight-rejected", - sched_id=sched.id, - intended=intended, - reason=f"missing-input:{missing}", - ) - if ( - st.candidates.get(sched.id) is not None - and st.candidates[sched.id].intended == intended - ): # type: ignore[union-attr] - st.candidates[sched.id] = None - st.consumed[sched.id] = intended - return None - # 1. overlap first (takes precedence over capacity) - active = self._active(sched.id) - if sched.overlap == "skip" and active: - self._record( - kind="skipped-overlap", - sched_id=sched.id, - intended=intended, - reason=f"active={[r.id for r in active]}", - ) - if ( - st.candidates.get(sched.id) is not None - and st.candidates[sched.id].intended == intended - ): # type: ignore[union-attr] - st.candidates[sched.id] = None - st.consumed[sched.id] = intended - return None - if sched.overlap == "parallel" and len(active) >= sched.max_active: - self._record( - kind="skipped-overlap", - sched_id=sched.id, - intended=intended, - reason=f"max_active={sched.max_active}", - ) - if ( - st.candidates.get(sched.id) is not None - and st.candidates[sched.id].intended == intended - ): # type: ignore[union-attr] - st.candidates[sched.id] = None - st.consumed[sched.id] = intended - return None - # 2. capacity: skip expires at deadline, latest holds one candidate - if self._task_load() >= self.capacity: - if sched.misfire == "latest": - st.candidates[sched.id] = Candidate(sched.id, intended, sched.rev) - st.consumed[sched.id] = intended - return "held" - if (now - intended).total_seconds() > sched.allowance_s: - self._record( - kind="skipped-misfire", - sched_id=sched.id, - intended=intended, - reason="capacity-deadline", - ) - st.consumed[sched.id] = intended - return None - return "held-undecided" # retry next poll, consumed NOT advanced - # 3. allocate + freeze (store-backed counter: no reuse on restart) - st.id_seq += 1 - run_id = f"run-{sched.id}-{st.id_seq}" - # 4. persist admission record BEFORE dispatch (fault boundary). - # The occurrence is decided here: history entry, candidate - # clearing, consumed progress, and one-shot exhaustion all belong - # to this persist, so a later crash can never re-decide it. - st.check_fault("admission") - st.admissions[run_id] = { - "sched": sched.id, - "intended": intended, - "rev": sched.rev, - "input": dict(frozen), - } - # The occurrence-status entry is part of the admission persist itself. - self._record( - kind="admitted", - sched_id=sched.id, - intended=intended, - run_id=run_id, - reason=f"rev={sched.rev}", - ) - # One-shot admission exhausts the schedule (admit exactly once). - if isinstance(self.sources.get(sched.id), OneShotSource): - sched.exhausted = True - if ( - st.candidates.get(sched.id) is not None - and st.candidates[sched.id].intended == intended - ): # type: ignore[union-attr] - st.candidates[sched.id] = None - st.consumed[sched.id] = intended - st.after_ok("admission") - # 5. materialize run view with same identity (fault boundary) - st.check_fault("materialize") - st.runs[run_id] = Run(run_id, sched.id, intended, sched.rev, dict(frozen)) - st.after_ok("materialize") - # 6. dispatch the captured invocation without re-resolving (faults) - st.check_fault("dispatch") - try: - self._dispatch(run_id, now) - finally: - st.after_ok("dispatch") - return run_id - - def _dispatch(self, run_id: str, now: datetime) -> None: - st = self.store - run = st.runs[run_id] - outcome = self.outcomes.get(run_id, self.outcomes.get("*", "complete")) - if outcome == "hang": - run.dispatched_unknown = True # task alive in-process; restart abandons it - return # stays RUNNING, occupies schedule + task slots - if outcome == "crash": - run.dispatched_unknown = True - raise ExecutorCrashed(run_id) - if outcome == "interrupt": - st.check_fault("interrupt-persist") - run.state = INTERRUPTED - st.after_ok("interrupt-persist") - self._record( - kind="interrupted", - sched_id=run.sched_id, - intended=run.intended, - run_id=run_id, - ) - return - assert outcome == "complete" - st.check_fault("complete") - run.state = COMPLETED - st.after_ok("complete") - self._record( - kind="completed", - sched_id=run.sched_id, - intended=run.intended, - run_id=run_id, - ) - - # -- polling ---------------------------------------------------------- - def poll(self, now: datetime) -> dict[str, str]: - st = self.store - if st.lock_holder != self.owner: - raise SecondOwnerError("lost ownership") - self._dispatch_pending(now) - ids = sorted(s.id for s in st.schedules.values()) - if not ids: - return {} - start = st.poll_cursor % len(ids) - order = ids[start:] + ids[:start] - st.poll_cursor += 1 - results: dict[str, str] = {} - for sid in order: - results[sid] = self._poll_one(st.schedules[sid], now) - return results - - def _dispatch_pending(self, now: datetime) -> None: - """Dispatch runs that recovery materialized but never executed. - - Recovery NEVER executes: it only completes missing views and flags - them pending. Execution happens here, through the same capacity - checks as normal admission — never inside recovery. Pending runs - of blocked schedules stay pending (fail closed). - """ - st = self.store - for run in sorted(st.runs.values(), key=lambda r: r.id): - if not run.needs_dispatch: - continue - sched = st.schedules.get(run.sched_id) - if sched is None or sched.blocked_reason: - continue - if self._task_load() >= self.capacity: - continue # stays pending until a slot frees - run.needs_dispatch = False - try: - self._dispatch(run.id, now) - except ExecutorCrashed: - # Dispatched with unknown outcome: next recovery fails it - # closed as abandoned (no replay), like any dispatch crash. - run.dispatched_unknown = True - - def _poll_one(self, sched: Schedule, now: datetime) -> str: - st = self.store - if sched.blocked_reason: - raise BlockedSchedule(sched.blocked_reason) - if sched.deleted: - # Deletion clears pending candidates; admitted runs/history stay. - if st.candidates.get(sched.id) is not None: - st.candidates[sched.id] = None - return "deleted" - if not sched.enabled: - # Resolved edge (no spec disable concept): administrative disable - # behaves like pause for catch-up — excluded, never backfilled. - if st.candidates.get(sched.id) is not None: - st.candidates[sched.id] = None - st.consumed[sched.id] = max(st.consumed.get(sched.id, EPOCH), now) - return "disabled" - if sched.exhausted: - return "exhausted" - if sched.paused: - # Explicit pause: clear unadmitted candidates, exclude interval. - if st.candidates.get(sched.id) is not None: - st.candidates[sched.id] = None - st.consumed[sched.id] = max(st.consumed.get(sched.id, EPOCH), now) - return "paused" - src = self.sources[sched.id] - consumed = st.consumed.get(sched.id, EPOCH) - if consumed > now: - return "clock-rollback-held" # never re-admit consumed instants - # Bounded scan: at most SCAN_CAP next_after calls, then jump. - due: list[datetime] = [] - cursor = consumed - jumped = False - while True: - nxt = src.next_after(cursor) # type: ignore[attr-defined] - if nxt is None or nxt > now: - if ( - not due - and isinstance(src, OneShotSource) - and nxt is None - and not sched.exhausted - and sched.misfire == "skip" - and self._is_oneshot_expired(src, now) - ): - self._record( - kind="exhausted", - sched_id=sched.id, - intended=src.at, - reason="oneshot-expired-skip", - ) - sched.exhausted = True - st.consumed[sched.id] = now - st.candidates[sched.id] = None - return "exhausted" - break - due.append(nxt) - cursor = nxt - if len(due) >= SCAN_CAP: - jumped = True - break - if jumped: - # Never enumerate further: one bounded query + interval summary. - if sched.misfire == "latest": - latest = src.prev_before(now) # type: ignore[attr-defined] - assert latest is not None - old = st.candidates.get(sched.id) - if old is not None and old.intended != latest: - self._record( - kind="superseded", - sched_id=sched.id, - intended=old.intended, - reason=f"coalesced-into:{latest.isoformat()}", - ) - st.candidates[sched.id] = Candidate(sched.id, latest, sched.rev) - self._record( - kind="interval-summary", - sched_id=sched.id, - interval=(consumed, now), - count=-1, - reason="coalesced-missed-span", - ) - st.consumed[sched.id] = now - return self._admit_held_candidate(sched, now) or "candidate-held" - self._record( - kind="interval-summary", - sched_id=sched.id, - interval=(consumed, now), - count=-1, - reason="skipped-missed-span", - ) - st.consumed[sched.id] = now - return "span-skipped" - # Normal path: decide each due instant in order. - last_result = "idle" - for instant in due: - if instant <= st.consumed.get(sched.id, EPOCH): - continue - age = (now - instant).total_seconds() - if age <= sched.allowance_s: - # A timely admission consumes any older held candidate: the - # newer due instant supersedes it (recorded, never admitted). - old = st.candidates.get(sched.id) - if old is not None and old.intended != instant: - self._record( - kind="superseded", - sched_id=sched.id, - intended=old.intended, - reason=f"admitted-newer:{instant.isoformat()}", - ) - st.candidates[sched.id] = None - r = self._admit(sched, instant, now) - last_result = f"admit:{r}" - elif sched.misfire == "latest": - old = st.candidates.get(sched.id) - if old is not None and old.intended != instant: - self._record( - kind="superseded", - sched_id=sched.id, - intended=old.intended, - reason=f"coalesced-into:{instant.isoformat()}", - ) - # A newer due time supersedes the unadmitted candidate; an - # admitted run is never touched (candidates only). - st.candidates[sched.id] = Candidate(sched.id, instant, sched.rev) - st.consumed[sched.id] = instant - last_result = self._admit_held_candidate(sched, now) or "candidate-held" - else: - if isinstance(src, OneShotSource) and not sched.exhausted: - self._record( - kind="exhausted", - sched_id=sched.id, - intended=instant, - reason="oneshot-expired-skip", - ) - sched.exhausted = True - st.consumed[sched.id] = instant - last_result = "exhausted" - continue - self._record( - kind="skipped-misfire", - sched_id=sched.id, - intended=instant, - reason=f"age={age:.0f}s", - ) - st.consumed[sched.id] = instant - last_result = "skipped-misfire" - # Also try a held candidate whose slot may have freed. - if last_result in ("idle",) and st.candidates.get(sched.id) is not None: - last_result = self._admit_held_candidate(sched, now) or "candidate-held" - return last_result - - def _is_oneshot_expired(self, src: OneShotSource, now: datetime) -> bool: - return src.at <= now - - def _admit_held_candidate(self, sched: Schedule, now: datetime) -> str | None: - st = self.store - cand = st.candidates.get(sched.id) - if cand is None or cand.rev != sched.rev: - if cand is not None and cand.rev != sched.rev: - self._record( - kind="superseded", - sched_id=sched.id, - intended=cand.intended, - reason="schedule-edit", - ) - st.candidates[sched.id] = None - return None - r = self._admit(sched, cand.intended, now) - return r if r != "held-undecided" else None - - # -- administration ---------------------------------------------------- - def resume_schedule(self, sid: str, now: datetime) -> None: - """Unpause: resume selects the next future occurrence; the paused - interval is excluded from catch-up under both policies.""" - sched = self.store.schedules[sid] - sched.paused = False - self.store.consumed[sid] = max(self.store.consumed.get(sid, EPOCH), now) - self.store.candidates[sid] = None - - def edit_schedule(self, sid: str, now: datetime) -> None: - """Definition edit: new revision, discard old candidates, begin at - the edit time (creation/edits never backfill).""" - st = self.store - sched = st.schedules[sid] - sched.rev += 1 - old = st.candidates.get(sid) - if old is not None: - self._record( - kind="superseded", - sched_id=sid, - intended=old.intended, - reason="schedule-edit", - ) - st.candidates[sid] = None - st.consumed[sid] = max(st.consumed.get(sid, EPOCH), now) - - # -- resume ------------------------------------------------------------ - def resume_run(self, run_id: str, now: datetime) -> str: - st = self.store - run = st.runs[run_id] - assert run.state == INTERRUPTED, "only waiting interruptions resume" - if st.resume_attempts.get(run_id) == "ACTIVE": - raise BlockedSchedule("ambiguous attempt already active") - if self._task_load() >= self.capacity: - return "blocked-capacity" # stays waiting; slot needed - # Durably mark the active attempt BEFORE executing again. The mark - # carries a fresh attempt identity that later stopped results echo - # back, so recovery can tell a new result from the old checkpoint. - st.check_fault("resume-mark") - st.attempt_seq += 1 - run.attempt_id = st.attempt_seq - st.resume_attempts[run_id] = "ACTIVE" - st.after_ok("resume-mark") - run.state = RUNNING - run.dispatched_unknown = True - st.check_fault("resume-execute") - try: - outcome = self.outcomes.get(run_id, self.outcomes.get("*", "complete")) - if outcome == "crash": - raise ExecutorCrashed(run_id) - if outcome == "interrupt": - return self._finish_resume_interrupted(run) - return self._finish_resume_completed(run) - finally: - st.after_ok("resume-execute") - - def _finish_resume_completed(self, run: Run) -> str: - # Completion, attempt-clearing, and history recording are separate - # persists with a fault boundary between each pair; a resumed run - # may also interrupt again instead of completing. - st = self.store - st.check_fault("resume-complete") - run.state = COMPLETED - run.dispatched_unknown = False - run.result_attempt = run.attempt_id - st.after_ok("resume-complete") - st.check_fault("resume-attempt-clear") - st.resume_attempts[run.id] = "DONE" - st.after_ok("resume-attempt-clear") - self._record( - kind="completed", - sched_id=run.sched_id, - intended=run.intended, - run_id=run.id, - reason="resumed-complete", - ) - return "resumed-complete" - - def _finish_resume_interrupted(self, run: Run) -> str: - st = self.store - st.check_fault("resume-interrupt") - run.state = INTERRUPTED - run.dispatched_unknown = False - run.result_attempt = run.attempt_id - st.after_ok("resume-interrupt") - st.check_fault("resume-attempt-clear") - st.resume_attempts[run.id] = "DONE" - st.after_ok("resume-attempt-clear") - self._record( - kind="interrupted", - sched_id=run.sched_id, - intended=run.intended, - run_id=run.id, - reason="resumed-reinterrupted", - ) - return "resumed-interrupted" - - # -- recovery ------------------------------------------------------------ - def recover(self, now: datetime) -> list[str]: - """Startup recovery under exclusive ownership. Returns diagnostics.""" - st = self.store - if st.lock_holder != self.owner: - raise SecondOwnerError("lost ownership") - diags: list[str] = [] - # Admission record is the recovery authority. Recovery NEVER - # executes: it only materializes missing views (flagged pending for - # the capacity-checked dispatch sweep in poll()), fails abandoned / - # ambiguous runs closed, reconciles terminal records, and blocks - # corrupt schedules. Any RUNNING run lost its in-memory task with - # the old process: with an admission record it is abandoned - # (failed, no replay); without one it is corrupt (fail closed). - # Interrupted runs carry the identity of the attempt that produced - # them: a result matching the ACTIVE attempt is fresh (resumable); - # anything else under an ACTIVE attempt is stale (ambiguous, failed). - for run_id, rec in st.admissions.items(): - if run_id not in st.runs: - # Admitted but never materialized/dispatched: complete the - # missing view and flag it pending. The next poll dispatches - # the captured invocation through capacity checks — exactly - # once, because the occurrence is already consumed. - st.runs[run_id] = Run( - run_id, - rec["sched"], - rec["intended"], - rec["rev"], - dict(rec["input"]), - ) - st.runs[run_id].needs_dispatch = True - diags.append(f"{run_id}:view-completed-pending-dispatch") - for run in st.runs.values(): - if run.state == INTERRUPTED: - if st.resume_attempts.get(run.id) == "ACTIVE": - if ( - run.result_attempt is not None - and run.result_attempt == run.attempt_id - ): - # The stopped result belongs to the active attempt: - # the re-interruption landed before the crash. Fresh, - # resumable; the attempt is done executing. - st.resume_attempts[run.id] = "DONE" - diags.append(f"{run.id}:fresh-result-resumable") - else: - run.state = FAILED - run.fail_reason = ( - "ambiguous resume attempt: may have executed; " - "external effects may already have occurred; no retry" - ) - self._record( - kind="failed", - sched_id=run.sched_id, - intended=run.intended, - run_id=run.id, - reason=run.fail_reason, - ) - diags.append(f"{run.id}:failed-closed") - else: - diags.append(f"{run.id}:waiting-resumable") - elif run.state == RUNNING: - if run.needs_dispatch: - # Admitted and materialized but provably never - # dispatched: stays pending for the poll sweep, never - # failed as abandoned. - diags.append(f"{run.id}:pending-dispatch-kept") - continue - if run.id not in st.admissions: - sched = st.schedules[run.sched_id] - sched.blocked_reason = ( - f"corrupt run view without admission: {run.id}" - ) - diags.append(f"{run.id}:corrupt-blocked") - continue - if st.resume_attempts.get(run.id) == "ACTIVE": - run.state = FAILED - run.fail_reason = ( - "ambiguous resume attempt: may have executed; " - "external effects may already have occurred; no retry" - ) - else: - run.state = FAILED - run.fail_reason = ( - "abandoned execution: outcome unknown; external " - "effects may already have occurred; no replay" - ) - self._record( - kind="failed", - sched_id=run.sched_id, - intended=run.intended, - run_id=run.id, - reason=run.fail_reason, - ) - diags.append(f"{run.id}:failed-closed") - # Reconcile terminal runs whose history entry was lost to a crash - # between the state persist and the record append (idempotent: only - # appends when no terminal record exists for the run). A COMPLETED - # run with a still-ACTIVE attempt crashed between completion and - # attempt-clearing: mark DONE, reconcile the record, never re-run. - for run in st.runs.values(): - if run.state == COMPLETED and st.resume_attempts.get(run.id) == "ACTIVE": - st.resume_attempts[run.id] = "DONE" - diags.append(f"{run.id}:attempt-reconciled") - if run.state == COMPLETED and not any( - r.kind == "completed" and r.run_id == run.id for r in st.history - ): - self._record( - kind="completed", - sched_id=run.sched_id, - intended=run.intended, - run_id=run.id, - reason="reconciled-on-recovery", - ) - diags.append(f"{run.id}:terminal-reconciled") - elif ( - run.state == INTERRUPTED - and st.resume_attempts.get(run.id) != "ACTIVE" - and not any( - r.kind == "interrupted" and r.run_id == run.id for r in st.history - ) - ): - self._record( - kind="interrupted", - sched_id=run.sched_id, - intended=run.intended, - run_id=run.id, - reason="reconciled-on-recovery", - ) - diags.append(f"{run.id}:terminal-reconciled") - return diags - - -def make_store() -> MemStore: - st = MemStore() - st.deployments["dep-1"] = {"rev": 1, "required": ["team", "report_time"]} - return st - - -def add_sched( - st: MemStore, - sid: str, - src: CountingMixin, - sources: dict, - start: datetime, - **kw: object, -) -> Schedule: - if sid in st.schedules: - raise ValueError(f"schedule id reused: {sid}") - sched = Schedule(id=sid, **kw) # type: ignore[arg-type] - st.schedules[sid] = sched - st.candidates[sid] = None - st.consumed[sid] = start - sources[sid] = src - return sched - - -# -------------------------------------------------------------------------- -# Tests -# -------------------------------------------------------------------------- - - -def test_defaults_skip_skip(): - s = Schedule(id="s") - assert s.overlap == "skip" and s.misfire == "skip" - - -def test_overlap_skip_x_misfire_skip_running_blocks_and_late_drops(): - st = make_store() - sources: dict = {} - t0 = ts(2026, 9, 8, 12, 0) - add_sched( - st, - "a", - PeriodicSource(timedelta(minutes=10), t0), - sources, - t0 - timedelta(minutes=10), - ) - with Scheduler(st, "p1", 4, sources, {"*": "hang"}) as sch: - sch.poll(t0) # admits 12:00, hangs (RUNNING) - (run_id,) = list(st.runs) - sch.poll(t0 + timedelta(minutes=10)) # 12:10 due while active - kinds = [(r.kind, r.intended) for r in st.history if r.sched_id == "a"] - assert ("skipped-overlap", ts(2026, 9, 8, 12, 10)) in kinds - # Late instant beyond allowance with skip: skipped-misfire. - sch.poll(t0 + timedelta(minutes=30)) # 12:20,12:30 missed (>60s) - kinds = [(r.kind, r.intended) for r in st.history if r.sched_id == "a"] - assert ("skipped-misfire", ts(2026, 9, 8, 12, 20)) in kinds - assert st.runs[run_id].state == RUNNING - - -def test_overlap_skip_x_misfire_latest_single_catchup_and_supersession(): - st = make_store() - sources: dict = {} - add_sched( - st, - "h", - PeriodicSource(timedelta(hours=1), ts(2026, 9, 8, 9, 0)), - sources, - ts(2026, 9, 8, 9, 0), - misfire="latest", - ) - with Scheduler(st, "p1", 0, sources) as sch: # no capacity: hold candidates - sch.poll(ts(2026, 9, 8, 12, 20)) # missed 10:00,11:00,12:00 - cand = st.candidates["h"] - assert cand is not None and cand.intended == ts(2026, 9, 8, 12, 0) - runs_for_h = [r for r in st.history if r.kind == "admitted"] - assert runs_for_h == [], "ONE pending candidate, not replay-all" - sup = [r for r in st.history if r.kind == "superseded"] - assert {r.intended for r in sup} == { - ts(2026, 9, 8, 10, 0), - ts(2026, 9, 8, 11, 0), - } - # Capacity returns at 13:00 while 13:00 is also due: exactly one admission, - # for the 13:00 instant (newer supersedes the unadmitted 12:00). - with Scheduler(st, "p1", 4, sources, {"*": "complete"}) as sch2: - sch2.poll(ts(2026, 9, 8, 13, 0)) - admitted = [r for r in st.history if r.kind == "admitted"] - assert len(admitted) == 1 - assert admitted[0].intended == ts(2026, 9, 8, 13, 0) - assert st.candidates["h"] is None, "no stale candidate survives admission" - assert any( - r.kind == "superseded" and r.intended == ts(2026, 9, 8, 12, 0) - for r in st.history - ) - assert admitted[0].run_id is not None - assert ( - st.runs[admitted[0].run_id].frozen_input["report_time"] - == "2026-09-08T13:00:00+00:00" - ) - - -def test_newer_due_never_supersedes_admitted_run(): - st = make_store() - sources: dict = {} - t0 = ts(2026, 9, 8, 12, 0) - add_sched( - st, - "a", - PeriodicSource(timedelta(minutes=10), t0), - sources, - t0 - timedelta(minutes=10), - misfire="latest", - ) - with Scheduler(st, "p1", 4, sources, {"*": "hang"}) as sch: - first = sch.poll(t0)["a"] - run_id = first.split(":", 1)[1] - assert st.runs[run_id].intended == t0 - sch.poll(t0 + timedelta(minutes=10)) - # 12:10 is skipped-overlap (terminal); the admitted 12:00 run is - # untouched and no candidate resurrects 12:10. - assert st.runs[run_id].state == RUNNING - assert st.candidates["a"] is None - assert any( - r.kind == "skipped-overlap" and r.intended == t0 + timedelta(minutes=10) - for r in st.history - ) - - -def test_terminal_overlap_skip_never_reappears_through_catchup(): - st = make_store() - sources: dict = {} - t0 = ts(2026, 9, 8, 12, 0) - add_sched( - st, - "a", - PeriodicSource(timedelta(minutes=10), t0), - sources, - t0 - timedelta(minutes=10), - misfire="latest", - ) - with Scheduler(st, "p1", 4, sources, {"*": "hang"}) as sch: - sch.poll(t0) - sch.poll(t0 + timedelta(minutes=10)) # terminal skipped-overlap @12:10 - # Restart abandons the hanging in-flight task (failed, no replay), but - # the terminal 12:10 skip is never reconstructed as a candidate. - with Scheduler(st, "p1", 0, sources) as sch2: - diags = sch2.recover(t0 + timedelta(minutes=11)) - assert any("failed-closed" in d for d in diags) - sch2.poll(t0 + timedelta(minutes=25)) # 12:20 missed -> held candidate - cand = st.candidates["a"] - assert cand is not None and cand.intended == t0 + timedelta(minutes=20) - intents = [ - r.intended - for r in st.history - if r.kind == "admitted" and r.intended == t0 + timedelta(minutes=10) - ] - assert intents == [], "terminally skipped 12:10 must never be admitted" - - -def test_parallel_limits_count_interrupted_and_waiting_frees_task_slot(): - st = make_store() - sources: dict = {} - t0 = ts(2026, 9, 8, 12, 0) - add_sched( - st, - "p", - PeriodicSource(timedelta(minutes=5), t0), - sources, - t0 - timedelta(minutes=5), - overlap="parallel", - max_active=2, - ) - with Scheduler(st, "owner", 1, sources, {"*": "hang"}) as sch: - sch.poll(t0) # run1 RUNNING (task slot taken) - assert sch._task_load() == 1 - st.runs["run-p-1"].state = INTERRUPTED # scripted durable interrupt - st.history.append( - Record(kind="interrupted", sched_id="p", intended=t0, run_id="run-p-1") - ) - assert sch._task_load() == 0, "waiting interruptions hold no task slot" - sch.poll(t0 + timedelta(minutes=5)) # run2 admitted (1 task slot free) - assert st.runs["run-p-2"].state == RUNNING - # max_active=2 reached (1 waiting + 1 running): 12:10 skipped-overlap. - sch.poll(t0 + timedelta(minutes=10)) - assert any( - r.kind == "skipped-overlap" and r.intended == t0 + timedelta(minutes=10) - for r in st.history - ) - # Lowering the limit never cancels; blocks new admission instead. - st.schedules["p"].max_active = 1 - st.runs["run-p-2"].state = INTERRUPTED - sch.poll(t0 + timedelta(minutes=15)) - assert any( - r.kind == "skipped-overlap" and r.intended == t0 + timedelta(minutes=15) - for r in st.history - ) - assert st.runs["run-p-1"].state == INTERRUPTED # not cancelled - # Resume needs a task slot: none free while... free one by completing. - st.runs["run-p-2"].state = COMPLETED - assert ( - sch.resume_run("run-p-1", t0 + timedelta(minutes=16)) == "resumed-complete" - ) - - -def test_explicit_pause_is_not_downtime(): - st = make_store() - sources: dict = {} - add_sched( - st, - "a", - PeriodicSource(timedelta(hours=1), ts(2026, 9, 8, 9, 0)), - sources, - ts(2026, 9, 8, 9, 0), - misfire="latest", - ) - with Scheduler(st, "p1", 4, sources, {"*": "complete"}) as sch: - sch.poll(ts(2026, 9, 8, 9, 0)) - st.schedules["a"].paused = True - sch.poll(ts(2026, 9, 8, 10, 30)) # paused polls exclude the interval - assert st.candidates["a"] is None - sch.poll(ts(2026, 9, 8, 11, 30)) - sch.resume_schedule("a", ts(2026, 9, 8, 12, 30)) # next future only - sch.poll(ts(2026, 9, 8, 12, 30)) - assert st.candidates["a"] is None, "paused times are not replayed" - admitted_intents = [r.intended for r in st.history if r.kind == "admitted"] - assert ts(2026, 9, 8, 10, 0) not in admitted_intents - assert ts(2026, 9, 8, 11, 0) not in admitted_intents - assert ts(2026, 9, 8, 12, 0) not in admitted_intents - # Contrast: enabled downtime DOES catch up under latest. - st2 = make_store() - sources2: dict = {} - add_sched( - st2, - "b", - PeriodicSource(timedelta(hours=1), ts(2026, 9, 8, 9, 0)), - sources2, - ts(2026, 9, 8, 9, 0), - misfire="latest", - ) - sch2 = Scheduler(st2, "p1", 0, sources2) - sch2.poll(ts(2026, 9, 8, 12, 20)) # 3 missed while "down", enabled - assert st2.candidates["b"] is not None - assert st2.candidates["b"].intended == ts(2026, 9, 8, 12, 0) # type: ignore[union-attr] - sch2.close() - - -def test_edit_discards_candidates_no_backfill_delete_keeps_history(): - st = make_store() - sources: dict = {} - t0 = ts(2026, 9, 8, 12, 0) - add_sched( - st, - "a", - PeriodicSource(timedelta(hours=1), ts(2026, 9, 8, 9, 0)), - sources, - ts(2026, 9, 8, 9, 0), - misfire="latest", - ) - with Scheduler(st, "p1", 0, sources) as sch: - sch.poll(ts(2026, 9, 8, 12, 20)) - assert st.candidates["a"] is not None - sch.edit_schedule("a", ts(2026, 9, 8, 12, 21)) # definition edit - sch.poll(ts(2026, 9, 8, 12, 21)) - assert st.candidates["a"] is None, "edits discard old-revision candidates" - assert any( - r.kind == "superseded" and r.reason == "schedule-edit" for r in st.history - ) - # No backfill before the edit: consumed advanced to edit time. - assert st.consumed["a"] >= ts(2026, 9, 8, 12, 21) - with Scheduler(st, "p1", 4, sources, {"*": "hang"}) as sch: - sch.poll(ts(2026, 9, 8, 13, 0)) - (run_id,) = [r for r in st.runs if st.runs[r].state == RUNNING] - st.schedules["a"].deleted = True # delete clears pending, keeps rest - sch.poll(ts(2026, 9, 8, 14, 0)) - assert st.candidates["a"] is None - assert st.runs[run_id].state == RUNNING, "delete must not cancel active runs" - n_history = len(st.history) - assert n_history > 0, "history remains intact" - # In-flight work can still finish after delete; history is appended. - st.runs[run_id].state = COMPLETED - st.history.append( - Record( - kind="completed", - sched_id="a", - intended=st.runs[run_id].intended, - run_id=run_id, - ) - ) - assert len(st.history) == n_history + 1 - with pytest.raises(ValueError, match="schedule id reused"): - add_sched(st, "a", PeriodicSource(timedelta(hours=1), t0), sources, t0) - - -def test_restart_survives_candidate_and_rollback_never_readmits(): - st = make_store() - sources: dict = {} - t0 = ts(2026, 9, 8, 12, 0) - add_sched( - st, - "a", - PeriodicSource(timedelta(minutes=10), t0), - sources, - t0 - timedelta(minutes=10), - misfire="latest", - ) - with Scheduler(st, "p1", 0, sources) as sch: - sch.poll(t0 + timedelta(minutes=25)) # candidate 12:20 - assert st.candidates["a"] is not None - with Scheduler(st, "p1", 0, sources) as sch2: # restart, still no capacity - sch2.recover(t0 + timedelta(minutes=26)) - sch2.poll(t0 + timedelta(minutes=26)) - assert st.candidates["a"] is not None - assert st.candidates["a"].intended == t0 + timedelta(minutes=20) # type: ignore[union-attr] - # Clock rollback: consumed instants are never re-admitted. - sch2.poll(t0 - timedelta(hours=1)) - admitted = [r for r in st.history if r.kind == "admitted"] - assert admitted == [] - - -def test_long_downtime_bounded_and_summarized(): - st = make_store() - sources: dict = {} - src = PeriodicSource(timedelta(minutes=1), ts(2023, 9, 8, 12, 0)) - add_sched(st, "m", src, sources, ts(2023, 9, 8, 12, 0), misfire="latest") - with Scheduler(st, "p1", 4, sources, {"*": "complete"}) as sch: - sch.poll(ts(2026, 9, 8, 12, 0, 0)) # 3 years of per-minute misses - total_calls = src.next_calls + src.prev_calls - assert total_calls <= SCAN_CAP + 2, f"must not enumerate: {total_calls} calls" - # Capacity available: latest means ONE prompt admission for the latest - # missed instant (11:59), using the candidate's intended time. - admitted = [r for r in st.history if r.kind == "admitted"] - assert len(admitted) == 1 - assert admitted[0].intended == ts(2026, 9, 8, 11, 59) - assert st.candidates["m"] is None # consumed by the prompt admission - summaries = [r for r in st.history if r.kind == "interval-summary"] - assert len(summaries) == 1, "one coalesced summary, not per-minute rows" - # Same gap with no capacity: exactly ONE held candidate, zero admissions. - st1b = make_store() - sources1b: dict = {} - src1b = PeriodicSource(timedelta(minutes=1), ts(2023, 9, 8, 12, 0)) - add_sched(st1b, "m", src1b, sources1b, ts(2023, 9, 8, 12, 0), misfire="latest") - with Scheduler(st1b, "p1", 0, sources1b, {"*": "complete"}) as sch: - sch.poll(ts(2026, 9, 8, 12, 0, 0)) - assert src1b.next_calls + src1b.prev_calls <= SCAN_CAP + 2 - assert st1b.candidates["m"] is not None - assert st1b.candidates["m"].intended == ts(2026, 9, 8, 11, 59) # type: ignore[union-attr] - assert [r for r in st1b.history if r.kind == "admitted"] == [] - # skip policy: same boundedness, zero admissions, next future selected. - st2 = make_store() - sources2: dict = {} - src2 = PeriodicSource(timedelta(minutes=1), ts(2023, 9, 8, 12, 0)) - add_sched(st2, "m", src2, sources2, ts(2023, 9, 8, 12, 0), misfire="skip") - with Scheduler(st2, "p1", 4, sources2, {"*": "complete"}) as sch: - sch.poll(ts(2026, 9, 8, 12, 0, 30)) - assert src2.next_calls + src2.prev_calls <= SCAN_CAP + 2 - assert [r for r in st2.history if r.kind == "admitted"] == [] - assert st2.candidates["m"] is None - - -def test_fairness_frequent_schedule_cannot_monopolize(): - st = make_store() - sources: dict = {} - t0 = ts(2026, 9, 8, 12, 0) - add_sched( - st, - "fast", - PeriodicSource(timedelta(minutes=1), t0), - sources, - t0 - timedelta(minutes=1), - ) - add_sched(st, "slow", OneShotSource(t0), sources, t0 - timedelta(hours=1)) - with Scheduler(st, "p1", 1, sources, {"*": "hang"}) as sch: - sch.poll(t0) # rotation starts at fast (sorted first): fast admitted - assert any(r.sched_id == "fast" and r.kind == "admitted" for r in st.history) - # Complete fast's run externally, next poll must serve slow first. - for r in st.runs.values(): - r.state = COMPLETED - res = sch.poll(t0 + timedelta(seconds=30)) - slow_admitted = [ - r for r in st.history if r.sched_id == "slow" and r.kind == "admitted" - ] - assert slow_admitted, f"slow schedule starved: {res}" - - -def test_fault_before_admission_never_dispatches(): - st = make_store() - sources: dict = {} - t0 = ts(2026, 9, 8, 12, 0) - add_sched(st, "a", OneShotSource(t0), sources, t0 - timedelta(hours=1)) - st.faults["admission"] = "before" - with Scheduler(st, "p1", 4, sources, {"*": "complete"}) as sch: - try: - sch.poll(t0) - assert False, "fault must propagate" - except InjectedFault: - pass - assert st.runs == {}, "no dispatch before durable admission" - assert st.admissions == {} - sch.poll(t0) # retry after the fault is clean - assert len(st.runs) == 1 - - -def test_fault_between_admission_and_view_recovers_without_redispatch(): - st = make_store() - sources: dict = {} - t0 = ts(2026, 9, 8, 12, 0) - add_sched(st, "a", OneShotSource(t0), sources, t0 - timedelta(hours=1)) - st.faults["materialize"] = "before" - with Scheduler(st, "p1", 4, sources, {"*": "complete"}) as sch: - try: - sch.poll(t0) - assert False - except InjectedFault: - pass - assert len(st.admissions) == 1 and len(st.runs) == 0 - # Recovery NEVER executes: it completes the view and flags it - # pending. No outcome exists yet. - diags = sch.recover(t0) - assert any("pending-dispatch" in d for d in diags) - assert len(st.runs) == 1 - run = st.runs["run-a-1"] - assert run.state == RUNNING and run.needs_dispatch - assert not run.dispatched_unknown - assert [r for r in st.history if r.kind == "completed"] == [] - # The next poll dispatches the captured invocation through the - # capacity checks — exactly once, no second admission. - sch.poll(t0) - assert run.state == COMPLETED and not run.needs_dispatch - n_admitted = len([r for r in st.history if r.kind == "admitted"]) - assert n_admitted == 1, "reconciliation must be idempotent" - assert len([r for r in st.history if r.kind == "completed"]) == 1 - - -def test_recovery_pending_dispatch_waits_for_capacity(): - st = make_store() - sources: dict = {} - t0 = ts(2026, 9, 8, 12, 0) - add_sched(st, "a", OneShotSource(t0), sources, t0 - timedelta(hours=1)) - st.faults["materialize"] = "before" - with Scheduler(st, "p1", 4, sources, {"*": "complete"}) as sch: - try: - sch.poll(t0) - assert False - except InjectedFault: - pass - sch.recover(t0) - sch.capacity = 0 # slots full: pending dispatch must wait - sch.poll(t0 + timedelta(seconds=1)) - assert st.runs["run-a-1"].needs_dispatch - assert [r for r in st.history if r.kind == "completed"] == [] - sch.capacity = 4 - sch.poll(t0 + timedelta(seconds=2)) - assert not st.runs["run-a-1"].needs_dispatch - assert st.runs["run-a-1"].state == COMPLETED - - -def test_run_ids_survive_restart_without_reuse(): - st = make_store() - sources: dict = {} - t0 = ts(2026, 9, 8, 12, 0) - add_sched( - st, - "a", - PeriodicSource(timedelta(minutes=10), t0), - sources, - t0 - timedelta(minutes=10), - ) - with Scheduler(st, "p1", 4, sources, {"*": "hang"}) as sch: - sch.poll(t0) - assert st.runs["run-a-1"].intended == t0 - # Restart: the counter lives in the store, so the next admission gets - # a fresh identity instead of overwriting run-a-1. - with Scheduler(st, "p1", 4, sources, {"*": "hang"}) as sch2: - sch2.recover(t0 + timedelta(seconds=1)) # hanging run abandoned - assert st.runs["run-a-1"].state == FAILED - sch2.poll(t0 + timedelta(minutes=10)) - assert st.runs["run-a-1"].intended == t0, "old run untouched" - assert st.runs["run-a-2"].intended == t0 + timedelta(minutes=10) - assert len(st.runs) == 2 - - -def test_crash_after_dispatch_marks_abandoned_failed_without_replay(): - st = make_store() - sources: dict = {} - t0 = ts(2026, 9, 8, 12, 0) - add_sched( - st, - "a", - PeriodicSource(timedelta(hours=1), t0), - sources, - t0 - timedelta(hours=1), - ) - with Scheduler(st, "p1", 4, sources, {"run-a-1": "crash", "*": "complete"}) as sch: - try: - sch.poll(t0) - assert False - except ExecutorCrashed: - pass - run = st.runs["run-a-1"] - assert run.state == RUNNING and run.dispatched_unknown - diags = sch.recover(t0 + timedelta(seconds=5)) - assert run.state == FAILED - assert "may already have occurred" in run.fail_reason - assert any("failed-closed" in d for d in diags) - # Future occurrences proceed; the failed one is never replayed. - sch.poll(t0 + timedelta(hours=1)) - intents = [r.intended for r in st.history if r.kind == "admitted"] - assert t0 not in intents[1:] or intents.count(t0) == 1 - - -def test_crash_during_resume_must_not_look_safe_to_retry(): - st = make_store() - sources: dict = {} - t0 = ts(2026, 9, 8, 12, 0) - add_sched(st, "a", OneShotSource(t0), sources, t0 - timedelta(hours=1)) - with Scheduler(st, "p1", 4, sources, {"*": "interrupt"}) as sch: - sch.poll(t0) - (run_id,) = list(st.runs) - assert st.runs[run_id].state == INTERRUPTED - # Restart: a merely-waiting interruption is safe to resume later... - with Scheduler(st, "p1", 4, sources, {"*": "crash"}) as sch2: - sch2.recover(t0) - assert st.runs[run_id].state == INTERRUPTED - assert run_id not in st.resume_attempts - try: - sch2.resume_run(run_id, t0) - assert False - except ExecutorCrashed: - pass - # ...but the crash left a durably marked ACTIVE attempt: recovery - # must fail it closed instead of presenting the old checkpoint again. - assert st.resume_attempts[run_id] == "ACTIVE" - diags = sch2.recover(t0) - assert st.runs[run_id].state == FAILED - assert "ambiguous" in st.runs[run_id].fail_reason - assert any("failed-closed" in d for d in diags) - - -def test_preflight_rejection_invents_no_run_and_freezes_input(): - st = make_store() - sources: dict = {} - t0 = ts(2026, 9, 8, 12, 0) - add_sched( - st, - "a", - PeriodicSource(timedelta(hours=1), t0), - sources, - t0 - timedelta(hours=1), - misfire="latest", - ) - with Scheduler(st, "p1", 0, sources) as sch: - sch.poll(t0 + timedelta(minutes=5)) # held candidate @12:00 - assert st.candidates["a"] is not None - # Deployment edit changes the expected input before admission. - st.deployments["dep-1"] = { - "rev": 2, - "required": ["team", "report_time", "region"], - } - sch2 = Scheduler(st, "p1", 4, sources, {"*": "complete"}) - sch2.poll(t0 + timedelta(minutes=6)) - assert st.runs == {}, "preflight rejection must invent no run" - assert any(r.kind == "preflight-rejected" for r in st.history) - sch2.close() - # Admitted runs freeze their invocation: later edits change nothing. - st3 = make_store() - sources3: dict = {} - add_sched(st3, "a", OneShotSource(t0), sources3, t0 - timedelta(hours=1)) - with Scheduler(st3, "p1", 4, sources3, {"*": "hang"}) as sch: - sch.poll(t0) - (run_id,) = list(st3.runs) - before = dict(st3.runs[run_id].frozen_input) - st3.deployments["dep-1"] = {"rev": 9, "required": ["team"]} - st3.runs[run_id].state = COMPLETED - assert st3.runs[run_id].frozen_input == before - - -def test_exclusive_ownership_and_unsupported_locking(): - st = make_store() - sources: dict = {} - sch1 = Scheduler(st, "proc-A", 1, sources) - try: - Scheduler(st, "proc-B", 1, sources) - assert False, "second owner must be rejected" - except SecondOwnerError: - pass - # A held lock never expires while the owner lives (no lease timeout). - assert not hasattr(st, "lease_expiry") - sch1.close() # process death releases - sch2 = Scheduler(st, "proc-B", 1, sources) # new owner starts, recovers - sch2.close() - bad = MemStore(lockable=False) - try: - Scheduler(bad, "proc-C", 1, {}) - assert False, "unsupported locking must reject scheduler startup" - except StartupRejected: - pass - - -def test_corrupt_view_without_admission_fails_closed(): - st = make_store() - sources: dict = {} - t0 = ts(2026, 9, 8, 12, 0) - add_sched( - st, - "a", - PeriodicSource(timedelta(hours=1), t0), - sources, - t0 - timedelta(hours=1), - ) - st.runs["ghost-1"] = Run("ghost-1", "a", t0, 1, {"team": "eng"}) - with Scheduler(st, "p1", 4, sources, {"*": "complete"}) as sch: - diags = sch.recover(t0) - assert any("corrupt-blocked" in d for d in diags) - assert st.schedules["a"].blocked_reason is not None - try: - sch.poll(t0 + timedelta(hours=1)) - assert False, "corrupt records fail closed with diagnostics" - except BlockedSchedule: - pass - - -def test_overlap_parallel_x_misfire_latest(): - st = make_store() - sources: dict = {} - t0 = ts(2026, 9, 8, 12, 0) - add_sched( - st, - "p", - PeriodicSource(timedelta(minutes=5), t0), - sources, - t0 - timedelta(minutes=5), - overlap="parallel", - max_active=2, - misfire="latest", - ) - with Scheduler(st, "p1", 0, sources) as sch: # no task slots: hold - sch.poll(t0 + timedelta(minutes=12)) # 12:00,05,10 missed - cand = st.candidates["p"] - assert cand is not None and cand.intended == t0 + timedelta(minutes=10) - assert [r for r in st.history if r.kind == "admitted"] == [] - with Scheduler(st, "p1", 1, sources, {"*": "hang"}) as sch: - sch.poll(t0 + timedelta(minutes=13)) # admits held 12:10, hangs - assert st.runs["run-p-1"].intended == t0 + timedelta(minutes=10) - # Task slot taken: 12:15 is held as the one latest candidate. - sch.poll(t0 + timedelta(minutes=15)) - assert list(st.runs) == ["run-p-1"] - assert st.candidates["p"] is not None - assert st.candidates["p"].intended == t0 + timedelta(minutes=15) # type: ignore[union-attr] - # Slot frees: the held 12:15 candidate is admitted (hangs). - st.runs["run-p-1"].state = COMPLETED - sch.poll(t0 + timedelta(minutes=16)) - assert st.runs["run-p-2"].intended == t0 + timedelta(minutes=15) - # run-p-2 waits durably: schedule slot held, task slot free. Admit a - # second hanging run to reach the cap, then the next due instant is - # terminally skipped-overlap. - st.runs["run-p-2"].state = INTERRUPTED - st.history.append( - Record( - kind="interrupted", - sched_id="p", - intended=t0 + timedelta(minutes=15), - run_id="run-p-2", - ) - ) - sch.poll(t0 + timedelta(minutes=20)) # admits 12:20, hangs - assert st.runs["run-p-3"].intended == t0 + timedelta(minutes=20) - sch.poll(t0 + timedelta(minutes=25)) # 2 active >= max: terminal skip - assert any( - r.kind == "skipped-overlap" and r.intended == t0 + timedelta(minutes=25) - for r in st.history - ) - # And a later catch-up admits the earliest missed instant ASAP - # (12:30), holds only the newest (12:45), and never resurrects the - # skipped 12:25. - st.runs["run-p-2"].state = COMPLETED - st.runs["run-p-3"].state = COMPLETED - sch.poll(t0 + timedelta(minutes=45)) # 12:30..45 missed - just = [ - r - for r in st.history - if r.kind == "admitted" and r.intended == t0 + timedelta(minutes=30) - ] - assert len(just) == 1 - cand = st.candidates["p"] - assert cand is not None and cand.intended == t0 + timedelta(minutes=45) - assert all( - r.intended != t0 + timedelta(minutes=25) - for r in st.history - if r.kind == "admitted" - ) - - -def test_manual_and_other_schedule_runs_are_overlap_independent(): - st = make_store() - sources: dict = {} - t0 = ts(2026, 9, 8, 12, 0) - add_sched(st, "a", OneShotSource(t0), sources, t0 - timedelta(hours=1)) - add_sched(st, "b", OneShotSource(t0), sources, t0 - timedelta(hours=1)) - # Manual runs and other schedules' runs never join this schedule's check. - st.runs["manual-1"] = Run("manual-1", "manual", t0, 1, {"team": "eng"}) - st.runs["run-b-0"] = Run("run-b-0", "b", t0, 1, {"team": "eng"}) - with Scheduler(st, "p1", 4, sources, {"*": "complete"}) as sch: - sch.poll(t0) - admitted_a = [ - r for r in st.history if r.kind == "admitted" and r.sched_id == "a" - ] - assert len(admitted_a) == 1, "overlap=skip ignores manual/other runs" - # ...while b is blocked by its OWN running run (per-schedule scope - # cuts both ways: b's check sees run-b-0, a's check does not). - skipped_b = [ - r for r in st.history if r.kind == "skipped-overlap" and r.sched_id == "b" - ] - assert len(skipped_b) == 1 and skipped_b[0].intended == t0 - - -def test_capacity_wait_then_admit_or_expire_for_skip(): - st = make_store() - sources: dict = {} - t0 = ts(2026, 9, 8, 12, 0) - add_sched(st, "a", OneShotSource(t0), sources, t0 - timedelta(hours=1)) - with Scheduler(st, "p1", 0, sources) as sch: # full: within allowance - assert sch.poll(t0) == {"a": "admit:held-undecided"} - assert st.consumed["a"] < t0, "undecided instants stay unconsumed" - assert st.candidates["a"] is None, "skip holds no candidate" - with Scheduler(st, "p1", 1, sources, {"*": "complete"}) as sch: - sch.poll(t0 + timedelta(seconds=30)) # still within allowance - admitted = [r for r in st.history if r.kind == "admitted"] - assert len(admitted) == 1 and admitted[0].intended == t0 - # Past the deadline instead: capacity-delayed skip expires. - st2 = make_store() - sources2: dict = {} - add_sched( - st2, - "a", - PeriodicSource(timedelta(minutes=10), t0), - sources2, - t0 - timedelta(hours=1), - ) - with Scheduler(st2, "p1", 0, sources2) as sch: - sch.poll(t0) - assert sch.poll(t0 + timedelta(seconds=61))["a"] == "skipped-misfire" - assert any( - r.kind == "skipped-misfire" and r.intended == t0 for r in st2.history - ) - - -def test_fault_after_admission_recovers_exactly_once(): - st = make_store() - sources: dict = {} - t0 = ts(2026, 9, 8, 12, 0) - add_sched(st, "a", OneShotSource(t0), sources, t0 - timedelta(hours=1)) - st.faults["admission"] = "after" # record persisted, crash before view - with Scheduler(st, "p1", 4, sources, {"*": "complete"}) as sch: - try: - sch.poll(t0) - assert False - except InjectedFault: - pass - assert len(st.admissions) == 1 and len(st.runs) == 0 - # Recovery materializes the view but never executes: pending. - diags = sch.recover(t0) - assert any("pending-dispatch" in d for d in diags) - run = st.runs["run-a-1"] - assert run.state == RUNNING and run.needs_dispatch - assert [r for r in st.history if r.kind == "completed"] == [] - sch.poll(t0) # sweep dispatches exactly once - assert run.state == COMPLETED and not run.needs_dispatch - assert len([r for r in st.history if r.kind == "admitted"]) == 1 - assert len([r for r in st.history if r.kind == "completed"]) == 1 - - -def test_fault_after_complete_reconciles_terminal_record(): - st = make_store() - sources: dict = {} - t0 = ts(2026, 9, 8, 12, 0) - add_sched(st, "a", OneShotSource(t0), sources, t0 - timedelta(hours=1)) - st.faults["complete"] = "after" # COMPLETED persisted, record lost - with Scheduler(st, "p1", 4, sources, {"*": "complete"}) as sch: - try: - sch.poll(t0) - assert False - except InjectedFault: - pass - assert st.runs["run-a-1"].state == COMPLETED - assert [r for r in st.history if r.kind == "completed"] == [] - diags = sch.recover(t0) - assert any("terminal-reconciled" in d for d in diags) - assert len([r for r in st.history if r.kind == "completed"]) == 1 - sch.poll(t0 + timedelta(hours=1)) # consumed advanced: no re-admit - assert len([r for r in st.history if r.kind == "admitted"]) == 1 - - -def test_fault_after_resume_mark_fails_closed_not_retried(): - st = make_store() - sources: dict = {} - t0 = ts(2026, 9, 8, 12, 0) - add_sched(st, "a", OneShotSource(t0), sources, t0 - timedelta(hours=1)) - with Scheduler(st, "p1", 4, sources, {"*": "interrupt"}) as sch: - sch.poll(t0) - (run_id,) = list(st.runs) - with Scheduler(st, "p1", 4, sources, {"*": "complete"}) as sch: - st.faults["resume-mark"] = "after" # ACTIVE persisted, never executed - try: - sch.resume_run(run_id, t0) - assert False - except InjectedFault: - pass - assert st.runs[run_id].state == INTERRUPTED # never re-ran - assert st.resume_attempts[run_id] == "ACTIVE" - diags = sch.recover(t0) # conservative: ambiguous, never retried - assert st.runs[run_id].state == FAILED - assert "ambiguous" in st.runs[run_id].fail_reason - assert any("failed-closed" in d for d in diags) - - -def test_fault_after_interrupt_reconciles_waiting_state(): - st = make_store() - sources: dict = {} - t0 = ts(2026, 9, 8, 12, 0) - add_sched(st, "a", OneShotSource(t0), sources, t0 - timedelta(hours=1)) - st.faults["interrupt-persist"] = "after" # INTERRUPTED kept, record lost - with Scheduler(st, "p1", 4, sources, {"*": "interrupt"}) as sch: - try: - sch.poll(t0) - assert False - except InjectedFault: - pass - assert st.runs["run-a-1"].state == INTERRUPTED - diags = sch.recover(t0) - assert st.runs["run-a-1"].state == INTERRUPTED # still resumable - assert any("terminal-reconciled" in d for d in diags) - # A resumed run may interrupt AGAIN: re-interruption is a durable - # terminal persist of its own, and the run stays resumable after it. - assert sch.resume_run("run-a-1", t0) == "resumed-interrupted" - assert st.runs["run-a-1"].state == INTERRUPTED - assert st.resume_attempts["run-a-1"] == "DONE" - assert any( - r.kind == "interrupted" and r.reason == "resumed-reinterrupted" - for r in st.history - ) - sch.outcomes["run-a-1"] = "complete" - assert sch.resume_run("run-a-1", t0) == "resumed-complete" - - -def _interrupted_run() -> tuple[MemStore, dict, str]: - st = make_store() - sources: dict = {} - t0 = ts(2026, 9, 8, 12, 0) - add_sched(st, "a", OneShotSource(t0), sources, t0 - timedelta(hours=1)) - with Scheduler(st, "p1", 4, sources, {"*": "interrupt"}) as sch: - sch.poll(t0) - (run_id,) = list(st.runs) - assert st.runs[run_id].state == INTERRUPTED - return st, sources, run_id - - -def test_fault_after_resume_complete_reconciles_without_rerun(): - st, sources, run_id = _interrupted_run() - t0 = ts(2026, 9, 8, 12, 0) - with Scheduler(st, "p1", 4, sources, {"*": "complete"}) as sch: - st.faults["resume-complete"] = "after" # COMPLETED kept, rest lost - try: - sch.resume_run(run_id, t0) - assert False - except InjectedFault: - pass - assert st.runs[run_id].state == COMPLETED - assert st.resume_attempts[run_id] == "ACTIVE" - assert [r for r in st.history if r.kind == "completed"] == [] - diags = sch.recover(t0) - # Reconciled, never re-executed: exactly one completion, DONE marker. - assert st.resume_attempts[run_id] == "DONE" - assert any("attempt-reconciled" in d for d in diags) - assert len([r for r in st.history if r.kind == "completed"]) == 1 - sch.poll(t0 + timedelta(hours=1)) - assert len([r for r in st.history if r.kind == "completed"]) == 1 - - -def test_fault_after_resume_attempt_clear_reconciles_record(): - st, sources, run_id = _interrupted_run() - t0 = ts(2026, 9, 8, 12, 0) - with Scheduler(st, "p1", 4, sources, {"*": "complete"}) as sch: - st.faults["resume-attempt-clear"] = "after" # DONE kept, record lost - try: - sch.resume_run(run_id, t0) - assert False - except InjectedFault: - pass - assert st.runs[run_id].state == COMPLETED - assert st.resume_attempts[run_id] == "DONE" - assert [r for r in st.history if r.kind == "completed"] == [] - diags = sch.recover(t0) - assert any("terminal-reconciled" in d for d in diags) - assert len([r for r in st.history if r.kind == "completed"]) == 1 - - -def test_fault_before_resume_interrupt_fails_closed(): - st, sources, run_id = _interrupted_run() - t0 = ts(2026, 9, 8, 12, 0) - with Scheduler(st, "p1", 4, sources, {"run-a-1": "interrupt"}) as sch: - st.faults["resume-interrupt"] = "before" # re-interrupt never persisted - try: - sch.resume_run(run_id, t0) - assert False - except InjectedFault: - pass - # The ACTIVE marker superseded the old waiting checkpoint, but the - # re-interruption never landed: fail closed, never retry the old one. - assert st.runs[run_id].state == RUNNING - assert st.resume_attempts[run_id] == "ACTIVE" - assert st.runs[run_id].result_attempt != st.runs[run_id].attempt_id - diags = sch.recover(t0) - assert st.runs[run_id].state == FAILED - assert "ambiguous" in st.runs[run_id].fail_reason - assert any("failed-closed" in d for d in diags) - - -def test_fault_after_resume_interrupt_stays_resumable(): - # Exact repro shape: the re-interruption persisted WITH the attempt - # identity, then the crash hit before attempt-clearing. Recovery must - # match result to attempt and keep the run resumable — not fail it. - st, sources, run_id = _interrupted_run() - t0 = ts(2026, 9, 8, 12, 0) - with Scheduler(st, "p1", 4, sources, {"run-a-1": "interrupt"}) as sch: - st.faults["resume-interrupt"] = "after" - try: - sch.resume_run(run_id, t0) - assert False - except InjectedFault: - pass - assert st.runs[run_id].state == INTERRUPTED - assert st.resume_attempts[run_id] == "ACTIVE" - assert st.runs[run_id].result_attempt == st.runs[run_id].attempt_id - assert st.runs[run_id].attempt_id > 0 - diags = sch.recover(t0) - assert st.runs[run_id].state == INTERRUPTED, "fresh result: resumable" - assert st.resume_attempts[run_id] == "DONE" - assert any("fresh-result-resumable" in d for d in diags) - assert len([r for r in st.history if r.kind == "interrupted"]) == 1 - sch.outcomes[run_id] = "complete" - assert sch.resume_run(run_id, t0) == "resumed-complete" diff --git a/src/wf_scheduling/poll.py b/src/wf_scheduling/poll.py index a2c2ae7f..0705450e 100644 --- a/src/wf_scheduling/poll.py +++ b/src/wf_scheduling/poll.py @@ -1,8 +1,9 @@ """Poll loop: overlap, misfire, candidates, fairness, capacity (T08). -Mirrors the reference state model (probes/deployment_scheduling_verify/ -test_schedule_state_model.py) against real file stores. Calendar iteration -uses the canonical :class:`wf_scheduling.calendar.OccurrenceSource` +Implements the scheduling state rules (first probed as a reference model, +retired to docs/historical now that tests/scheduling/ pins them) against +real file stores. Calendar iteration uses the canonical +:class:`wf_scheduling.calendar.OccurrenceSource` (``next_after``/``prev_before`` only, never enumeration); latest-missed catch-up is one bounded ``prev_before`` query (F1). Overlap decisions precede capacity checks; terminal skips never reappear; ``latest`` retains diff --git a/src/wf_server/cli.py b/src/wf_server/cli.py index 62e63faa..067d3da7 100644 --- a/src/wf_server/cli.py +++ b/src/wf_server/cli.py @@ -69,6 +69,7 @@ def serve( server = None workflow_config: WorkflowConfigFile | None = None + mcp_backed = mcp_config is not None if mcp_config is not None: server = build_workflow_server_from_legacy_mcp_config(mcp_config) @@ -118,6 +119,22 @@ def serve( server = build_local_static_workflow_server(resolved_store_root, drafts=True) sched_config = server_scheduler_config(workflow_config, enable_scheduler) + if sched_config is not None: + config_mcp_sources = ( + workflow_config is not None + and any( + getattr(source, "kind", None) == "mcp" + for source in workflow_config.server.sources + ) + ) + if mcp_backed or config_mcp_sources: + # The scheduler is verified over local/static servers only: + # refuse to run it over an MCP-backed runtime instead of + # operating it untested. + raise typer.BadParameter( + "--enable-scheduler requires a local/static server; " + "MCP-backed servers are not supported yet" + ) rpc_app = create_rpc_app( server, rpc_path=resolved_rpc_path, diff --git a/tests/artifacts/test_run_store.py b/tests/artifacts/test_run_store.py index b5290d20..03e2047c 100644 --- a/tests/artifacts/test_run_store.py +++ b/tests/artifacts/test_run_store.py @@ -117,3 +117,12 @@ def test_file_run_store_rejects_unsafe_run_id(tmp_path) -> None: def test_workflow_run_record_validates_latest_checkpoint_id() -> None: with pytest.raises(ValidationError): run_record("run_123", "../outside") + + +def test_allocate_run_id_survives_fresh_instances_without_reuse(tmp_path) -> None: + first = FileRunStore(tmp_path) + assert first.allocate_run_id() == "run-000001" + assert first.allocate_run_id() == "run-000002" + # A fresh instance (restart boundary) must not reuse identities. + second = FileRunStore(tmp_path) + assert second.allocate_run_id() == "run-000003" diff --git a/tests/examples/test_scheduled_deployment_example.py b/tests/examples/test_scheduled_deployment_example.py new file mode 100644 index 00000000..f02a6f41 --- /dev/null +++ b/tests/examples/test_scheduled_deployment_example.py @@ -0,0 +1,136 @@ +"""Executable scheduling example: one deployment run on a timer. + +This is the runnable companion to +``docs/deployment_scheduling.md`` ("hypothetically used as follows"): +a local server plus the opt-in scheduler admit and complete one +one-shot scheduled run of a constant workflow, then show its occurrence +history. Run it with:: + + uv run pytest -q tests/examples/test_scheduled_deployment_example.py +""" + +from __future__ import annotations + +import asyncio +from datetime import UTC, datetime, timedelta +from typing import Any + +from wf_api.models import RawWorkflowPlan +from wf_core import END +from wf_server import build_local_static_workflow_server +from wf_server.scheduling import build_scheduler_service +from wf_scheduling.lifecycle import SchedulerServiceConfig +from wf_scheduling.store import FileScheduleStore + + +def _constant_plan() -> RawWorkflowPlan: + return RawWorkflowPlan.model_validate( + { + "name": "scheduled_hello", + "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": "hello on a schedule", + "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"]}, + } + ], + } + ) + + +async def _wait_for(condition: Any, timeout: float = 30.0) -> None: + async with asyncio.timeout(timeout): + while not condition(): + await asyncio.sleep(0.02) + + +async def test_scheduled_deployment_completes_on_a_timer(tmp_path: Any) -> None: + root = tmp_path / "store" + server = build_local_static_workflow_server(root, schedules=True) + await server.api.create_artifact_from_plan( + artifact_id="scheduled_hello", + version=1, + title="Scheduled Hello", + plan=_constant_plan(), + outcomes=["ok"], + source_bindings={}, + ) + await server.api.save_deployment( + { + "id": "scheduled_hello.default", + "artifact_id": "scheduled_hello", + "artifact_version": 1, + "bindings": {}, + } + ) + due = datetime.now(UTC) + timedelta(seconds=0.5) + created = await server.api.schedules.create_schedule( + schedule_id="hello-once", + deployment_id="scheduled_hello.default", + trigger={"kind": "oneshot", "at": due.isoformat()}, + ) + assert created["revision"] == 1 + + service = build_scheduler_service( + server, + SchedulerServiceConfig(poll_interval_s=0.05, capacity=2), + ) + try: + await service.start() + + def _completed() -> bool: + runs = list(server.stores.run_store.list_runs()) + return any(run.status.value == "completed" for run in runs) + + await _wait_for(_completed) + run_id = next( + run.id + for run in server.stores.run_store.list_runs() + if run.status.value == "completed" + ) + inspected = await server.api.inspect_run(run_id=run_id) + assert inspected["status"] == "completed" + assert (inspected["output"] or {})["result"] == "hello on a schedule" + + page = await server.api.schedules.list_schedule_occurrences( + schedule_id="hello-once" + ) + kinds = [row["kind"] for row in page["occurrences"]] + assert "admitted" in kinds + assert "completed" in kinds + finally: + await service.stop() + + # The schedule definition survives its run; history is retained. + assert FileScheduleStore(root).get_schedule("hello-once").exhausted is True diff --git a/tests/scheduling/test_poll.py b/tests/scheduling/test_poll.py index 70073589..0d2c1ff8 100644 --- a/tests/scheduling/test_poll.py +++ b/tests/scheduling/test_poll.py @@ -353,3 +353,34 @@ def test_poll_one_uses_fresh_definition_after_edit(tmp_path: Path) -> None: assert admission.resolved_input["team"] == "new" assert admission.schedule_revision == 2 sched.ownership.release() + + +def test_manual_and_other_schedule_runs_are_overlap_independent( + tmp_path: Path, +) -> None: + from wf_api.run_lifecycle import ( + materialize_admitted_view, + persist_admission, + ) + + sched, store, runs, sources = _harness(tmp_path, script={"*": "hang"}) + t0 = ts(2026, 9, 8, 12, 0) + # An active run of another schedule occupies only its own slot. + _add(sched, store, sources, "b", OneShotSource(t0), t0 - timedelta(hours=1)) + assert sched.poll(t0)["b"].startswith("admit:run-") + # A manual run (no schedule owner) participates in no overlap check. + manual_id = runs.allocate_run_id() + manual = persist_admission( + store=runs, + run_id=manual_id, + environment=fixture_environment(object()), + resolved_input={}, + max_steps=None, + ) + materialize_admitted_view(store=runs, admission=manual) + # A due one-shot admits despite both unrelated active runs. + _add(sched, store, sources, "a", OneShotSource(t0), t0 - timedelta(hours=1)) + result = sched.poll(t0 + timedelta(seconds=1)) + assert result["a"].startswith("admit:run-") + assert result["b"] == "exhausted" + sched.ownership.release() diff --git a/tests/wf_server/test_cli.py b/tests/wf_server/test_cli.py index 2b1c5248..8d3a8243 100644 --- a/tests/wf_server/test_cli.py +++ b/tests/wf_server/test_cli.py @@ -565,3 +565,81 @@ def test_rpc_server_cli_flag_overrides_disabled_config_scheduler( assert result.exit_code == 0, result.output assert captured["lifespan"] is not None + + +def test_rpc_server_cli_enable_scheduler_rejects_mcp_backed_server( + monkeypatch, tmp_path +) -> None: + config_path = tmp_path / "wf_mcp.config.json" + config_path.write_text( + json.dumps({"store_root": str(tmp_path / "store"), "connections": []}), + encoding="utf-8", + ) + + def fake_build_mcp_server(path): + return object() + + monkeypatch.setattr( + "wf_server.cli.build_workflow_server_from_legacy_mcp_config", + fake_build_mcp_server, + ) + + result = CliRunner().invoke( + app, + [ + "--mcp-config", + str(config_path), + "--enable-scheduler", + ], + ) + + assert result.exit_code != 0 + assert "requires a local/static server" in result.output + + +def test_rpc_server_cli_config_mcp_sources_reject_scheduler( + monkeypatch, tmp_path +) -> None: + config_path = tmp_path / "wf.json" + config_path.write_text( + json.dumps( + { + "version": 1, + "server": { + "store": {"kind": "filesystem", "root": ".wf_store"}, + "sources": [ + { + "kind": "mcp", + "id": "everything.default", + "provider": "everything", + "account": "default", + "transport": { + "kind": "stdio", + "command": "uvx", + "args": ["mcp-server-everything"], + }, + } + ], + "scheduler": {"enabled": True}, + }, + } + ), + encoding="utf-8", + ) + + def fake_build_server(config, *, drafts=False): + return object() + + def fake_create_rpc_app(server, *, rpc_path="/rpc", drafts=False, lifespan=None): + return object() + + monkeypatch.setattr( + "wf_server.cli.build_workflow_server_from_workflow_config", + fake_build_server, + ) + monkeypatch.setattr("wf_server.cli.create_rpc_app", fake_create_rpc_app) + + result = CliRunner().invoke(app, ["--config", str(config_path)]) + + assert result.exit_code != 0 + assert "requires a local/static server" in result.output