test: verify pinned dataflow publication composition
Some checks failed
Dependency Audit / dependency-audit (push) Has been cancelled
Security Audit / security-audit (push) Has been cancelled

This commit is contained in:
2026-07-28 13:48:11 +02:00
parent 6afb8fea76
commit 5b79e7d377
3 changed files with 118 additions and 10 deletions

View File

@@ -1,5 +1,5 @@
#!/usr/bin/env python3
"""Exercise the connector -> datasource -> dataflow capability path."""
"""Exercise connector -> datasource -> dataflow publication capabilities."""
from __future__ import annotations
@@ -9,7 +9,17 @@ from sqlalchemy.orm import sessionmaker
from govoplan_connectors.backend.db.models import ConnectorTabularSource
from govoplan_core.auth import ApiPrincipal
from govoplan_core.core.access import PrincipalRef
from govoplan_core.core.datasources import datasource_catalogue, datasource_lifecycle
from govoplan_core.core.dataflows import (
DataflowPublicationTarget,
DataflowRunRequest,
dataflow_run_lifecycle,
)
from govoplan_core.core.datasources import (
DatasourceReadRequest,
datasource_catalogue,
datasource_lifecycle,
datasource_publication,
)
from govoplan_core.core.modules import ModuleContext
from govoplan_core.core.tabular_sources import (
TabularSnapshotInput,
@@ -22,11 +32,18 @@ from govoplan_dataflow.backend.schemas import (
GraphNode,
GraphPosition,
PipelineGraph,
PipelineCreateRequest,
PipelinePreviewRequest,
)
from govoplan_dataflow.backend.service import preview_pipeline
from govoplan_dataflow.backend.db.models import (
DataflowPipeline,
DataflowPipelineRevision,
DataflowRun,
)
from govoplan_dataflow.backend.service import create_pipeline, preview_pipeline
from govoplan_datasources.backend.db.models import (
DatasourceMaterializationRecord,
DatasourcePublicationRecord,
DatasourceRecord,
DatasourceStageRecord,
)
@@ -47,6 +64,10 @@ def main() -> int:
DatasourceRecord.__table__,
DatasourceMaterializationRecord.__table__,
DatasourceStageRecord.__table__,
DatasourcePublicationRecord.__table__,
DataflowPipeline.__table__,
DataflowPipelineRevision.__table__,
DataflowRun.__table__,
],
)
session_factory = sessionmaker(bind=engine)
@@ -55,7 +76,15 @@ def main() -> int:
writer = tabular_snapshot_writer(registry)
lifecycle = datasource_lifecycle(registry)
catalogue = datasource_catalogue(registry)
if writer is None or lifecycle is None or catalogue is None:
publisher = datasource_publication(registry)
runner = dataflow_run_lifecycle(registry)
if (
writer is None
or lifecycle is None
or catalogue is None
or publisher is None
or runner is None
):
raise RuntimeError("Datasource composition capabilities are incomplete.")
origin = writer.create_snapshot(
@@ -102,8 +131,73 @@ def main() -> int:
raise RuntimeError(f"Unexpected Dataflow rows: {result.rows!r}")
if result.source_fingerprints[0]["source_ref"] != datasource.ref:
raise RuntimeError("Dataflow lineage did not retain the datasource reference.")
pipeline = create_pipeline(
session,
tenant_id="tenant-1",
actor_id="account-1",
payload=PipelineCreateRequest(
name="Monthly case output",
status="active",
graph=_graph(
datasource_ref=datasource.ref,
fingerprint=datasource.fingerprint,
),
editor_mode="graph",
),
)
run_request = DataflowRunRequest(
pipeline_ref=f"pipeline:{pipeline.id}",
revision=1,
idempotency_key="composition-run-1",
publication=DataflowPublicationTarget(
name="Monthly case result",
source_name="monthly_case_result",
freeze=True,
frozen_label="Composition evidence",
),
)
published = runner.start_run(
session,
principal,
request=run_request,
)
replayed = runner.start_run(
session,
principal,
request=run_request,
)
if published.status != "succeeded":
raise RuntimeError(f"Dataflow publication failed: {published.error}")
if replayed.ref != published.ref or not replayed.replayed:
raise RuntimeError("Dataflow run idempotency did not replay the prior run.")
if (
not published.output_datasource_ref
or not published.output_materialization_ref
):
raise RuntimeError("Dataflow publication did not retain output references.")
output = catalogue.read_datasource(
session,
principal,
request=DatasourceReadRequest(
datasource_ref=published.output_datasource_ref,
),
)
if list(output.rows) != expected_rows:
raise RuntimeError(
f"Unexpected published Dataflow rows: {list(output.rows)!r}"
)
if (
output.materialization is None
or output.materialization.ref != published.output_materialization_ref
or output.materialization.frozen_at is None
):
raise RuntimeError(
"Published Datasource materialization is not pinned and frozen."
)
engine.dispose()
print("Connector -> Datasources -> Dataflow composition passed.")
print(
"Connector -> Datasources -> pinned Dataflow publication composition passed."
)
return 0