import os import pytest import requests from app.core.common.identifiers import new_governance_uid from app.core.orchestration.compilers import compile_kestra_flow from app.core.orchestration.engines import KestraAdapter pytestmark = pytest.mark.skipif( os.getenv("KESTRA_INTEGRATION") != "1", reason="set KESTRA_INTEGRATION=1 for the local Kestra contract", ) def test_compiled_flow_validates_and_is_deployed_disabled(): base_url = os.getenv( "KESTRA_BASE_URL", "http://127.0.0.1:18080/api/v1" ).rstrip("/") username = os.getenv("KESTRA_USERNAME", "admin@dataops.local") password = os.getenv("KESTRA_PASSWORD", "DataOpsKestra1!") spec = { "schema_version": "1.0", "dataflow_uid": new_governance_uid(), "name": "V51 runtime contract", "nodes": [ { "id": "read_orders", "type": "sql.query", "data_source_uid": new_governance_uid(), "purpose": "read", "config": {"statement": "SELECT 1"}, } ], "edges": [], "parameters": {}, } plan = { "schema_version": "1.0", "timezone": "Asia/Shanghai", "triggers": [{"type": "cron", "expression": "0 2 * * *"}], "max_concurrency": 1, "conflict_policy": "skip", "timeout_seconds": 60, "retry": {"max_attempts": 2, "delay_seconds": 1}, "backfill": {"max_days": 1, "max_runs": 1}, } compiled = compile_kestra_flow(spec, plan, "test", 1) adapter = KestraAdapter( base_url=base_url, username=username, password=password, ) auth = (username, password) flow_url = ( f"{base_url}/main/flows/{compiled.namespace}/{compiled.flow_id}" ) try: validation = adapter.validate(compiled.yaml) assert not any( item.get("constraints") for item in validation ), validation adapter.deploy_disabled(compiled.yaml) response = requests.get(flow_url, auth=auth, timeout=30) response.raise_for_status() deployed = response.json() assert deployed["disabled"] is True assert deployed["triggers"][0]["disabled"] is True assert deployed["namespace"] == compiled.namespace assert deployed["id"] == compiled.flow_id finally: cleanup = requests.delete(flow_url, auth=auth, timeout=30) assert cleanup.status_code in {204, 404}