| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209 |
- import copy
- import pytest
- from app.core.orchestration.migration.dual_run import (
- DualRunCoordinator,
- DualRunMode,
- MigrationBlocked,
- N8nWorkflowConverter,
- ShadowIsolation,
- validate_mode_transition,
- )
- DATAFLOW_UID = "01900000-0000-7000-8000-000000000091"
- SOURCE_UID = "01900000-0000-7000-8000-000000000092"
- def n8n_read_workflow():
- return {
- "id": "n8n-read-orders",
- "name": "Read orders",
- "active": True,
- "nodes": [
- {
- "id": "trigger-1",
- "name": "Manual",
- "type": "n8n-nodes-base.manualTrigger",
- "parameters": {},
- },
- {
- "id": "query-1",
- "name": "Read Orders",
- "type": "n8n-nodes-base.postgres",
- "parameters": {
- "operation": "executeQuery",
- "query": "SELECT id, amount FROM orders ORDER BY id",
- "purpose": "read",
- },
- },
- ],
- "connections": {
- "Manual": {
- "main": [[{"node": "Read Orders", "type": "main", "index": 0}]]
- }
- },
- }
- def test_dual_run_modes_allow_only_guarded_forward_and_rollback_paths():
- assert validate_mode_transition(
- DualRunMode.N8N_PRIMARY,
- DualRunMode.N8N_PRIMARY_KESTRA_SHADOW,
- )
- assert validate_mode_transition(
- DualRunMode.N8N_PRIMARY_KESTRA_SHADOW,
- DualRunMode.KESTRA_PRIMARY_N8N_STANDBY,
- )
- assert validate_mode_transition(
- DualRunMode.KESTRA_PRIMARY_N8N_STANDBY,
- DualRunMode.KESTRA_PRIMARY,
- )
- assert validate_mode_transition(
- DualRunMode.KESTRA_PRIMARY_N8N_STANDBY,
- DualRunMode.N8N_PRIMARY_KESTRA_SHADOW,
- )
- with pytest.raises(ValueError, match="unsupported migration transition"):
- validate_mode_transition(
- DualRunMode.N8N_PRIMARY,
- DualRunMode.KESTRA_PRIMARY,
- )
- def test_converter_builds_read_only_workflow_spec_without_credentials():
- converted = N8nWorkflowConverter().convert(
- n8n_read_workflow(),
- dataflow_uid=DATAFLOW_UID,
- data_source_uids={"Read Orders": SOURCE_UID},
- )
- assert converted.workflow_spec["nodes"] == [
- {
- "id": "read_orders",
- "type": "sql.query",
- "data_source_uid": SOURCE_UID,
- "purpose": "read",
- "config": {
- "statement": "SELECT id, amount FROM orders ORDER BY id",
- "parameters": {},
- },
- }
- ]
- assert converted.workflow_spec["edges"] == []
- assert converted.schedule_plan["triggers"] == [{"type": "manual"}]
- assert converted.definition_hash
- assert "credential" not in repr(converted).lower()
- def test_converter_blocks_community_nodes_with_actionable_report():
- workflow = n8n_read_workflow()
- workflow["nodes"].append(
- {
- "id": "community-1",
- "name": "Community magic",
- "type": "@vendor/n8n-nodes-secret-magic",
- "parameters": {"operation": "run"},
- }
- )
- with pytest.raises(MigrationBlocked) as error:
- N8nWorkflowConverter().convert(
- workflow,
- dataflow_uid=DATAFLOW_UID,
- data_source_uids={"Read Orders": SOURCE_UID},
- )
- assert error.value.report["status"] == "blocked"
- assert error.value.report["workflow_id"] == "n8n-read-orders"
- assert error.value.report["blockers"] == [
- {
- "node_name": "Community magic",
- "node_type": "@vendor/n8n-nodes-secret-magic",
- "reason_code": "unsupported_node_type",
- }
- ]
- assert "script" not in repr(error.value.report).lower()
- class Engine:
- def __init__(self, name):
- self.name = name
- self.calls = []
- def execute(self, definition_id, inputs):
- self.calls.append((definition_id, copy.deepcopy(inputs)))
- return {
- "execution_id": f"{self.name}-execution-1",
- "status": "success",
- }
- class Store:
- def __init__(self):
- self.records = []
- def record_dual_run(self, record):
- self.records.append(copy.deepcopy(record))
- return copy.deepcopy(record)
- def test_shadow_execution_reuses_input_snapshot_and_never_changes_primary():
- converted = N8nWorkflowConverter().convert(
- n8n_read_workflow(),
- dataflow_uid=DATAFLOW_UID,
- data_source_uids={"Read Orders": SOURCE_UID},
- )
- n8n = Engine("n8n")
- kestra = Engine("kestra")
- store = Store()
- coordinator = DualRunCoordinator(
- n8n_engine=n8n,
- kestra_engine=kestra,
- store=store,
- )
- result = coordinator.execute_shadow(
- dataflow_uid=DATAFLOW_UID,
- environment="test",
- mode=DualRunMode.N8N_PRIMARY_KESTRA_SHADOW,
- n8n_definition_id="n8n-read-orders",
- kestra_definition_id="kestra-read-orders",
- workflow_spec=converted.workflow_spec,
- inputs={"day": "2026-07-19"},
- correlation_id="01900000-0000-7000-8000-000000000093",
- isolation=ShadowIsolation.read_only(),
- )
- assert n8n.calls[0][1] == kestra.calls[0][1]
- assert result["input_snapshot_hash"] == store.records[0]["input_snapshot_hash"]
- assert result["primary_engine"] == "n8n"
- assert result["shadow_engine"] == "kestra"
- assert result["formal_engine_changed"] is False
- def test_shadow_write_requires_explicit_isolated_target():
- write_spec = N8nWorkflowConverter().convert(
- n8n_read_workflow(),
- dataflow_uid=DATAFLOW_UID,
- data_source_uids={"Read Orders": SOURCE_UID},
- ).workflow_spec
- write_spec["nodes"][0].update(
- {
- "type": "sql.execute",
- "purpose": "write",
- "idempotency": {
- "strategy": "partition_replace",
- "key": "orders:${parameters.day}",
- },
- }
- )
- with pytest.raises(ValueError, match="isolated shadow target"):
- ShadowIsolation.read_only().validate(write_spec)
- write_spec["nodes"][0]["config"]["shadow_target"] = (
- "dataops_shadow.orders_20260719"
- )
- isolated = ShadowIsolation.isolated_target("dataops_shadow.")
- assert isolated.validate(write_spec)["strategy"] == "isolated_target"
|