Assistant: codex Assistant-Model: gpt-6-astra Assistant-Session: 01a0e429-436c-73f1-9144-bac4336c06e8
140 lines
6.2 KiB
Python
140 lines
6.2 KiB
Python
"""Durable dispatch and audit recovery without production authority or transports."""
|
|
import pytest
|
|
from temporalio.exceptions import ActivityError, ApplicationError
|
|
from temporalio.worker.workflow_sandbox import SandboxedWorkflowRunner
|
|
from temporalio.workflow import _Definition
|
|
|
|
from activity_core.release_broker import Broker, TERMINAL
|
|
from activity_core.release_dispatch import ReleaseActivities, ReleaseDispatchWorkflow
|
|
from tests.test_release_broker import NOW, FakeBackend, admit, bundle # noqa: F401
|
|
|
|
|
|
@pytest.fixture
|
|
def dispatch(bundle, monkeypatch):
|
|
broker, *_ = bundle
|
|
rid = admit(bundle)
|
|
backend = FakeBackend()
|
|
advance = Broker.advance
|
|
monkeypatch.setattr(Broker, 'advance', lambda self, rid, backend: advance(self, rid, backend, NOW))
|
|
events = {}
|
|
sink = lambda event: events.update({event['event_id']: event})
|
|
activities = ReleaseActivities(broker, lambda plan: backend, sink)
|
|
return activities, rid, backend, events, sink
|
|
|
|
|
|
async def test_restart_delivers_all_committed_transitions_and_deduplicates(dispatch):
|
|
activities, rid, backend, events, sink = dispatch
|
|
assert await activities.advance(rid) == 'publish_pending'
|
|
assert await activities.advance(rid) == 'published'
|
|
# Restart before audit delivery. The ledger contains both undelivered steps.
|
|
restarted = ReleaseActivities(
|
|
Broker(activities.broker.database, activities.broker.receipts, admitted=True),
|
|
lambda plan: backend, sink,
|
|
)
|
|
assert await restarted.deliver(rid) == 3
|
|
assert await restarted.deliver(rid) == 0
|
|
while await restarted.advance(rid) not in TERMINAL:
|
|
pass
|
|
assert await restarted.deliver(rid) == 2
|
|
assert [e['phase'] for e in events.values()] == [
|
|
'planned', 'publish_pending', 'published', 'synced', 'complete',
|
|
]
|
|
assert all(set(e) == {'event_id', 'event_type', 'release_id', 'application',
|
|
'phase', 'observed_at', 'candidate', 'rollback'}
|
|
for e in events.values())
|
|
|
|
|
|
async def test_lost_audit_ack_replays_same_event_without_losing_it(dispatch):
|
|
activities, rid, _, events, sink = dispatch
|
|
def lost_ack(event):
|
|
sink(event)
|
|
raise RuntimeError('secret sink token must not reach Temporal')
|
|
activities.sink = lost_ack
|
|
with pytest.raises(ApplicationError) as error:
|
|
await activities.deliver(rid)
|
|
assert str(error.value) == 'ReleaseAuditUnavailable: release audit unavailable'
|
|
assert error.value.__suppress_context__
|
|
activities.sink = sink
|
|
assert await activities.deliver(rid) == 1
|
|
assert len(events) == 1
|
|
|
|
|
|
async def test_lost_publication_response_resumes_and_rolls_back(dispatch):
|
|
activities, rid, backend, events, _ = dispatch
|
|
backend.lose_response = True
|
|
await activities.advance(rid)
|
|
with pytest.raises(ApplicationError, match='release step unavailable'):
|
|
await activities.advance(rid)
|
|
backend.fail_health = True
|
|
for _ in range(10):
|
|
phase = await activities.advance(rid)
|
|
await activities.deliver(rid)
|
|
if phase in TERMINAL:
|
|
break
|
|
assert phase == 'rolled_back'
|
|
assert backend.revision == 'b' * 40
|
|
assert list(events.values())[-1]['phase'] == 'rolled_back'
|
|
|
|
|
|
@pytest.mark.parametrize('rid', ['untrusted', 'f' * 64, {'key': 'caller config'}])
|
|
async def test_only_existing_admitted_id_accepted(dispatch, rid):
|
|
activities, _, backend, events, _ = dispatch
|
|
with pytest.raises(ApplicationError) as error:
|
|
await activities.advance(rid)
|
|
assert error.value.non_retryable
|
|
assert not backend.calls and not events
|
|
|
|
|
|
async def test_revocation_blocks_resume_and_delivery(dispatch):
|
|
activities, rid, backend, events, _ = dispatch
|
|
activities.broker.admitted = False
|
|
for call in (activities.advance, activities.deliver):
|
|
with pytest.raises(ApplicationError) as error:
|
|
await call(rid)
|
|
assert error.value.non_retryable
|
|
assert not backend.calls and not events
|
|
|
|
|
|
async def test_temporal_sandbox_and_workflow_rollback(dispatch, monkeypatch):
|
|
# Validate actual Temporal sandbox imports, then exercise orchestration with
|
|
# real ledger activities. No Temporal server or simulated passage of soak time.
|
|
SandboxedWorkflowRunner().prepare_workflow(_Definition.must_from_class(ReleaseDispatchWorkflow))
|
|
activities, rid, backend, events, _ = dispatch
|
|
backend.fail_health = True
|
|
async def execute(name, arg, **kwargs):
|
|
return await (activities.advance(arg) if name == 'advance_admitted_release'
|
|
else activities.deliver(arg))
|
|
async def sleep(_):
|
|
pass
|
|
monkeypatch.setattr('activity_core.release_dispatch.workflow.execute_activity', execute)
|
|
monkeypatch.setattr('activity_core.release_dispatch.workflow.sleep', sleep)
|
|
assert await ReleaseDispatchWorkflow().run(rid) == 'rolled_back'
|
|
assert list(events.values())[-1]['phase'] == 'rolled_back'
|
|
|
|
|
|
async def test_sink_outage_cannot_prevent_rollback(dispatch, monkeypatch):
|
|
activities, rid, backend, events, _ = dispatch
|
|
backend.fail_health = True
|
|
phases = []
|
|
async def execute(name, arg, **kwargs):
|
|
if name == 'advance_admitted_release':
|
|
phase = await activities.advance(arg)
|
|
phases.append(phase)
|
|
return phase
|
|
if phases[-1] not in TERMINAL:
|
|
assert kwargs['retry_policy'].maximum_attempts == 1
|
|
raise ActivityError('audit unavailable', scheduled_event_id=1,
|
|
started_event_id=2, identity='fixture',
|
|
activity_type=name, activity_id='fixture', retry_state=None)
|
|
assert kwargs['retry_policy'].maximum_attempts == 0
|
|
return await activities.deliver(arg)
|
|
async def sleep(_):
|
|
pass
|
|
monkeypatch.setattr('activity_core.release_dispatch.workflow.execute_activity', execute)
|
|
monkeypatch.setattr('activity_core.release_dispatch.workflow.sleep', sleep)
|
|
# A plain logger avoids requiring a live workflow runtime for this unit test.
|
|
import logging
|
|
monkeypatch.setattr('activity_core.release_dispatch.workflow.logger', logging.getLogger(__name__))
|
|
assert await ReleaseDispatchWorkflow().run(rid) == 'rolled_back'
|
|
assert backend.revision == 'b' * 40
|
|
assert len(events) == len(phases) + 1 # Includes the admission transition.
|