| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223 |
- """Opt-in end-to-end Runner test against governed PostgreSQL/MySQL pools."""
- import os
- import pytest
- import requests
- from neo4j import GraphDatabase
- from sqlalchemy import create_engine, text
- from app.core.common.identifiers import new_governance_uid
- from app.core.data_source.credentials import (
- CredentialCodec,
- DataSourceCredentialRepository,
- )
- from app.core.data_source.definitions import DataSourceDefinitionRepository
- from app.core.data_source.models import (
- DataSourceCredential,
- DataSourceDefinition,
- )
- from app.runner.auth import TaskTokenIssuer
- pytestmark = pytest.mark.integration
- def _require_runner():
- if os.environ.get("RUN_RUNNER_INTEGRATION") != "1":
- pytest.skip("set RUN_RUNNER_INTEGRATION=1 to test the Runner container")
- def _seed_source(graph_driver, platform_engine, definition, password):
- codec = CredentialCodec.from_base64(
- os.environ.get(
- "DATASOURCE_CREDENTIAL_MASTER_KEY",
- "MDEyMzQ1Njc4OWFiY2RlZjAxMjM0NTY3ODlhYmNkZWY=",
- ),
- "v1",
- )
- credentials = DataSourceCredentialRepository(codec)
- with platform_engine.begin() as connection:
- sealed = credentials.create_version(
- connection,
- data_source_uid=definition.uid,
- credential=DataSourceCredential("source_reader", password),
- actor_uid=None,
- )
- saved = DataSourceDefinition(
- **{
- **definition.__dict__,
- "credential_ref": definition.uid,
- "credential_version": sealed.credential_version,
- }
- )
- with graph_driver.session() as session:
- DataSourceDefinitionRepository(session).save(saved)
- def _cleanup(graph_driver, platform_engine, dataflow_uid, source_uids):
- with graph_driver.session() as session:
- repository = DataSourceDefinitionRepository(session)
- for uid in source_uids:
- repository.delete(uid)
- with platform_engine.begin() as connection:
- connection.execute(
- text(
- """
- DELETE FROM public.runner_task_executions
- WHERE dataflow_uid = CAST(:dataflow_uid AS uuid)
- """
- ),
- {"dataflow_uid": dataflow_uid},
- )
- connection.execute(
- text(
- """
- DELETE FROM public.datasource_credential_audit_events
- WHERE data_source_uid = ANY(CAST(:uids AS uuid[]))
- """
- ),
- {"uids": list(source_uids)},
- )
- connection.execute(
- text(
- """
- DELETE FROM public.datasource_credentials
- WHERE data_source_uid = ANY(CAST(:uids AS uuid[]))
- """
- ),
- {"uids": list(source_uids)},
- )
- def test_runner_reuses_governed_pools_and_rejects_duplicate_submit():
- _require_runner()
- runner_url = os.environ.get(
- "TEST_RUNNER_URL",
- "http://127.0.0.1:15600",
- )
- platform_engine = create_engine(
- os.environ.get(
- "TEST_DATABASE_URL",
- "postgresql://dataops:dataops-test-password"
- "@127.0.0.1:15432/dataops",
- )
- )
- graph_driver = GraphDatabase.driver(
- os.environ.get("TEST_NEO4J_URI", "bolt://127.0.0.1:17687"),
- auth=("neo4j", "Passw0rd"),
- )
- dataflow_uid = new_governance_uid()
- source_uids = (new_governance_uid(), new_governance_uid())
- definitions = (
- DataSourceDefinition(
- uid=source_uids[0],
- name_en="runner-acceptance-postgres",
- database_type="postgresql",
- host="source-postgres",
- port=5432,
- database="acceptance",
- pool_size=1,
- max_overflow=1,
- ),
- DataSourceDefinition(
- uid=source_uids[1],
- name_en="runner-acceptance-mysql",
- database_type="mysql",
- host="source-mysql",
- port=3306,
- database="acceptance",
- pool_size=1,
- max_overflow=1,
- ),
- )
- issuer = TaskTokenIssuer(
- os.environ.get(
- "RUNNER_TASK_TOKEN_SECRET",
- "dataops-local-runner-task-token-secret-change-me",
- )
- )
- try:
- _seed_source(
- graph_driver,
- platform_engine,
- definitions[0],
- "source-test-password",
- )
- _seed_source(
- graph_driver,
- platform_engine,
- definitions[1],
- "source-test-password",
- )
- for definition in definitions:
- task_uid = new_governance_uid()
- node = {
- "id": f"read_{definition.database_type}",
- "type": "sql.query",
- "data_source_uid": definition.uid,
- "purpose": "read",
- "config": {
- "statement": (
- "SELECT customer_name FROM acceptance_customers "
- "WHERE id >= :minimum_id ORDER BY id"
- ),
- "parameters": {
- "minimum_id": "${parameters.minimum_id}"
- },
- },
- }
- token = issuer.issue(
- task_uid=task_uid,
- dataflow_uid=dataflow_uid,
- workflow_version=1,
- correlation_id=new_governance_uid(),
- node=node,
- )
- payload = {
- "task_token": token,
- "node": node,
- "parameters": {"minimum_id": 1},
- }
- response = requests.post(
- f"{runner_url}/v1/tasks/execute",
- json=payload,
- timeout=20,
- )
- response.raise_for_status()
- assert response.json()["result"]["rows"] == [
- {"customer_name": "Alpha"},
- {"customer_name": "Beta"},
- ]
- replay = requests.post(
- f"{runner_url}/v1/tasks/execute",
- json=payload,
- timeout=20,
- )
- assert replay.status_code == 409
- with platform_engine.connect() as connection:
- rows = connection.execute(
- text(
- """
- SELECT status, commit_outcome, COUNT(*) AS count
- FROM public.runner_task_executions
- WHERE dataflow_uid = CAST(:uid AS uuid)
- GROUP BY status, commit_outcome
- """
- ),
- {"uid": dataflow_uid},
- ).mappings().all()
- assert [dict(row) for row in rows] == [
- {
- "status": "success",
- "commit_outcome": "not_applicable",
- "count": 2,
- }
- ]
- finally:
- _cleanup(graph_driver, platform_engine, dataflow_uid, source_uids)
- graph_driver.close()
- platform_engine.dispose()
|