| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294 |
- 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"),
- ]
|