From 993ed07fd33ab77861d6dbcc1b8c71c832215967 Mon Sep 17 00:00:00 2001 From: lda Date: Tue, 8 Sep 2026 11:24:40 +0700 Subject: [PATCH] sched: require proven ownership at poll/admin/recovery entries (F5) --- src/wf_scheduling/ownership.py | 10 ++ src/wf_scheduling/poll.py | 18 +++ src/wf_scheduling/recovery.py | 8 + tests/scheduling/test_ownership_guard.py | 191 +++++++++++++++++++++++ tests/scheduling/test_poll.py | 9 ++ tests/scheduling/test_preparation.py | 8 + tests/scheduling/test_recovery.py | 71 ++++++--- 7 files changed, 293 insertions(+), 22 deletions(-) create mode 100644 tests/scheduling/test_ownership_guard.py diff --git a/src/wf_scheduling/ownership.py b/src/wf_scheduling/ownership.py index 38dfb976..6b4ca0c6 100644 --- a/src/wf_scheduling/ownership.py +++ b/src/wf_scheduling/ownership.py @@ -39,6 +39,16 @@ class SchedulerOwnership: def lock_path(self) -> Path: return self.root / "scheduler.lock" + @property + def held(self) -> bool: + """Whether this process holds the lock through this object. + + Only a successful :meth:`acquire` sets this: a second process can + never observe ``held`` while another owner holds the OS lock, so + entry-point guards can treat it as proof of exclusive ownership. + """ + return self._locked + def acquire(self) -> SchedulerOwnership: """Acquire the held lock non-blockingly or raise SecondOwnerError.""" self.lock_path.parent.mkdir(parents=True, exist_ok=True) diff --git a/src/wf_scheduling/poll.py b/src/wf_scheduling/poll.py index 87462829..9de3801b 100644 --- a/src/wf_scheduling/poll.py +++ b/src/wf_scheduling/poll.py @@ -27,6 +27,7 @@ from wf_scheduling.calendar import ( from wf_scheduling.dispatch import RunDispatcher, StillRunning from wf_scheduling.models import OccurrenceRecord, PendingCandidate from wf_scheduling.occurrences import occurrence_id +from wf_scheduling.ownership import SchedulerOwnership, SecondOwnerError from wf_scheduling.prepare import InvocationPreparer, PreparationRejected UTC = timezone.utc @@ -84,6 +85,7 @@ class Scheduler: capacity: int, preparer: InvocationPreparer, dispatcher: RunDispatcher, + ownership: SchedulerOwnership, ) -> None: self.schedule_store = schedule_store self.run_store = run_store @@ -91,8 +93,21 @@ class Scheduler: self.capacity = capacity self.preparer = preparer self.dispatcher = dispatcher + self.ownership = ownership self._poll_cursor = 0 + def _require_ownership(self) -> None: + """Reject schedule mutation/dispatch without proven live ownership. + + Runs before any store write or dispatcher side effect: without a + held lock this process cannot prove exclusive ownership, so polling + or administering schedules would risk double admission. + """ + if self.ownership is None or not self.ownership.held: + raise SecondOwnerError( + "scheduler ownership is required before polling or mutating schedules" + ) + # -- helpers ------------------------------------------------------ @staticmethod def _status_value(run: Any) -> Any: @@ -342,6 +357,7 @@ class Scheduler: # -- polling -------------------------------------------------------- def poll(self, now: datetime) -> dict[str, str]: + self._require_ownership() self._dispatch_pending(now) schedules = self.schedule_store.list_schedules(include_deleted=True) ids = sorted(item.id for item in schedules) @@ -616,6 +632,7 @@ class Scheduler: # -- administration --------------------------------------------------- def resume_schedule(self, sid: str, now: datetime) -> None: """Unpause: resume selects the next future occurrence.""" + self._require_ownership() sched = self.schedule_store.get_schedule(sid) sched.paused = False self.schedule_store.save_schedule(sched) @@ -625,6 +642,7 @@ class Scheduler: def edit_schedule(self, sid: str, now: datetime) -> None: """Definition edit: new revision, discard old candidates, no backfill.""" + self._require_ownership() sched = self.schedule_store.get_schedule(sid) sched.revision += 1 sched.updated_at = now diff --git a/src/wf_scheduling/recovery.py b/src/wf_scheduling/recovery.py index fef29198..b396b45f 100644 --- a/src/wf_scheduling/recovery.py +++ b/src/wf_scheduling/recovery.py @@ -13,6 +13,8 @@ from __future__ import annotations from datetime import datetime, timezone from typing import Any +from wf_scheduling.ownership import SchedulerOwnership, SecondOwnerError + UTC = timezone.utc ABANDONED_REASON = ( @@ -30,11 +32,17 @@ def recover( schedule_store: Any, run_store: Any, now: datetime, + ownership: SchedulerOwnership, record_history: Any | None = None, ) -> list[str]: """Reconcile durable state after a restart without executing work.""" from wf_artifacts.runs.models import StoredRunStatus + if ownership is None or not ownership.held: + raise SecondOwnerError( + "scheduler ownership is required before recovery: an unowned " + "recovery could abandon or redispatch another owner's work" + ) diags: list[str] = [] # Admission record is the recovery authority: admitted but never # materialized views are completed here and flagged pending for the diff --git a/tests/scheduling/test_ownership_guard.py b/tests/scheduling/test_ownership_guard.py new file mode 100644 index 00000000..fd93e7c5 --- /dev/null +++ b/tests/scheduling/test_ownership_guard.py @@ -0,0 +1,191 @@ +"""Ownership guards at scheduler mutation/recovery/dispatch entries (R4/F5). + +The held lock must be proven at the entry points, not merely exist as a +class: ``poll``, schedule administration, and ``recover`` reject when the +caller cannot prove live ownership, before any write or side effect. +""" + +from __future__ import annotations + +from datetime import UTC, datetime, timedelta +from pathlib import Path +from typing import Any + +import pytest + +from tests.scheduling.controlled import ( + DictDeployments, + ScriptedDispatcher, + fixture_environment, +) +from wf_artifacts.runs.store import FileRunStore +from wf_scheduling import recovery as sched_recovery +from wf_scheduling.calendar import OneShotSource +from wf_scheduling.models import Schedule +from wf_scheduling.ownership import SchedulerOwnership, SecondOwnerError +from wf_scheduling.poll import Scheduler +from wf_scheduling.prepare import SchedulePreparer +from wf_scheduling.store import FileScheduleStore + + +def ts(y: int, mo: int, d: int, h: int = 0, mi: int = 0) -> datetime: + return datetime(y, mo, d, h, mi, tzinfo=UTC) + + +def _sched_model(sid: str, **kw: Any) -> Schedule: + now = ts(2026, 9, 8, 12, 0) + base: dict[str, Any] = { + "id": sid, + "deployment_id": "dep-1", + "trigger": {"kind": "cron", "expression": "0 * * * *", "timezone": "UTC"}, + "input_bindings": [], + "created_at": now.isoformat(), + "updated_at": now.isoformat(), + } + base.update(kw) + return Schedule.model_validate(base) + + +def _scheduler( + tmp_path: Path, + ownership: SchedulerOwnership, + *, + script: dict | None = None, +) -> tuple[Scheduler, FileScheduleStore, FileRunStore]: + sched_store = FileScheduleStore(tmp_path / "sched") + run_store = FileRunStore(tmp_path / "runs") + sched = Scheduler( + schedule_store=sched_store, + run_store=run_store, + sources={}, + capacity=4, + preparer=SchedulePreparer( + DictDeployments({"dep-1": {"rev": 1, "required": []}}), + fixture_environment, + ), + dispatcher=ScriptedDispatcher(script), + ownership=ownership, + ) + return sched, sched_store, run_store + + +def _due_setup(sched: Scheduler, store: FileScheduleStore, intended: datetime) -> None: + store.create_schedule(_sched_model("a")) + store.save_consumed("a", intended - timedelta(hours=1)) + sched.sources["a"] = OneShotSource(intended) + + +def test_poll_without_ownership_rejects_before_writes(tmp_path: Path) -> None: + calls: list[str] = [] + + def spy(admission: Any, now: datetime) -> Any: + calls.append(admission.id) + raise AssertionError("dispatcher must not run without ownership") + + ownership = SchedulerOwnership(tmp_path / "sched", owner="never-acquired") + sched, store, runs = _scheduler(tmp_path, ownership) + sched.dispatcher = ScriptedDispatcher({"*": spy}) + intended = ts(2026, 9, 8, 12, 0) + _due_setup(sched, store, intended) + with pytest.raises(SecondOwnerError): + sched.poll(intended) + assert calls == [] + assert runs.list_runs() == [] + assert runs.list_admissions() == [] + assert store.list_occurrences("a", limit=100)["total"] == 0 + assert store.get_consumed("a") == intended - timedelta(hours=1) + + +def test_poll_with_released_ownership_rejects(tmp_path: Path) -> None: + ownership = SchedulerOwnership(tmp_path / "sched", owner="test").acquire() + sched, store, runs = _scheduler(tmp_path, ownership) + intended = ts(2026, 9, 8, 12, 0) + _due_setup(sched, store, intended) + ownership.release() + assert not ownership.held + with pytest.raises(SecondOwnerError): + sched.poll(intended) + assert runs.list_runs() == [] + assert store.list_occurrences("a", limit=100)["total"] == 0 + + +def test_recover_without_ownership_rejects_before_writes(tmp_path: Path) -> None: + from tests.artifacts.test_run_store import artifact as _artifact + from tests.artifacts.test_run_store import deployment as _deployment + from wf_api.run_lifecycle import materialize_admitted_view, persist_admission + from wf_artifacts import PinnedRunEnvironment + + sched_store = FileScheduleStore(tmp_path / "sched") + run_store = FileRunStore(tmp_path / "runs") + sched_store.create_schedule(_sched_model("a")) + env = PinnedRunEnvironment( + deployment=_deployment(), root_artifact=_artifact(), child_artifacts=[] + ) + admission = persist_admission( + store=run_store, + run_id=run_store.allocate_run_id(), + environment=env, + resolved_input={}, + max_steps=None, + scheduled_at=ts(2026, 9, 8, 12, 0), + schedule_id="a", + schedule_revision=1, + ) + materialize_admitted_view(store=run_store, admission=admission) + ownership = SchedulerOwnership(tmp_path / "sched", owner="never-acquired") + with pytest.raises(SecondOwnerError): + sched_recovery.recover( + schedule_store=sched_store, + run_store=run_store, + now=ts(2026, 9, 8, 12, 0), + ownership=ownership, + ) + # Rejected before any write: the admitted run is untouched. + assert run_store.get_run(admission.id).status.value == "admitted" + + +def test_held_ownership_allows_poll_and_recovery(tmp_path: Path) -> None: + ownership = SchedulerOwnership(tmp_path / "sched", owner="test").acquire() + try: + sched, store, runs = _scheduler(tmp_path, ownership, script={"*": "hang"}) + intended = ts(2026, 9, 8, 12, 0) + _due_setup(sched, store, intended) + assert sched.poll(intended) == {"a": f"admit:{runs.list_runs()[0].id}"} + diags = sched_recovery.recover( + schedule_store=store, + run_store=runs, + now=intended, + ownership=ownership, + ) + assert diags != [] + finally: + ownership.release() + + +def test_second_owner_cannot_acquire_for_poll(tmp_path: Path) -> None: + first = SchedulerOwnership(tmp_path / "sched", owner="first").acquire() + try: + second = SchedulerOwnership(tmp_path / "sched", owner="second") + with pytest.raises(SecondOwnerError): + second.acquire() + assert not second.held + # A scheduler proving the loser's (unheld) ownership is rejected. + sched, store, _ = _scheduler(tmp_path, second) + intended = ts(2026, 9, 8, 12, 0) + _due_setup(sched, store, intended) + with pytest.raises(SecondOwnerError): + sched.poll(intended) + finally: + first.release() + + +def test_admin_mutations_require_ownership(tmp_path: Path) -> None: + ownership = SchedulerOwnership(tmp_path / "sched", owner="test") + sched, store, _ = _scheduler(tmp_path, ownership) + store.create_schedule(_sched_model("a", paused=True)) + with pytest.raises(SecondOwnerError): + sched.resume_schedule("a", ts(2026, 9, 8, 12, 0)) + with pytest.raises(SecondOwnerError): + sched.edit_schedule("a", ts(2026, 9, 8, 12, 0)) + assert store.get_schedule("a").paused is True + assert store.get_schedule("a").revision == 1 diff --git a/tests/scheduling/test_poll.py b/tests/scheduling/test_poll.py index 6ddb68cc..4d63ada8 100644 --- a/tests/scheduling/test_poll.py +++ b/tests/scheduling/test_poll.py @@ -15,6 +15,7 @@ from wf_artifacts.runs.models import StoredRunStatus from wf_artifacts.runs.store import FileRunStore from wf_scheduling.calendar import OneShotSource from wf_scheduling.models import Schedule +from wf_scheduling.ownership import SchedulerOwnership from wf_scheduling.poll import SCAN_CAP, Scheduler from wf_scheduling.prepare import SchedulePreparer from wf_scheduling.store import FileScheduleStore @@ -87,6 +88,7 @@ def _harness( capacity=capacity, preparer=preparer, dispatcher=ScriptedDispatcher(script), + ownership=SchedulerOwnership(tmp_path / "sched", owner="test").acquire(), ) return sched, sched_store, run_store, sources @@ -135,6 +137,7 @@ def test_overlap_skip_blocks_and_late_drops() -> None: sched.poll(t0 + timedelta(minutes=30)) kinds = [(r["kind"], r["resolved_at"]) for r in _history(store, "a")] assert any(k == "skipped-misfire" for k, _ in kinds) + sched.ownership.release() def test_latest_coalesces_to_one_candidate_and_no_double_admit() -> None: @@ -167,6 +170,7 @@ def test_latest_coalesces_to_one_candidate_and_no_double_admit() -> None: 2026, 9, 8, 13, 0 ) assert store.get_candidate("h") is None + sched.ownership.release() def test_parallel_limits_and_interrupted_slots() -> None: @@ -206,6 +210,7 @@ def test_parallel_limits_and_interrupted_slots() -> None: r["kind"] == "skipped-overlap" and "12:15" in str(r["resolved_at"]) for r in _history(store, "p") ) + sched.ownership.release() def test_pause_is_not_downtime() -> None: @@ -238,6 +243,7 @@ def test_pause_is_not_downtime() -> None: ] assert "2026-09-08T10:00:00+00:00" not in admitted assert "2026-09-08T11:00:00+00:00" not in admitted + sched.ownership.release() def test_long_downtime_is_bounded(tmp_path: Path) -> None: @@ -254,6 +260,7 @@ def test_long_downtime_is_bounded(tmp_path: Path) -> None: assert ( len([r for r in _history(store, "m") if r["kind"] == "interval-summary"]) == 1 ) + sched.ownership.release() def test_fairness_slow_schedule_not_starved(tmp_path: Path) -> None: @@ -277,6 +284,7 @@ def test_fairness_slow_schedule_not_starved(tmp_path: Path) -> None: runs.save_run(run.model_copy(update={"status": StoredRunStatus.COMPLETED})) sched.poll(t0 + timedelta(seconds=30)) assert any(r["kind"] == "admitted" for r in _history(store, "slow")) + sched.ownership.release() def test_capacity_wait_then_expire_for_skip(tmp_path: Path) -> None: @@ -289,3 +297,4 @@ def test_capacity_wait_then_expire_for_skip(tmp_path: Path) -> None: cast(ScriptedDispatcher, sched.dispatcher).script = {"*": "complete"} sched.poll(t0 + timedelta(seconds=30)) assert len([r for r in _history(store, "a") if r["kind"] == "admitted"]) == 1 + sched.ownership.release() diff --git a/tests/scheduling/test_preparation.py b/tests/scheduling/test_preparation.py index 6edbdc58..54735da9 100644 --- a/tests/scheduling/test_preparation.py +++ b/tests/scheduling/test_preparation.py @@ -23,6 +23,7 @@ from wf_artifacts.runs.store import FileRunStore from wf_scheduling.calendar import OneShotSource from wf_scheduling.dispatch import RunDispatcher from wf_scheduling.models import Schedule +from wf_scheduling.ownership import SchedulerOwnership from wf_scheduling.poll import Scheduler from wf_scheduling.prepare import ( InvocationPreparer, @@ -66,6 +67,7 @@ def _scheduler( capacity=4, preparer=preparer, dispatcher=ScriptedDispatcher(script), + ownership=SchedulerOwnership(tmp_path / "sched", owner="test").acquire(), ) return sched, sched_store, run_store @@ -80,6 +82,7 @@ def test_scheduler_requires_typed_collaborators() -> None: "capacity", "preparer", "dispatcher", + "ownership", } @@ -137,6 +140,7 @@ def test_occurrence_bindings_resolve_into_admission(tmp_path: Path) -> None: assert admission.resolved_input["which"] == "b" assert admission.deployment_revision == 3 assert admission.schedule_revision == 1 + sched.ownership.release() def test_unknown_deployment_rejects_without_a_run(tmp_path: Path) -> None: @@ -151,6 +155,7 @@ def test_unknown_deployment_rejects_without_a_run(tmp_path: Path) -> None: page = store.list_occurrences("gone", limit=100) kinds = [r["kind"] for r in cast(list[dict[str, Any]], page["occurrences"])] assert kinds == ["preflight-rejected"] + sched.ownership.release() def test_missing_required_input_rejects_without_a_run(tmp_path: Path) -> None: @@ -167,6 +172,7 @@ def test_missing_required_input_rejects_without_a_run(tmp_path: Path) -> None: capacity=4, preparer=preparer, dispatcher=ScriptedDispatcher({"*": "hang"}), + ownership=SchedulerOwnership(tmp_path / "sched", owner="test").acquire(), ) intended = ts(2026, 9, 8, 12, 0) sched_store.create_schedule(_sched_model("need")) @@ -178,6 +184,7 @@ def test_missing_required_input_rejects_without_a_run(tmp_path: Path) -> None: entry = cast(list[dict[str, Any]], page["occurrences"])[0] assert entry["kind"] == "preflight-rejected" assert "missing-input" in entry["reason"] + sched.ownership.release() def test_conflicting_schedule_targets_reject_without_a_run(tmp_path: Path) -> None: @@ -206,6 +213,7 @@ def test_conflicting_schedule_targets_reject_without_a_run(tmp_path: Path) -> No entry = cast(list[dict[str, Any]], page["occurrences"])[0] assert entry["kind"] == "preflight-rejected" assert entry["reason"].startswith("invalid-input:") + sched.ownership.release() def test_preparer_rejection_type_shape() -> None: diff --git a/tests/scheduling/test_recovery.py b/tests/scheduling/test_recovery.py index c54b7e21..8f4ee123 100644 --- a/tests/scheduling/test_recovery.py +++ b/tests/scheduling/test_recovery.py @@ -15,6 +15,7 @@ from wf_artifacts.runs.store import FileRunStore from wf_scheduling import recovery as sched_recovery from wf_scheduling.calendar import OneShotSource from wf_scheduling.models import Schedule +from wf_scheduling.ownership import SchedulerOwnership from wf_scheduling.poll import Scheduler from wf_scheduling.prepare import SchedulePreparer from wf_scheduling.store import FileScheduleStore @@ -57,9 +58,16 @@ def test_recovery_materializes_missing_view_as_pending(tmp_path: Path) -> None: resolved_input={}, max_steps=None, ) - diags = sched_recovery.recover( - schedule_store=sched_store, run_store=run_store, now=ts(2026, 9, 8, 12, 0) - ) + ownership = SchedulerOwnership(tmp_path / "sched", owner="test").acquire() + try: + diags = sched_recovery.recover( + schedule_store=sched_store, + run_store=run_store, + now=ts(2026, 9, 8, 12, 0), + ownership=ownership, + ) + finally: + ownership.release() assert any("pending-dispatch" in d for d in diags) assert run_store.get_run(admission.id).status.value == "admitted" assert sched_recovery._is_pending(run_store, admission.id) @@ -88,9 +96,16 @@ def test_recovery_fails_abandoned_admitted_without_replay(tmp_path: Path) -> Non schedule_revision=1, ) materialize_admitted_view(store=run_store, admission=admission) - diags = sched_recovery.recover( - schedule_store=sched_store, run_store=run_store, now=ts(2026, 9, 8, 12, 0) - ) + ownership = SchedulerOwnership(tmp_path / "sched", owner="test").acquire() + try: + diags = sched_recovery.recover( + schedule_store=sched_store, + run_store=run_store, + now=ts(2026, 9, 8, 12, 0), + ownership=ownership, + ) + finally: + ownership.release() assert any("failed-closed" in d for d in diags) assert run_store.get_run(admission.id).status.value == "failed" @@ -118,25 +133,37 @@ def test_recovery_never_executes_pending_until_poll(tmp_path: Path) -> None: schedule_id="a", schedule_revision=1, ) - diags = sched_recovery.recover( - schedule_store=sched_store, run_store=run_store, now=ts(2026, 9, 8, 12, 0) - ) + ownership = SchedulerOwnership(tmp_path / "sched", owner="test").acquire() + try: + diags = sched_recovery.recover( + schedule_store=sched_store, + run_store=run_store, + now=ts(2026, 9, 8, 12, 0), + ownership=ownership, + ) + finally: + ownership.release() assert any("pending-dispatch" in d for d in diags) # Recovery itself produced no terminal history; the poll sweep dispatches. assert sched_store.list_occurrences("a", limit=100)["total"] == 0 sources = {"a": OneShotSource(ts(2026, 9, 8, 12, 0))} - sched = Scheduler( - schedule_store=sched_store, - run_store=run_store, - sources=sources, # type: ignore[arg-type] - capacity=4, - preparer=SchedulePreparer( - DictDeployments({"dep-1": {"rev": 1, "required": []}}), - fixture_environment, - ), - dispatcher=ScriptedDispatcher({"*": "complete"}), - ) - sched_store.save_consumed("a", ts(2026, 9, 8, 12, 0)) - sched.poll(ts(2026, 9, 8, 12, 1)) + ownership = SchedulerOwnership(tmp_path / "sched", owner="test").acquire() + try: + sched = Scheduler( + schedule_store=sched_store, + run_store=run_store, + sources=sources, # type: ignore[arg-type] + capacity=4, + preparer=SchedulePreparer( + DictDeployments({"dep-1": {"rev": 1, "required": []}}), + fixture_environment, + ), + dispatcher=ScriptedDispatcher({"*": "complete"}), + ownership=ownership, + ) + sched_store.save_consumed("a", ts(2026, 9, 8, 12, 0)) + sched.poll(ts(2026, 9, 8, 12, 1)) + finally: + ownership.release() assert run_store.get_run(admission.id).status.value == "completed" assert not sched_recovery._is_pending(run_store, admission.id)