-
Notifications
You must be signed in to change notification settings - Fork 557
Expand file tree
/
Copy pathtest_reservation_release_race.py
More file actions
90 lines (79 loc) · 3.83 KB
/
Copy pathtest_reservation_release_race.py
File metadata and controls
90 lines (79 loc) · 3.83 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
"""Cancellation cleanup must fence ownership at the database write boundary."""
from __future__ import annotations
import asyncio
from datetime import UTC, datetime
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
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=()):
if sql.lstrip().startswith(("INSERT", "UPDATE")) and "scheduler_jobs" in sql:
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
changed.name = "edited while cleanup waited"
changed.payload = {"message": "keep this new configuration"}
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"}
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()