fix: make schedule mutations synchronous
This commit is contained in:
+73
-8
@@ -49,13 +49,13 @@ _PENDING_PAGE_CURSOR = "__pending__"
|
|||||||
|
|
||||||
|
|
||||||
def _serialize_schedule_write(method: Any) -> Any:
|
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)
|
@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()
|
store = self._schedule_store()
|
||||||
with schedule_store_transaction(store):
|
with schedule_store_transaction(store):
|
||||||
return await method(self, *args, **kwargs)
|
return method(self, *args, **kwargs)
|
||||||
|
|
||||||
return wrapped
|
return wrapped
|
||||||
|
|
||||||
@@ -218,7 +218,6 @@ class WorkflowScheduleApi:
|
|||||||
self._check_trigger(schedule)
|
self._check_trigger(schedule)
|
||||||
self._validate_sample(schedule)
|
self._validate_sample(schedule)
|
||||||
|
|
||||||
@_serialize_schedule_write
|
|
||||||
async def create_schedule(
|
async def create_schedule(
|
||||||
self,
|
self,
|
||||||
*,
|
*,
|
||||||
@@ -240,6 +239,34 @@ class WorkflowScheduleApi:
|
|||||||
``ValueError``; duplicate ids (deleted ids are never reusable)
|
``ValueError``; duplicate ids (deleted ids are never reusable)
|
||||||
raise ``ScheduleExistsError``.
|
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()
|
store = self._schedule_store()
|
||||||
now = datetime.now(UTC)
|
now = datetime.now(UTC)
|
||||||
# KeyError first: the deployment must exist before any validation.
|
# KeyError first: the deployment must exist before any validation.
|
||||||
@@ -303,7 +330,6 @@ class WorkflowScheduleApi:
|
|||||||
}
|
}
|
||||||
)
|
)
|
||||||
|
|
||||||
@_serialize_schedule_write
|
|
||||||
async def update_schedule(
|
async def update_schedule(
|
||||||
self,
|
self,
|
||||||
*,
|
*,
|
||||||
@@ -335,6 +361,36 @@ class WorkflowScheduleApi:
|
|||||||
wins via the store's authoritative revision check; our already
|
wins via the store's authoritative revision check; our already
|
||||||
applied watermark advance is a safe over-skip in that case too.
|
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()
|
store = self._schedule_store()
|
||||||
now = datetime.now(UTC)
|
now = datetime.now(UTC)
|
||||||
current = store.get_schedule(schedule_id)
|
current = store.get_schedule(schedule_id)
|
||||||
@@ -394,7 +450,6 @@ class WorkflowScheduleApi:
|
|||||||
stored = store.update_schedule(updated, expected_revision=expected_revision)
|
stored = store.update_schedule(updated, expected_revision=expected_revision)
|
||||||
return _PROJECT_SCHEDULE(stored.model_dump(mode="json"))
|
return _PROJECT_SCHEDULE(stored.model_dump(mode="json"))
|
||||||
|
|
||||||
@_serialize_schedule_write
|
|
||||||
async def pause_schedule(self, *, schedule_id: str) -> ScheduleResult:
|
async def pause_schedule(self, *, schedule_id: str) -> ScheduleResult:
|
||||||
"""Pause one schedule (mirror the poll-loop paused branch).
|
"""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
|
schedule unpaused (retryable) and never a paused flag whose span
|
||||||
backfills on resume.
|
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()
|
store = self._schedule_store()
|
||||||
now = datetime.now(UTC)
|
now = datetime.now(UTC)
|
||||||
# Existence first: unknown ids raise KeyError before any write.
|
# Existence first: unknown ids raise KeyError before any write.
|
||||||
@@ -420,7 +479,6 @@ class WorkflowScheduleApi:
|
|||||||
store.save_schedule(schedule)
|
store.save_schedule(schedule)
|
||||||
return _PROJECT_SCHEDULE(schedule.model_dump(mode="json"))
|
return _PROJECT_SCHEDULE(schedule.model_dump(mode="json"))
|
||||||
|
|
||||||
@_serialize_schedule_write
|
|
||||||
async def resume_schedule(self, *, schedule_id: str) -> ScheduleResult:
|
async def resume_schedule(self, *, schedule_id: str) -> ScheduleResult:
|
||||||
"""Resume one schedule (mirror ``Scheduler.resume_schedule``).
|
"""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
|
at least now. No history row is written. Crash-safe ordering like
|
||||||
pause: candidate and watermark first, flag flip last.
|
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()
|
store = self._schedule_store()
|
||||||
now = datetime.now(UTC)
|
now = datetime.now(UTC)
|
||||||
# Existence first: unknown ids raise KeyError before any write.
|
# Existence first: unknown ids raise KeyError before any write.
|
||||||
@@ -443,7 +505,6 @@ class WorkflowScheduleApi:
|
|||||||
store.save_schedule(schedule)
|
store.save_schedule(schedule)
|
||||||
return _PROJECT_SCHEDULE(schedule.model_dump(mode="json"))
|
return _PROJECT_SCHEDULE(schedule.model_dump(mode="json"))
|
||||||
|
|
||||||
@_serialize_schedule_write
|
|
||||||
async def delete_schedule(self, *, schedule_id: str) -> ScheduleResult:
|
async def delete_schedule(self, *, schedule_id: str) -> ScheduleResult:
|
||||||
"""Soft-delete one schedule (mirror the poll-loop deleted branch).
|
"""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
|
dropped by the poll loop, and a cleared candidate on a live
|
||||||
schedule is rebuilt from the untouched watermark.
|
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()
|
store = self._schedule_store()
|
||||||
# Existence first: unknown ids raise KeyError before any write.
|
# Existence first: unknown ids raise KeyError before any write.
|
||||||
store.get_schedule(schedule_id)
|
store.get_schedule(schedule_id)
|
||||||
|
|||||||
Reference in New Issue
Block a user