test_n8n_kestra_dual_run.py 6.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209
  1. import copy
  2. import pytest
  3. from app.core.orchestration.migration.dual_run import (
  4. DualRunCoordinator,
  5. DualRunMode,
  6. MigrationBlocked,
  7. N8nWorkflowConverter,
  8. ShadowIsolation,
  9. validate_mode_transition,
  10. )
  11. DATAFLOW_UID = "01900000-0000-7000-8000-000000000091"
  12. SOURCE_UID = "01900000-0000-7000-8000-000000000092"
  13. def n8n_read_workflow():
  14. return {
  15. "id": "n8n-read-orders",
  16. "name": "Read orders",
  17. "active": True,
  18. "nodes": [
  19. {
  20. "id": "trigger-1",
  21. "name": "Manual",
  22. "type": "n8n-nodes-base.manualTrigger",
  23. "parameters": {},
  24. },
  25. {
  26. "id": "query-1",
  27. "name": "Read Orders",
  28. "type": "n8n-nodes-base.postgres",
  29. "parameters": {
  30. "operation": "executeQuery",
  31. "query": "SELECT id, amount FROM orders ORDER BY id",
  32. "purpose": "read",
  33. },
  34. },
  35. ],
  36. "connections": {
  37. "Manual": {
  38. "main": [[{"node": "Read Orders", "type": "main", "index": 0}]]
  39. }
  40. },
  41. }
  42. def test_dual_run_modes_allow_only_guarded_forward_and_rollback_paths():
  43. assert validate_mode_transition(
  44. DualRunMode.N8N_PRIMARY,
  45. DualRunMode.N8N_PRIMARY_KESTRA_SHADOW,
  46. )
  47. assert validate_mode_transition(
  48. DualRunMode.N8N_PRIMARY_KESTRA_SHADOW,
  49. DualRunMode.KESTRA_PRIMARY_N8N_STANDBY,
  50. )
  51. assert validate_mode_transition(
  52. DualRunMode.KESTRA_PRIMARY_N8N_STANDBY,
  53. DualRunMode.KESTRA_PRIMARY,
  54. )
  55. assert validate_mode_transition(
  56. DualRunMode.KESTRA_PRIMARY_N8N_STANDBY,
  57. DualRunMode.N8N_PRIMARY_KESTRA_SHADOW,
  58. )
  59. with pytest.raises(ValueError, match="unsupported migration transition"):
  60. validate_mode_transition(
  61. DualRunMode.N8N_PRIMARY,
  62. DualRunMode.KESTRA_PRIMARY,
  63. )
  64. def test_converter_builds_read_only_workflow_spec_without_credentials():
  65. converted = N8nWorkflowConverter().convert(
  66. n8n_read_workflow(),
  67. dataflow_uid=DATAFLOW_UID,
  68. data_source_uids={"Read Orders": SOURCE_UID},
  69. )
  70. assert converted.workflow_spec["nodes"] == [
  71. {
  72. "id": "read_orders",
  73. "type": "sql.query",
  74. "data_source_uid": SOURCE_UID,
  75. "purpose": "read",
  76. "config": {
  77. "statement": "SELECT id, amount FROM orders ORDER BY id",
  78. "parameters": {},
  79. },
  80. }
  81. ]
  82. assert converted.workflow_spec["edges"] == []
  83. assert converted.schedule_plan["triggers"] == [{"type": "manual"}]
  84. assert converted.definition_hash
  85. assert "credential" not in repr(converted).lower()
  86. def test_converter_blocks_community_nodes_with_actionable_report():
  87. workflow = n8n_read_workflow()
  88. workflow["nodes"].append(
  89. {
  90. "id": "community-1",
  91. "name": "Community magic",
  92. "type": "@vendor/n8n-nodes-secret-magic",
  93. "parameters": {"operation": "run"},
  94. }
  95. )
  96. with pytest.raises(MigrationBlocked) as error:
  97. N8nWorkflowConverter().convert(
  98. workflow,
  99. dataflow_uid=DATAFLOW_UID,
  100. data_source_uids={"Read Orders": SOURCE_UID},
  101. )
  102. assert error.value.report["status"] == "blocked"
  103. assert error.value.report["workflow_id"] == "n8n-read-orders"
  104. assert error.value.report["blockers"] == [
  105. {
  106. "node_name": "Community magic",
  107. "node_type": "@vendor/n8n-nodes-secret-magic",
  108. "reason_code": "unsupported_node_type",
  109. }
  110. ]
  111. assert "script" not in repr(error.value.report).lower()
  112. class Engine:
  113. def __init__(self, name):
  114. self.name = name
  115. self.calls = []
  116. def execute(self, definition_id, inputs):
  117. self.calls.append((definition_id, copy.deepcopy(inputs)))
  118. return {
  119. "execution_id": f"{self.name}-execution-1",
  120. "status": "success",
  121. }
  122. class Store:
  123. def __init__(self):
  124. self.records = []
  125. def record_dual_run(self, record):
  126. self.records.append(copy.deepcopy(record))
  127. return copy.deepcopy(record)
  128. def test_shadow_execution_reuses_input_snapshot_and_never_changes_primary():
  129. converted = N8nWorkflowConverter().convert(
  130. n8n_read_workflow(),
  131. dataflow_uid=DATAFLOW_UID,
  132. data_source_uids={"Read Orders": SOURCE_UID},
  133. )
  134. n8n = Engine("n8n")
  135. kestra = Engine("kestra")
  136. store = Store()
  137. coordinator = DualRunCoordinator(
  138. n8n_engine=n8n,
  139. kestra_engine=kestra,
  140. store=store,
  141. )
  142. result = coordinator.execute_shadow(
  143. dataflow_uid=DATAFLOW_UID,
  144. environment="test",
  145. mode=DualRunMode.N8N_PRIMARY_KESTRA_SHADOW,
  146. n8n_definition_id="n8n-read-orders",
  147. kestra_definition_id="kestra-read-orders",
  148. workflow_spec=converted.workflow_spec,
  149. inputs={"day": "2026-07-19"},
  150. correlation_id="01900000-0000-7000-8000-000000000093",
  151. isolation=ShadowIsolation.read_only(),
  152. )
  153. assert n8n.calls[0][1] == kestra.calls[0][1]
  154. assert result["input_snapshot_hash"] == store.records[0]["input_snapshot_hash"]
  155. assert result["primary_engine"] == "n8n"
  156. assert result["shadow_engine"] == "kestra"
  157. assert result["formal_engine_changed"] is False
  158. def test_shadow_write_requires_explicit_isolated_target():
  159. write_spec = N8nWorkflowConverter().convert(
  160. n8n_read_workflow(),
  161. dataflow_uid=DATAFLOW_UID,
  162. data_source_uids={"Read Orders": SOURCE_UID},
  163. ).workflow_spec
  164. write_spec["nodes"][0].update(
  165. {
  166. "type": "sql.execute",
  167. "purpose": "write",
  168. "idempotency": {
  169. "strategy": "partition_replace",
  170. "key": "orders:${parameters.day}",
  171. },
  172. }
  173. )
  174. with pytest.raises(ValueError, match="isolated shadow target"):
  175. ShadowIsolation.read_only().validate(write_spec)
  176. write_spec["nodes"][0]["config"]["shadow_target"] = (
  177. "dataops_shadow.orders_20260719"
  178. )
  179. isolated = ShadowIsolation.isolated_target("dataops_shadow.")
  180. assert isolated.validate(write_spec)["strategy"] == "isolated_target"