import copy import pytest from app.core.orchestration.migration.cutover import ( CutoverGates, EngineRoleController, WorkflowCutoverService, ) from app.core.orchestration.migration.dual_run import DualRunMode from app.core.orchestration.migration.retirement import ( N8nRetirementChecklist, N8nRetirementGate, ) class Store: def __init__(self): self.mode = DualRunMode.N8N_PRIMARY_KESTRA_SHADOW self.operations = {} self.events = [] self.completed = [] self.rolled_back = [] def get_state(self, dataflow_uid, environment): return { "dataflow_uid": dataflow_uid, "environment": environment, "mode": self.mode.value, "n8n_definition_id": "n8n-orders", "kestra_definition_id": "kestra-orders", "status": "stable", } def request_transition(self, record): previous = self.operations.get(record["idempotency_key"]) if previous: return False, copy.deepcopy(previous) operation = { **copy.deepcopy(record), "operation_id": "cutover-operation-1", "status": "claimed", } self.operations[record["idempotency_key"]] = operation self.events.append( { "event_type": "workflow.engine.cutover.requested", "payload": copy.deepcopy(operation), } ) return True, copy.deepcopy(operation) def complete_transition(self, operation_id, target_mode, result): self.mode = DualRunMode(target_mode) self.completed.append((operation_id, target_mode, copy.deepcopy(result))) def rollback_transition(self, operation_id, source_mode, result): self.mode = DualRunMode(source_mode) self.rolled_back.append((operation_id, source_mode, copy.deepcopy(result))) class Engines: def __init__(self): self.calls = [] self.fail_kestra_activate = False def deactivate_n8n(self, definition_id): self.calls.append(("deactivate_n8n", definition_id)) return {"status": "disabled"} def activate_n8n(self, definition_id): self.calls.append(("activate_n8n", definition_id)) return {"status": "enabled"} def deactivate_kestra(self, definition_id): self.calls.append(("deactivate_kestra", definition_id)) return {"status": "disabled"} def activate_kestra(self, definition_id): self.calls.append(("activate_kestra", definition_id)) if self.fail_kestra_activate: raise RuntimeError("Kestra unavailable") return {"status": "enabled"} class Adapter: def __init__(self): self.calls = [] def activate(self, namespace, definition_id): self.calls.append(("activate", namespace, definition_id)) return {"status": "enabled"} def deactivate(self, namespace, definition_id): self.calls.append(("deactivate", namespace, definition_id)) return {"status": "disabled"} def test_engine_role_controller_binds_fixed_namespaces(): n8n = Adapter() kestra = Adapter() controller = EngineRoleController( n8n=n8n, kestra=kestra, n8n_namespace="dataops", kestra_namespace="dataops.migration", ) controller.deactivate_n8n("n8n-orders") controller.activate_kestra("kestra-orders") assert n8n.calls == [("deactivate", "dataops", "n8n-orders")] assert kestra.calls == [ ("activate", "dataops.migration", "kestra-orders") ] def passing_gates(): return CutoverGates( inventory_complete=True, workflow_spec_valid=True, rollback_target_present=True, shadow_runs=3, required_shadow_runs=3, failure_recovery_observed=True, reconciliation_passed=True, connection_budget_passed=True, sla_passed=True, n8n_baseline_unchanged=True, ) def test_cutover_request_is_transactional_outbox_only(): store = Store() engines = Engines() service = WorkflowCutoverService(store=store, engines=engines) result = service.request_cutover( dataflow_uid="01900000-0000-7000-8000-000000000094", environment="test", target_mode=DualRunMode.KESTRA_PRIMARY_N8N_STANDBY, gates=passing_gates(), idempotency_key="cutover:orders:v55", correlation_id="01900000-0000-7000-8000-000000000095", ) assert result["status"] == "claimed" assert result["source_mode"] == "n8n_primary_kestra_shadow" assert result["target_mode"] == "kestra_primary_n8n_standby" assert engines.calls == [] assert store.events[0]["event_type"] == "workflow.engine.cutover.requested" assert "password" not in repr(store.events[0]).lower() def test_cutover_consumer_disables_n8n_before_enabling_kestra(): store = Store() engines = Engines() service = WorkflowCutoverService(store=store, engines=engines) operation = service.request_cutover( dataflow_uid="01900000-0000-7000-8000-000000000094", environment="test", target_mode=DualRunMode.KESTRA_PRIMARY_N8N_STANDBY, gates=passing_gates(), idempotency_key="cutover:orders:v55", correlation_id="01900000-0000-7000-8000-000000000095", ) result = service.apply_requested_cutover(operation) assert result["status"] == "succeeded" assert result["mode"] == "kestra_primary_n8n_standby" assert engines.calls == [ ("deactivate_n8n", "n8n-orders"), ("activate_kestra", "kestra-orders"), ] assert store.completed[0][0] == "cutover-operation-1" def test_kestra_activation_failure_automatically_restores_n8n(): store = Store() engines = Engines() engines.fail_kestra_activate = True service = WorkflowCutoverService(store=store, engines=engines) operation = service.request_cutover( dataflow_uid="01900000-0000-7000-8000-000000000094", environment="test", target_mode=DualRunMode.KESTRA_PRIMARY_N8N_STANDBY, gates=passing_gates(), idempotency_key="cutover:orders:v55", correlation_id="01900000-0000-7000-8000-000000000095", ) result = service.apply_requested_cutover(operation) assert result["status"] == "rolled_back" assert result["mode"] == "n8n_primary_kestra_shadow" assert engines.calls == [ ("deactivate_n8n", "n8n-orders"), ("activate_kestra", "kestra-orders"), ("activate_n8n", "n8n-orders"), ] assert store.rolled_back[0][1] == "n8n_primary_kestra_shadow" def test_same_cutover_key_is_idempotent(): store = Store() service = WorkflowCutoverService(store=store, engines=Engines()) arguments = { "dataflow_uid": "01900000-0000-7000-8000-000000000094", "environment": "test", "target_mode": DualRunMode.KESTRA_PRIMARY_N8N_STANDBY, "gates": passing_gates(), "idempotency_key": "cutover:orders:v55", "correlation_id": "01900000-0000-7000-8000-000000000095", } first = service.request_cutover(**arguments) second = service.request_cutover(**arguments) assert first == second assert len(store.events) == 1 def test_cutover_gates_fail_closed(): failed = CutoverGates( **{ **passing_gates().as_dict(), "reconciliation_passed": False, } ) with pytest.raises(ValueError, match="reconciliation_passed"): WorkflowCutoverService(store=Store(), engines=Engines()).request_cutover( dataflow_uid="01900000-0000-7000-8000-000000000094", environment="test", target_mode=DualRunMode.KESTRA_PRIMARY_N8N_STANDBY, gates=failed, idempotency_key="cutover:orders:v55", correlation_id="01900000-0000-7000-8000-000000000095", ) def test_kestra_primary_requires_independent_retirement_approval(): store = Store() store.mode = DualRunMode.KESTRA_PRIMARY_N8N_STANDBY service = WorkflowCutoverService(store=store, engines=Engines()) with pytest.raises(ValueError, match="retirement approval"): service.request_cutover( dataflow_uid="01900000-0000-7000-8000-000000000094", environment="test", target_mode=DualRunMode.KESTRA_PRIMARY, gates=passing_gates(), idempotency_key="retire:orders:v55", correlation_id="01900000-0000-7000-8000-000000000095", ) def test_approved_retirement_can_finalize_kestra_primary_mode(): store = Store() store.mode = DualRunMode.KESTRA_PRIMARY_N8N_STANDBY engines = Engines() service = WorkflowCutoverService(store=store, engines=engines) decision = N8nRetirementGate.evaluate( N8nRetirementChecklist( inventory_disposition_complete=True, all_dataflows_kestra_primary=True, business_cycle_and_recovery_complete=True, standby_observation_no_rollback=True, monitoring_and_alerting_complete=True, dependency_failures_safe=True, secure_archive_complete=True, independent_review_approved=True, final_l4_passed=True, ) ) operation = service.request_cutover( dataflow_uid="01900000-0000-7000-8000-000000000094", environment="test", target_mode=DualRunMode.KESTRA_PRIMARY, gates=passing_gates(), idempotency_key="retire:orders:v55", correlation_id="01900000-0000-7000-8000-000000000095", retirement_decision=decision, ) result = service.apply_requested_cutover(operation) assert result["status"] == "succeeded" assert result["mode"] == "kestra_primary" assert engines.calls == [ ("deactivate_n8n", "n8n-orders"), ("activate_kestra", "kestra-orders"), ]