test_kestra_runtime_contract.py 2.4 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576
  1. import os
  2. import pytest
  3. import requests
  4. from app.core.common.identifiers import new_governance_uid
  5. from app.core.orchestration.compilers import compile_kestra_flow
  6. from app.core.orchestration.engines import KestraAdapter
  7. pytestmark = pytest.mark.skipif(
  8. os.getenv("KESTRA_INTEGRATION") != "1",
  9. reason="set KESTRA_INTEGRATION=1 for the local Kestra contract",
  10. )
  11. def test_compiled_flow_validates_and_is_deployed_disabled():
  12. base_url = os.getenv(
  13. "KESTRA_BASE_URL", "http://127.0.0.1:18080/api/v1"
  14. ).rstrip("/")
  15. username = os.getenv("KESTRA_USERNAME", "admin@dataops.local")
  16. password = os.getenv("KESTRA_PASSWORD", "DataOpsKestra1!")
  17. spec = {
  18. "schema_version": "1.0",
  19. "dataflow_uid": new_governance_uid(),
  20. "name": "V51 runtime contract",
  21. "nodes": [
  22. {
  23. "id": "read_orders",
  24. "type": "sql.query",
  25. "data_source_uid": new_governance_uid(),
  26. "purpose": "read",
  27. "config": {"statement": "SELECT 1"},
  28. }
  29. ],
  30. "edges": [],
  31. "parameters": {},
  32. }
  33. plan = {
  34. "schema_version": "1.0",
  35. "timezone": "Asia/Shanghai",
  36. "triggers": [{"type": "cron", "expression": "0 2 * * *"}],
  37. "max_concurrency": 1,
  38. "conflict_policy": "skip",
  39. "timeout_seconds": 60,
  40. "retry": {"max_attempts": 2, "delay_seconds": 1},
  41. "backfill": {"max_days": 1, "max_runs": 1},
  42. }
  43. compiled = compile_kestra_flow(spec, plan, "test", 1)
  44. adapter = KestraAdapter(
  45. base_url=base_url,
  46. username=username,
  47. password=password,
  48. )
  49. auth = (username, password)
  50. flow_url = (
  51. f"{base_url}/main/flows/{compiled.namespace}/{compiled.flow_id}"
  52. )
  53. try:
  54. validation = adapter.validate(compiled.yaml)
  55. assert not any(
  56. item.get("constraints") for item in validation
  57. ), validation
  58. adapter.deploy_disabled(compiled.yaml)
  59. response = requests.get(flow_url, auth=auth, timeout=30)
  60. response.raise_for_status()
  61. deployed = response.json()
  62. assert deployed["disabled"] is True
  63. assert deployed["triggers"][0]["disabled"] is True
  64. assert deployed["namespace"] == compiled.namespace
  65. assert deployed["id"] == compiled.flow_id
  66. finally:
  67. cleanup = requests.delete(flow_url, auth=auth, timeout=30)
  68. assert cleanup.status_code in {204, 404}