| 12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576 |
- 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}
|