diff --git a/.github/scripts/windows_test_assignments.json b/.github/scripts/windows_test_assignments.json index ad64ac38f..61b1ed1e8 100644 --- a/.github/scripts/windows_test_assignments.json +++ b/.github/scripts/windows_test_assignments.json @@ -828,6 +828,7 @@ "tests/test_scheduler/test_parser.py", "tests/test_scheduler/test_parser_tz.py", "tests/test_scheduler/test_persistence.py", + "tests/test_scheduler/test_reservation_release_race.py", "tests/test_scheduler/test_rpc_webhook_delivery.py", "tests/test_scheduler/test_schedule_compute_notify.py", "tests/test_scheduler/test_schedule_normalizer.py", diff --git a/.github/scripts/windows_test_durations.json b/.github/scripts/windows_test_durations.json index 4a40c9f53..058bfb4a4 100644 --- a/.github/scripts/windows_test_durations.json +++ b/.github/scripts/windows_test_durations.json @@ -1159,6 +1159,7 @@ "tests/test_scheduler/test_parser.py": 0.033, "tests/test_scheduler/test_parser_tz.py": 0.608, "tests/test_scheduler/test_persistence.py": 1.236, + "tests/test_scheduler/test_reservation_release_race.py": 0.01, "tests/test_scheduler/test_rpc_webhook_delivery.py": 0.051, "tests/test_scheduler/test_schedule_compute_notify.py": 0.01, "tests/test_scheduler/test_schedule_normalizer.py": 0.017, diff --git a/src/opensquilla/scheduler/persistence.py b/src/opensquilla/scheduler/persistence.py index 0fb15a2bd..81bb4cdc3 100644 --- a/src/opensquilla/scheduler/persistence.py +++ b/src/opensquilla/scheduler/persistence.py @@ -830,15 +830,31 @@ async def release_reservation( job_id: str, reservation_token: str, ) -> bool: - current = await self.get(job_id) - if current is None or current.reservation_token != reservation_token: - return False - clear_reservation(current) - if current.status == JobStatus.RUNNING: - current.status = JobStatus.PENDING - current.updated_at = datetime.now(UTC) - await self.save(current) - return True + # Fence the cleanup in the write itself. A read followed by save() + # can overwrite a new owner or resurrect a concurrently deleted job. + async with self._db().execute( + """ + UPDATE scheduler_jobs + SET status = CASE WHEN status = ? THEN ? ELSE status END, + reservation_token = '', + reserved_at = NULL, + reserved_by = '', + reservation_source = '', + scheduled_run_at = NULL, + updated_at = MAX(updated_at, ?) + WHERE id = ? AND reservation_token = ? + """, + ( + JobStatus.RUNNING.value, + JobStatus.PENDING.value, + datetime.now(UTC).isoformat(), + job_id, + reservation_token, + ), + ) as cur: + released = cur.rowcount == 1 + await self._db().commit() + return released async def delete(self, job_id: str) -> None: await self._db().execute("DELETE FROM scheduler_jobs WHERE id = ?", (job_id,)) diff --git a/tests/test_scheduler/test_reservation_release_race.py b/tests/test_scheduler/test_reservation_release_race.py new file mode 100644 index 000000000..c293bace9 --- /dev/null +++ b/tests/test_scheduler/test_reservation_release_race.py @@ -0,0 +1,101 @@ +"""Cancellation cleanup must fence ownership at the database write boundary.""" + +from __future__ import annotations + +import asyncio +from datetime import UTC, datetime, timedelta +from pathlib import Path + +import pytest + +from opensquilla.scheduler.persistence import JobStore +from opensquilla.scheduler.types import CronJob, JobReservation, JobStatus + + +@pytest.mark.parametrize("concurrent_edit", ["delete", "replace_owner", "pause", "edit"]) +async def test_release_does_not_overwrite_concurrent_job_changes( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch, concurrent_edit: str +) -> None: + db_path = str(tmp_path / "scheduler.db") + store, writer = JobStore(db_path), JobStore(db_path) + await store.open() + await writer.open() + cleanup: asyncio.Task[bool] | None = None + cleanup_updated_at: datetime | None = None + concurrent_updated_at: datetime | None = None + try: + job = CronJob(id="job", name="original", cron_expr="*/5 * * * *", handler_key="test") + await store.save(job) + reservation = await store.reserve_manual_job(job.id, datetime.now(UTC)) + assert isinstance(reservation, JobReservation) + + # The competing writer owns the SQLite write lock. Cleanup may read an + # old snapshot, but its write cannot land until this writer commits. + await writer._db().execute("BEGIN IMMEDIATE") + write_started = asyncio.Event() + execute = store._db().execute + + def observe_write(sql, params=()): + nonlocal cleanup_updated_at + statement = sql.lstrip() + if statement.startswith("UPDATE") and "scheduler_jobs" in statement: + cleanup_updated_at = datetime.fromisoformat(tuple(params)[2]) + if statement.startswith(("INSERT", "UPDATE")) and "scheduler_jobs" in statement: + write_started.set() + return execute(sql, params) + + monkeypatch.setattr(store._db(), "execute", observe_write) + cleanup = asyncio.create_task(store.release_reservation(job.id, reservation.token)) + async with asyncio.timeout(2): + await write_started.wait() + + if concurrent_edit == "delete": + await writer.delete(job.id) + else: + changed = await writer.get(job.id) + assert changed is not None + assert cleanup_updated_at is not None + changed.name = "edited while cleanup waited" + changed.payload = {"message": "keep this new configuration"} + changed.updated_at = cleanup_updated_at + timedelta(seconds=1) + concurrent_updated_at = changed.updated_at + if concurrent_edit == "replace_owner": + changed.reservation_token = "new-owner" + changed.reserved_by = "replacement-worker" + elif concurrent_edit == "pause": + changed.status = JobStatus.PAUSED + changed.enabled = False + await writer.save(changed) + + released = await asyncio.wait_for(cleanup, timeout=2) + current = await writer.get(job.id) + if concurrent_edit == "delete": + assert current is None, "cleanup must not recreate a deleted job" + assert released is False + else: + assert current is not None + assert current.name == "edited while cleanup waited" + assert current.payload == {"message": "keep this new configuration"} + assert concurrent_updated_at is not None + assert current.updated_at == concurrent_updated_at + if concurrent_edit == "replace_owner": + assert released is False + assert current.reservation_token == "new-owner" + assert current.reserved_by == "replacement-worker" + assert current.status == JobStatus.RUNNING + else: + assert released is True + assert current.reservation_token == "" + assert current.reserved_at is None + assert current.reserved_by == "" + assert current.reservation_source == "" + assert current.scheduled_run_at is None + expected = JobStatus.PAUSED if concurrent_edit == "pause" else JobStatus.PENDING + assert current.status == expected + assert current.enabled is (concurrent_edit != "pause") + finally: + await writer._db().rollback() + if cleanup is not None: + await asyncio.gather(cleanup, return_exceptions=True) + await writer.close() + await store.close()