"""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()