"""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.