From 1042884faa28d7e545301cef5c085be54bb513fd Mon Sep 17 00:00:00 2001 From: lda Date: Thu, 10 Sep 2026 02:16:08 +0700 Subject: [PATCH] fix: make schedule mutations synchronous --- src/wf_api/schedules.py | 81 +++++++++++++++++++++++++++++++++++++---- 1 file changed, 73 insertions(+), 8 deletions(-) diff --git a/src/wf_api/schedules.py b/src/wf_api/schedules.py index 13118d4f..899f3066 100644 --- a/src/wf_api/schedules.py +++ b/src/wf_api/schedules.py @@ -49,13 +49,13 @@ _PENDING_PAGE_CURSOR = "__pending__" def _serialize_schedule_write(method: Any) -> Any: - """Serialize one API schedule mutation with the store's local lock.""" + """Run one synchronous schedule mutation under the store transaction.""" @wraps(method) - async def wrapped(self: Any, *args: Any, **kwargs: Any) -> Any: + def wrapped(self: Any, *args: Any, **kwargs: Any) -> Any: store = self._schedule_store() with schedule_store_transaction(store): - return await method(self, *args, **kwargs) + return method(self, *args, **kwargs) return wrapped @@ -218,7 +218,6 @@ class WorkflowScheduleApi: self._check_trigger(schedule) self._validate_sample(schedule) - @_serialize_schedule_write async def create_schedule( self, *, @@ -240,6 +239,34 @@ class WorkflowScheduleApi: ``ValueError``; duplicate ids (deleted ids are never reusable) raise ``ScheduleExistsError``. """ + return self._create_schedule( + schedule_id=schedule_id, + deployment_id=deployment_id, + trigger=trigger, + input_bindings=input_bindings, + overlap=overlap, + misfire=misfire, + max_active_runs=max_active_runs, + lateness_allowance_s=lateness_allowance_s, + max_steps=max_steps, + enabled=enabled, + ) + + @_serialize_schedule_write + def _create_schedule( + self, + *, + schedule_id: str, + deployment_id: str, + trigger: dict[str, Any], + input_bindings: list[dict[str, Any]] | None = None, + overlap: str = "skip", + misfire: str = "skip", + max_active_runs: int = 1, + lateness_allowance_s: float = 60.0, + max_steps: int | None = None, + enabled: bool = True, + ) -> ScheduleResult: store = self._schedule_store() now = datetime.now(UTC) # KeyError first: the deployment must exist before any validation. @@ -303,7 +330,6 @@ class WorkflowScheduleApi: } ) - @_serialize_schedule_write async def update_schedule( self, *, @@ -335,6 +361,36 @@ class WorkflowScheduleApi: wins via the store's authoritative revision check; our already applied watermark advance is a safe over-skip in that case too. """ + return self._update_schedule( + schedule_id=schedule_id, + expected_revision=expected_revision, + deployment_id=deployment_id, + trigger=trigger, + input_bindings=input_bindings, + overlap=overlap, + misfire=misfire, + max_active_runs=max_active_runs, + lateness_allowance_s=lateness_allowance_s, + max_steps=max_steps, + enabled=enabled, + ) + + @_serialize_schedule_write + def _update_schedule( + self, + *, + schedule_id: str, + expected_revision: int, + deployment_id: str | None = None, + trigger: dict[str, Any] | None = None, + input_bindings: list[dict[str, Any]] | None = None, + overlap: str | None = None, + misfire: str | None = None, + max_active_runs: int | None = None, + lateness_allowance_s: float | None = None, + max_steps: int | None = None, + enabled: bool | None = None, + ) -> ScheduleResult: store = self._schedule_store() now = datetime.now(UTC) current = store.get_schedule(schedule_id) @@ -394,7 +450,6 @@ class WorkflowScheduleApi: stored = store.update_schedule(updated, expected_revision=expected_revision) return _PROJECT_SCHEDULE(stored.model_dump(mode="json")) - @_serialize_schedule_write async def pause_schedule(self, *, schedule_id: str) -> ScheduleResult: """Pause one schedule (mirror the poll-loop paused branch). @@ -406,6 +461,10 @@ class WorkflowScheduleApi: schedule unpaused (retryable) and never a paused flag whose span backfills on resume. """ + return self._pause_schedule(schedule_id=schedule_id) + + @_serialize_schedule_write + def _pause_schedule(self, *, schedule_id: str) -> ScheduleResult: store = self._schedule_store() now = datetime.now(UTC) # Existence first: unknown ids raise KeyError before any write. @@ -420,7 +479,6 @@ class WorkflowScheduleApi: store.save_schedule(schedule) return _PROJECT_SCHEDULE(schedule.model_dump(mode="json")) - @_serialize_schedule_write async def resume_schedule(self, *, schedule_id: str) -> ScheduleResult: """Resume one schedule (mirror ``Scheduler.resume_schedule``). @@ -429,6 +487,10 @@ class WorkflowScheduleApi: at least now. No history row is written. Crash-safe ordering like pause: candidate and watermark first, flag flip last. """ + return self._resume_schedule(schedule_id=schedule_id) + + @_serialize_schedule_write + def _resume_schedule(self, *, schedule_id: str) -> ScheduleResult: store = self._schedule_store() now = datetime.now(UTC) # Existence first: unknown ids raise KeyError before any write. @@ -443,7 +505,6 @@ class WorkflowScheduleApi: store.save_schedule(schedule) return _PROJECT_SCHEDULE(schedule.model_dump(mode="json")) - @_serialize_schedule_write async def delete_schedule(self, *, schedule_id: str) -> ScheduleResult: """Soft-delete one schedule (mirror the poll-loop deleted branch). @@ -454,6 +515,10 @@ class WorkflowScheduleApi: dropped by the poll loop, and a cleared candidate on a live schedule is rebuilt from the untouched watermark. """ + return self._delete_schedule(schedule_id=schedule_id) + + @_serialize_schedule_write + def _delete_schedule(self, *, schedule_id: str) -> ScheduleResult: store = self._schedule_store() # Existence first: unknown ids raise KeyError before any write. store.get_schedule(schedule_id)