fix: reconcile superseded schedule reservations

This commit is contained in:
NotoriousRebel
2026-08-20 10:02:40 -04:00
parent a1d794d5af
commit ca1cb72ec6
4 changed files with 181 additions and 12 deletions
+23
View File
@@ -110,6 +110,29 @@ def test_schedule_page_creates_and_manages_two_target_passive_schedule(
expect(page.locator('#schedule-empty')).to_be_visible()
def test_schedule_page_rejects_a_non_loopback_host_in_browser(
harvestview_server_url: str,
page: Page,
browser_failures,
) -> None:
browser_failures.allow_console_error('Failed to load resource: the server responded with a status of 403 (Forbidden)')
remote_url = harvestview_server_url.replace('127.0.0.1', 'attacker.example') + '/schedules'
def forward_non_loopback_host(route: Route) -> None:
rejected = page.context.request.get(
f'{harvestview_server_url}/schedules',
headers={'Host': 'attacker.example'},
)
route.fulfill(status=rejected.status, headers=rejected.headers, body=rejected.body())
page.route(remote_url, forward_non_loopback_host)
response = page.goto(remote_url)
assert response is not None
assert response.status == 403
expect(page.locator('body')).to_contain_text('theHarvester is available only on localhost')
def test_schedule_clone_requires_fresh_authorization_for_active_work(
harvestview_server_url: str,
page: Page,
+128 -2
View File
@@ -372,7 +372,8 @@ def test_dispatch_recovery_reuses_reservations_and_run_ids(tmp_path, monkeypatch
asyncio.run(scenario())
def test_skip_policy_recovers_current_occurrence_reservations(tmp_path, monkeypatch) -> None:
@pytest.mark.parametrize('overlap_policy', ['skip', 'queue'])
def test_overlap_policy_recovers_current_occurrence_reservations(tmp_path, monkeypatch, overlap_policy: str) -> None:
from theHarvester.lib.api import run_worker
from theHarvester.lib.api.run_models import RunRequest
from theHarvester.lib.api.run_store import RunStore
@@ -390,7 +391,9 @@ def test_skip_policy_recovers_current_occurrence_reservations(tmp_path, monkeypa
store = ScheduleStore()
run_store = RunStore()
scheduled_for = datetime(2026, 8, 20, 9, tzinfo=UTC)
schedule = await store.create(ScheduleCreate.model_validate(_payload(start_at=scheduled_for.isoformat())))
payload = _payload(start_at=scheduled_for.isoformat())
payload['overlap_policy'] = overlap_policy
schedule = await store.create(ScheduleCreate.model_validate(payload))
first_run_id = '55555555-5555-4555-8555-555555555555'
second_run_id = '66666666-6666-4666-8666-666666666666'
await store.reserve_dispatch(schedule.schedule_id, scheduled_for, 'example.test', first_run_id)
@@ -450,6 +453,129 @@ def test_recovered_queued_occurrence_completes_without_a_false_error(tmp_path, m
asyncio.run(scenario())
@pytest.mark.parametrize('transition', ['pause-resume', 'replace'])
@pytest.mark.parametrize('overlap_policy', ['skip', 'queue'])
def test_newer_occurrence_terminalizes_a_stale_runless_reservation(
tmp_path,
monkeypatch,
transition: str,
overlap_policy: str,
) -> None:
from theHarvester.lib.api import run_worker
from theHarvester.lib.api.run_store import RunStore
from theHarvester.lib.api.schedule_models import ScheduleCreate, parse_utc
from theHarvester.lib.api.schedule_service import _dispatch_claimed
from theHarvester.lib.api.schedule_store import ScheduleStore
monkeypatch.setenv('THEHARVESTER_RUN_DB', str(tmp_path / 'runs.sqlite'))
monkeypatch.setenv('THEHARVESTER_SCHEDULE_DB', str(tmp_path / 'schedules.sqlite'))
monkeypatch.setattr(run_worker, 'worker_enabled', lambda: True)
monkeypatch.setattr(run_worker, 'worker_available', lambda: True)
monkeypatch.setattr(run_worker, 'wake_worker', lambda: None)
async def scenario() -> None:
store = ScheduleStore()
first = datetime(2020, 1, 1, tzinfo=UTC)
stale = datetime(2020, 1, 1, 1, tzinfo=UTC)
payload = _payload(start_at=first.isoformat(), targets=['example.test'])
payload['timing']['frequency'] = 'hourly'
payload['overlap_policy'] = overlap_policy
schedule = await store.create(ScheduleCreate.model_validate(payload))
assert (await store.claim_due('first-owner', now=first))[0].schedule_id == schedule.schedule_id
assert await store.complete_claim(
schedule.schedule_id,
'first-owner',
scheduled_for=first,
next_run_at=stale,
)
stale_run_id = '88888888-8888-4888-8888-888888888888'
await store.reserve_dispatch(schedule.schedule_id, stale, 'example.test', stale_run_id)
if transition == 'pause-resume':
assert await store.set_enabled(schedule.schedule_id, False) is not None
changed = await store.set_enabled(schedule.schedule_id, True)
else:
replacement = _payload(start_at='2030-01-01T09:00:00+00:00', targets=['example.test'])
replacement['timing']['frequency'] = 'hourly'
replacement['overlap_policy'] = overlap_policy
changed = await store.replace(schedule.schedule_id, ScheduleCreate.model_validate(replacement))
assert changed is not None
assert changed.next_run_at is not None
current = parse_utc(changed.next_run_at)
claimed = await store.claim_due('current-owner', now=current)
await _dispatch_claimed(claimed[0], 'current-owner')
dispatches = await store.list_dispatches(schedule.schedule_id)
stale_dispatch = next(dispatch for dispatch in dispatches if dispatch.run_id == stale_run_id)
assert stale_dispatch.state == 'failed'
assert stale_dispatch.error == 'Reserved occurrence was superseded before run creation'
current_dispatches = [dispatch for dispatch in dispatches if dispatch.scheduled_for == changed.next_run_at]
assert len(current_dispatches) == 1
assert current_dispatches[0].state == 'queued'
assert [run['run_id'] for run in await RunStore().list_runs(limit=10)] == [current_dispatches[0].run_id]
asyncio.run(scenario())
def test_large_reservation_reconciliation_renews_the_schedule_claim(tmp_path, monkeypatch) -> None:
from theHarvester.lib.api import run_worker
from theHarvester.lib.api.schedule_models import ScheduleCreate, parse_utc
from theHarvester.lib.api.schedule_service import _dispatch_claimed
from theHarvester.lib.api.schedule_store import ScheduleStore
monkeypatch.setenv('THEHARVESTER_RUN_DB', str(tmp_path / 'runs.sqlite'))
monkeypatch.setenv('THEHARVESTER_SCHEDULE_DB', str(tmp_path / 'schedules.sqlite'))
monkeypatch.setattr(run_worker, 'worker_enabled', lambda: True)
monkeypatch.setattr(run_worker, 'worker_available', lambda: True)
monkeypatch.setattr(run_worker, 'wake_worker', lambda: None)
renewals = 0
renew_claim = ScheduleStore.renew_claim
async def track_renewal(self, schedule_id: str, owner_id: str, *, lease_seconds: int = 60) -> bool:
nonlocal renewals
renewals += 1
return await renew_claim(self, schedule_id, owner_id, lease_seconds=lease_seconds)
monkeypatch.setattr(ScheduleStore, 'renew_claim', track_renewal)
async def scenario() -> None:
store = ScheduleStore()
first = datetime(2020, 1, 1, tzinfo=UTC)
stale = datetime(2020, 1, 1, 1, tzinfo=UTC)
payload = _payload(start_at=first.isoformat(), targets=['example.test'])
payload['timing']['frequency'] = 'hourly'
payload['overlap_policy'] = 'queue'
schedule = await store.create(ScheduleCreate.model_validate(payload))
assert (await store.claim_due('first-owner', now=first))[0].schedule_id == schedule.schedule_id
assert await store.complete_claim(
schedule.schedule_id,
'first-owner',
scheduled_for=first,
next_run_at=stale,
)
for index in range(100):
await store.reserve_dispatch(
schedule.schedule_id,
stale,
f'example-{index}.test',
f'00000000-0000-4000-8000-{index:012d}',
)
assert await store.set_enabled(schedule.schedule_id, False) is not None
changed = await store.set_enabled(schedule.schedule_id, True)
assert changed is not None and changed.next_run_at is not None
current = parse_utc(changed.next_run_at)
claimed = await store.claim_due('current-owner', now=current)
await _dispatch_claimed(claimed[0], 'current-owner')
dispatches = await store.list_dispatches(schedule.schedule_id)
assert renewals == 1
assert sum(dispatch.state == 'failed' for dispatch in dispatches) == 100
assert sum(dispatch.state == 'queued' for dispatch in dispatches) == 1
asyncio.run(scenario())
def test_deferred_claim_keeps_occurrence_identity(tmp_path, monkeypatch) -> None:
from theHarvester.lib.api.schedule_models import ScheduleCreate, parse_utc
from theHarvester.lib.api.schedule_store import ScheduleStore
+22 -6
View File
@@ -33,20 +33,33 @@ def wake_scheduler() -> None:
_scheduler_wakeup.set()
async def _refresh_pending_dispatches(
async def _reconcile_pending_dispatches(
schedule_id: str,
schedule_store: ScheduleStore,
run_store: RunStore,
*,
current_occurrence: datetime | None = None,
reserved_only: bool = False,
claim_owner: str | None = None,
) -> bool:
active = False
for dispatch in await schedule_store.pending_dispatches(schedule_id):
dispatches = await schedule_store.pending_dispatches(schedule_id, reserved_only=reserved_only)
for index, dispatch in enumerate(dispatches, start=1):
if index % 100 == 0:
if claim_owner is not None and not await schedule_store.renew_claim(schedule_id, claim_owner):
raise RuntimeError('Scheduler lost its occurrence claim while reconciling scheduled runs')
await asyncio.sleep(0)
run = await run_store.get(dispatch.run_id)
if run is None:
if dispatch.state == 'reserved':
if current_occurrence is None or parse_utc(dispatch.scheduled_for) != current_occurrence:
if current_occurrence is None:
active = True
elif parse_utc(dispatch.scheduled_for) != current_occurrence:
await schedule_store.set_dispatch_state(
dispatch.run_id,
'failed',
'Reserved occurrence was superseded before run creation',
)
continue
await schedule_store.set_dispatch_state(dispatch.run_id, 'failed', 'Scheduled run record is missing')
continue
@@ -129,7 +142,7 @@ async def _enqueue_targets(
async def refresh_schedule_dispatches(schedule_id: str) -> None:
await _refresh_pending_dispatches(schedule_id, ScheduleStore(), RunStore())
await _reconcile_pending_dispatches(schedule_id, ScheduleStore(), RunStore())
async def dispatch_schedule_now(schedule_id: str) -> ScheduleDispatchResponse | None:
@@ -170,12 +183,15 @@ async def _dispatch_claimed(schedule: ScheduleResponse, owner_id: str) -> None:
)
return
if schedule.overlap_policy == 'skip' and await _refresh_pending_dispatches(
active = await _reconcile_pending_dispatches(
schedule.schedule_id,
schedule_store,
run_store,
current_occurrence=scheduled_for,
):
reserved_only=schedule.overlap_policy == 'queue',
claim_owner=owner_id,
)
if schedule.overlap_policy == 'skip' and active:
next_run = schedule.timing.next_future_after(scheduled_for)
await schedule_store.complete_claim(
schedule.schedule_id,
+8 -4
View File
@@ -444,14 +444,18 @@ class ScheduleStore:
finally:
await connection.close()
async def pending_dispatches(self, schedule_id: str) -> list[ScheduleDispatchRecord]:
async def pending_dispatches(
self,
schedule_id: str,
*,
reserved_only: bool = False,
) -> list[ScheduleDispatchRecord]:
await self.initialize()
connection = await self._connect()
try:
states = "state = 'reserved'" if reserved_only else "state IN ('reserved', 'queued', 'running', 'cancelling')"
cursor = await connection.execute(
'SELECT * FROM schedule_dispatches WHERE schedule_id = ? '
"AND state IN ('reserved', 'queued', 'running', 'cancelling') "
'ORDER BY scheduled_for, target',
f'SELECT * FROM schedule_dispatches WHERE schedule_id = ? AND {states} ORDER BY scheduled_for, target',
(schedule_id,),
)
return [ScheduleDispatchRecord.model_validate(dict(row)) for row in await cursor.fetchall()]