test_workflow_engine_cutover.py 9.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294
  1. import copy
  2. import pytest
  3. from app.core.orchestration.migration.cutover import (
  4. CutoverGates,
  5. EngineRoleController,
  6. WorkflowCutoverService,
  7. )
  8. from app.core.orchestration.migration.dual_run import DualRunMode
  9. from app.core.orchestration.migration.retirement import (
  10. N8nRetirementChecklist,
  11. N8nRetirementGate,
  12. )
  13. class Store:
  14. def __init__(self):
  15. self.mode = DualRunMode.N8N_PRIMARY_KESTRA_SHADOW
  16. self.operations = {}
  17. self.events = []
  18. self.completed = []
  19. self.rolled_back = []
  20. def get_state(self, dataflow_uid, environment):
  21. return {
  22. "dataflow_uid": dataflow_uid,
  23. "environment": environment,
  24. "mode": self.mode.value,
  25. "n8n_definition_id": "n8n-orders",
  26. "kestra_definition_id": "kestra-orders",
  27. "status": "stable",
  28. }
  29. def request_transition(self, record):
  30. previous = self.operations.get(record["idempotency_key"])
  31. if previous:
  32. return False, copy.deepcopy(previous)
  33. operation = {
  34. **copy.deepcopy(record),
  35. "operation_id": "cutover-operation-1",
  36. "status": "claimed",
  37. }
  38. self.operations[record["idempotency_key"]] = operation
  39. self.events.append(
  40. {
  41. "event_type": "workflow.engine.cutover.requested",
  42. "payload": copy.deepcopy(operation),
  43. }
  44. )
  45. return True, copy.deepcopy(operation)
  46. def complete_transition(self, operation_id, target_mode, result):
  47. self.mode = DualRunMode(target_mode)
  48. self.completed.append((operation_id, target_mode, copy.deepcopy(result)))
  49. def rollback_transition(self, operation_id, source_mode, result):
  50. self.mode = DualRunMode(source_mode)
  51. self.rolled_back.append((operation_id, source_mode, copy.deepcopy(result)))
  52. class Engines:
  53. def __init__(self):
  54. self.calls = []
  55. self.fail_kestra_activate = False
  56. def deactivate_n8n(self, definition_id):
  57. self.calls.append(("deactivate_n8n", definition_id))
  58. return {"status": "disabled"}
  59. def activate_n8n(self, definition_id):
  60. self.calls.append(("activate_n8n", definition_id))
  61. return {"status": "enabled"}
  62. def deactivate_kestra(self, definition_id):
  63. self.calls.append(("deactivate_kestra", definition_id))
  64. return {"status": "disabled"}
  65. def activate_kestra(self, definition_id):
  66. self.calls.append(("activate_kestra", definition_id))
  67. if self.fail_kestra_activate:
  68. raise RuntimeError("Kestra unavailable")
  69. return {"status": "enabled"}
  70. class Adapter:
  71. def __init__(self):
  72. self.calls = []
  73. def activate(self, namespace, definition_id):
  74. self.calls.append(("activate", namespace, definition_id))
  75. return {"status": "enabled"}
  76. def deactivate(self, namespace, definition_id):
  77. self.calls.append(("deactivate", namespace, definition_id))
  78. return {"status": "disabled"}
  79. def test_engine_role_controller_binds_fixed_namespaces():
  80. n8n = Adapter()
  81. kestra = Adapter()
  82. controller = EngineRoleController(
  83. n8n=n8n,
  84. kestra=kestra,
  85. n8n_namespace="dataops",
  86. kestra_namespace="dataops.migration",
  87. )
  88. controller.deactivate_n8n("n8n-orders")
  89. controller.activate_kestra("kestra-orders")
  90. assert n8n.calls == [("deactivate", "dataops", "n8n-orders")]
  91. assert kestra.calls == [
  92. ("activate", "dataops.migration", "kestra-orders")
  93. ]
  94. def passing_gates():
  95. return CutoverGates(
  96. inventory_complete=True,
  97. workflow_spec_valid=True,
  98. rollback_target_present=True,
  99. shadow_runs=3,
  100. required_shadow_runs=3,
  101. failure_recovery_observed=True,
  102. reconciliation_passed=True,
  103. connection_budget_passed=True,
  104. sla_passed=True,
  105. n8n_baseline_unchanged=True,
  106. )
  107. def test_cutover_request_is_transactional_outbox_only():
  108. store = Store()
  109. engines = Engines()
  110. service = WorkflowCutoverService(store=store, engines=engines)
  111. result = service.request_cutover(
  112. dataflow_uid="01900000-0000-7000-8000-000000000094",
  113. environment="test",
  114. target_mode=DualRunMode.KESTRA_PRIMARY_N8N_STANDBY,
  115. gates=passing_gates(),
  116. idempotency_key="cutover:orders:v55",
  117. correlation_id="01900000-0000-7000-8000-000000000095",
  118. )
  119. assert result["status"] == "claimed"
  120. assert result["source_mode"] == "n8n_primary_kestra_shadow"
  121. assert result["target_mode"] == "kestra_primary_n8n_standby"
  122. assert engines.calls == []
  123. assert store.events[0]["event_type"] == "workflow.engine.cutover.requested"
  124. assert "password" not in repr(store.events[0]).lower()
  125. def test_cutover_consumer_disables_n8n_before_enabling_kestra():
  126. store = Store()
  127. engines = Engines()
  128. service = WorkflowCutoverService(store=store, engines=engines)
  129. operation = service.request_cutover(
  130. dataflow_uid="01900000-0000-7000-8000-000000000094",
  131. environment="test",
  132. target_mode=DualRunMode.KESTRA_PRIMARY_N8N_STANDBY,
  133. gates=passing_gates(),
  134. idempotency_key="cutover:orders:v55",
  135. correlation_id="01900000-0000-7000-8000-000000000095",
  136. )
  137. result = service.apply_requested_cutover(operation)
  138. assert result["status"] == "succeeded"
  139. assert result["mode"] == "kestra_primary_n8n_standby"
  140. assert engines.calls == [
  141. ("deactivate_n8n", "n8n-orders"),
  142. ("activate_kestra", "kestra-orders"),
  143. ]
  144. assert store.completed[0][0] == "cutover-operation-1"
  145. def test_kestra_activation_failure_automatically_restores_n8n():
  146. store = Store()
  147. engines = Engines()
  148. engines.fail_kestra_activate = True
  149. service = WorkflowCutoverService(store=store, engines=engines)
  150. operation = service.request_cutover(
  151. dataflow_uid="01900000-0000-7000-8000-000000000094",
  152. environment="test",
  153. target_mode=DualRunMode.KESTRA_PRIMARY_N8N_STANDBY,
  154. gates=passing_gates(),
  155. idempotency_key="cutover:orders:v55",
  156. correlation_id="01900000-0000-7000-8000-000000000095",
  157. )
  158. result = service.apply_requested_cutover(operation)
  159. assert result["status"] == "rolled_back"
  160. assert result["mode"] == "n8n_primary_kestra_shadow"
  161. assert engines.calls == [
  162. ("deactivate_n8n", "n8n-orders"),
  163. ("activate_kestra", "kestra-orders"),
  164. ("activate_n8n", "n8n-orders"),
  165. ]
  166. assert store.rolled_back[0][1] == "n8n_primary_kestra_shadow"
  167. def test_same_cutover_key_is_idempotent():
  168. store = Store()
  169. service = WorkflowCutoverService(store=store, engines=Engines())
  170. arguments = {
  171. "dataflow_uid": "01900000-0000-7000-8000-000000000094",
  172. "environment": "test",
  173. "target_mode": DualRunMode.KESTRA_PRIMARY_N8N_STANDBY,
  174. "gates": passing_gates(),
  175. "idempotency_key": "cutover:orders:v55",
  176. "correlation_id": "01900000-0000-7000-8000-000000000095",
  177. }
  178. first = service.request_cutover(**arguments)
  179. second = service.request_cutover(**arguments)
  180. assert first == second
  181. assert len(store.events) == 1
  182. def test_cutover_gates_fail_closed():
  183. failed = CutoverGates(
  184. **{
  185. **passing_gates().as_dict(),
  186. "reconciliation_passed": False,
  187. }
  188. )
  189. with pytest.raises(ValueError, match="reconciliation_passed"):
  190. WorkflowCutoverService(store=Store(), engines=Engines()).request_cutover(
  191. dataflow_uid="01900000-0000-7000-8000-000000000094",
  192. environment="test",
  193. target_mode=DualRunMode.KESTRA_PRIMARY_N8N_STANDBY,
  194. gates=failed,
  195. idempotency_key="cutover:orders:v55",
  196. correlation_id="01900000-0000-7000-8000-000000000095",
  197. )
  198. def test_kestra_primary_requires_independent_retirement_approval():
  199. store = Store()
  200. store.mode = DualRunMode.KESTRA_PRIMARY_N8N_STANDBY
  201. service = WorkflowCutoverService(store=store, engines=Engines())
  202. with pytest.raises(ValueError, match="retirement approval"):
  203. service.request_cutover(
  204. dataflow_uid="01900000-0000-7000-8000-000000000094",
  205. environment="test",
  206. target_mode=DualRunMode.KESTRA_PRIMARY,
  207. gates=passing_gates(),
  208. idempotency_key="retire:orders:v55",
  209. correlation_id="01900000-0000-7000-8000-000000000095",
  210. )
  211. def test_approved_retirement_can_finalize_kestra_primary_mode():
  212. store = Store()
  213. store.mode = DualRunMode.KESTRA_PRIMARY_N8N_STANDBY
  214. engines = Engines()
  215. service = WorkflowCutoverService(store=store, engines=engines)
  216. decision = N8nRetirementGate.evaluate(
  217. N8nRetirementChecklist(
  218. inventory_disposition_complete=True,
  219. all_dataflows_kestra_primary=True,
  220. business_cycle_and_recovery_complete=True,
  221. standby_observation_no_rollback=True,
  222. monitoring_and_alerting_complete=True,
  223. dependency_failures_safe=True,
  224. secure_archive_complete=True,
  225. independent_review_approved=True,
  226. final_l4_passed=True,
  227. )
  228. )
  229. operation = service.request_cutover(
  230. dataflow_uid="01900000-0000-7000-8000-000000000094",
  231. environment="test",
  232. target_mode=DualRunMode.KESTRA_PRIMARY,
  233. gates=passing_gates(),
  234. idempotency_key="retire:orders:v55",
  235. correlation_id="01900000-0000-7000-8000-000000000095",
  236. retirement_decision=decision,
  237. )
  238. result = service.apply_requested_cutover(operation)
  239. assert result["status"] == "succeeded"
  240. assert result["mode"] == "kestra_primary"
  241. assert engines.calls == [
  242. ("deactivate_n8n", "n8n-orders"),
  243. ("activate_kestra", "kestra-orders"),
  244. ]