From 97ca670bfe87d2eadaef682ca80e5b94fe98e20f Mon Sep 17 00:00:00 2001 From: Albrecht Degering Date: Fri, 21 Aug 2026 17:36:19 +0200 Subject: [PATCH] feat(datasources): publish durable artifact outputs --- README.md | 12 +- docs/QUALITY_POLICY.md | 20 ++ pyproject.toml | 2 +- src/govoplan_datasources/__init__.py | 2 +- src/govoplan_datasources/backend/manifest.py | 38 ++- src/govoplan_datasources/backend/payloads.py | 89 +++++ src/govoplan_datasources/backend/service.py | 313 ++++++++++++++++-- tests/test_lifecycle.py | 195 +++++++++++ webui/package.json | 2 +- .../features/datasources/DatasourcesPage.tsx | 4 +- 10 files changed, 639 insertions(+), 38 deletions(-) diff --git a/README.md b/README.md index 48010ca..336b7fc 100644 --- a/README.md +++ b/README.md @@ -22,8 +22,11 @@ explicit frozen states, previews, retirement, and atomic producer publication. Origin discovery retains each provider's source mode, structured health, and declared pushdown support. Live previews preserve the provider's effective row, serialized-byte, and elapsed-time limits and its redacted diagnostics. -Producer modules can append a bounded tabular result or create a new static -datasource through an idempotent capability. The publication ledger retains the +Producer modules can append a bounded inline tabular result or pin a larger +durable artifact through the same idempotent capability. Artifact references +declare a backend, locator, SHA-256 checksum, schema, fingerprint, row and byte +counts; the installed provider verifies integrity and serves bounded reads. +The publication ledger retains the producer run, output materialization, provenance, and replay identity. On PostgreSQL, a transaction-scoped advisory lock serializes each tenant, producer, and idempotency identity before any output side effect, so retries from multiple @@ -37,8 +40,9 @@ materialization provenance. The supported contract is documented in The contracts already model database, HTTP/REST, directory, file, feed, document, binary, directory, and stream sources so providers can be added -without changing consumers. Larger durable artifact-backed publications remain -a later storage-provider slice. +without changing consumers. Storage modules contribute artifact backends +through the provider-neutral `datasources.artifactBackends` capability; the +Datasources module never imports their internals. See [docs/CONCEPT.md](docs/CONCEPT.md) for ownership and lifecycle details. diff --git a/docs/QUALITY_POLICY.md b/docs/QUALITY_POLICY.md index fb2a32f..ca468d5 100644 --- a/docs/QUALITY_POLICY.md +++ b/docs/QUALITY_POLICY.md @@ -96,6 +96,26 @@ This makes concurrent retries from separate API or worker nodes converge on the same publication and materialization rather than relying on a late uniqueness failure after output rows have already been persisted. +## Durable artifact publications + +Outputs larger than the inline row and byte limits use an immutable artifact +reference. The reference pins its backend and locator together with SHA-256 +checksum, schema, datasource fingerprint, row count, byte count, media type, +and optional resume checkpoint. Datasources persists that reference as the +materialization payload and asks the installed Core-contract artifact backend +to verify it before creating catalogue state. Reads remain bounded and are +re-authorized by Datasources before reaching the backend. + +Schema rules are evaluated by Datasources. Content-level rules such as +uniqueness or range require producer evidence bound to the exact payload +checksum and current quality-policy hash, including all evaluated rule IDs. +Missing or explicitly deferred evidence produces a `review_required` +publication and an immutable, addressable materialization, but it does not +replace the Datasource's current state. Valid warnings produce +`published_with_warnings`; failed evidence blocks the publication without a +catalogue side effect. These terminal states are preserved for Dataflow and +Workflow handoffs instead of being collapsed into generic success. + Approval authority, approval expiry, and retention/deletion execution remain separate work under `govoplan-datasources#2`. Until those contracts are added, no JSON flag is treated as an approval and no stage is deleted automatically. diff --git a/pyproject.toml b/pyproject.toml index 73991f3..5b776df 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta" [project] name = "govoplan-datasources" -version = "0.1.18" +version = "0.1.19" description = "Governed datasource catalogue, staging, and materialization lifecycle for GovOPlaN." readme = "README.md" requires-python = ">=3.12" diff --git a/src/govoplan_datasources/__init__.py b/src/govoplan_datasources/__init__.py index 47c27e5..9c78b7b 100644 --- a/src/govoplan_datasources/__init__.py +++ b/src/govoplan_datasources/__init__.py @@ -1,3 +1,3 @@ """GovOPlaN Datasources module.""" -__version__ = "0.1.18" +__version__ = "0.1.19" diff --git a/src/govoplan_datasources/backend/manifest.py b/src/govoplan_datasources/backend/manifest.py index b0786ca..6bb5827 100644 --- a/src/govoplan_datasources/backend/manifest.py +++ b/src/govoplan_datasources/backend/manifest.py @@ -7,10 +7,12 @@ from govoplan_core.core.access import ( CAPABILITY_AUTH_PRINCIPAL_RESOLVER, ) from govoplan_core.core.datasources import ( + CAPABILITY_DATASOURCE_ARTIFACT_BACKENDS, CAPABILITY_DATASOURCE_CATALOGUE, CAPABILITY_DATASOURCE_LIFECYCLE, CAPABILITY_DATASOURCE_ORIGINS, CAPABILITY_DATASOURCE_PUBLICATION, + datasource_artifact_backend_provider, ) from govoplan_core.core.module_guards import ( drop_table_retirement_provider, @@ -55,12 +57,13 @@ from govoplan_datasources.backend.service import ( STAGE_WRITE_SCOPE, SqlDatasourceProvider, ) +from govoplan_datasources.backend.payloads import ExternalArtifactPayloadBackend MODULE_ID = "datasources" MODULE_NAME = "Datasources" -MODULE_VERSION = "0.1.18" -DATASOURCE_INTERFACE_VERSION = "0.1.0" +MODULE_VERSION = "0.1.19" +DATASOURCE_INTERFACE_VERSION = "0.2.0" ARCHITECTURE = ModuleArchitectureDeclaration( layer="data_reporting_integration", @@ -85,7 +88,7 @@ ARCHITECTURE = ModuleArchitectureDeclaration( ), known_limits=( "Governance references are stable provider-neutral refs; dedicated selectors depend on the owning optional modules.", - "Quality and freshness policies are stored and snapshotted but enforcement remains provider-specific.", + "External artifact content validation depends on an installed payload backend and checksum-bound producer evidence.", ), supported_authority_modes=( "native_authoritative", @@ -182,7 +185,19 @@ def _router(context: ModuleContext): def _provider(context: ModuleContext) -> SqlDatasourceProvider: - return SqlDatasourceProvider(registry=context.registry) + provider = datasource_artifact_backend_provider(context.registry) + backends = ( + tuple( + ExternalArtifactPayloadBackend(backend) + for backend in provider.artifact_backends() + ) + if provider is not None + else () + ) + return SqlDatasourceProvider( + registry=context.registry, + payload_backends=backends, + ) def _dsar_provider(context: ModuleContext) -> DatasourcesDsarProvider: @@ -240,6 +255,7 @@ manifest = ModuleManifest( CAPABILITY_AUTH_PRINCIPAL_RESOLVER, CAPABILITY_AUTH_PERMISSION_EVALUATOR, CAPABILITY_DATASOURCE_ORIGINS, + CAPABILITY_DATASOURCE_ARTIFACT_BACKENDS, ), provides_interfaces=( ModuleInterfaceProvider( @@ -500,7 +516,10 @@ manifest = ModuleManifest( "Authority mode states whether GovOPlaN, an external system, a synchronized projection, an overlay, or a linked reference " "controls the data. The authoritative source, owner, steward, responsible organization/function, schema owner, privacy " "profile, retention policy, transfer agreement, legal basis, holds, correction procedure, purposes, official keys, and " - "known limits provide discoverable institutional context. Freshness and quality policies are typed JSON contracts retained " + "known limits provide discoverable institutional context. A retention-policy reference identifies the owning Policy rule; it " + "does not itself delete materializations, override legal holds, or prove that a scheduler is active. A transfer-agreement " + "reference records the governed agreement for external exchange; it does not grant connector credentials, recipient access, " + "or authority to export classified rows. Freshness and quality policies are typed JSON contracts retained " "with materialization evidence. Datasources enforces the declared bounded tabular stage rules and schema policy; origin-specific " "or consumer-specific controls still remain with the provider or consuming control that declares support. Metadata " "visibility never grants row access." @@ -527,6 +546,8 @@ manifest = ModuleManifest( "datasources.field.publication-state", "datasources.field.freshness-policy", "datasources.field.quality-policy", + "datasources.field.retention-policy", + "datasources.field.transfer-agreement", ], }, ), @@ -540,7 +561,11 @@ manifest = ModuleManifest( "promotion. Updates compare the detected schema with the current target and classify each change as compatible, warning, or " "breaking. Diagnostics expose counts and bounded row numbers, never field values. The policy version and hash, diagnostics, and " "schema diff are copied into immutable materialization provenance when promotion succeeds. Producer publication applies the same " - "gate before any catalogue effect, retains validation evidence on the immutable output revision, serializes a publication identity across PostgreSQL worker nodes before replay lookup, and emits a transactional terminal event. Approval and retention execution " + "gate before any catalogue effect, retains validation evidence on the immutable output revision, serializes a publication identity across PostgreSQL worker nodes before replay lookup, and emits a transactional terminal event. Large outputs may instead supply a " + "durable provider-neutral artifact reference with a pinned locator, SHA-256 checksum, declared schema, fingerprint, and size. " + "The configured payload backend verifies the artifact and provides bounded reads. Content-level rules require checksum- and " + "policy-bound producer evidence; absent evidence creates an immutable review-required materialization without changing the " + "Datasource's current state. Warning and review-required outcomes are preserved for Workflow handoffs. Approval and retention execution " "are not inferred from arbitrary JSON flags and remain separate governed lifecycle work." ), layer="available", @@ -572,6 +597,7 @@ manifest = ModuleManifest( "limitations": [ "Referential rules currently use a bounded embedded value set rather than reading another protected Datasource.", "Approval authority and automatic retention execution are not part of the current stage contract.", + "Artifact bytes remain owned by their payload backend; Datasources stores an immutable reference and integrity evidence.", ], }, ), diff --git a/src/govoplan_datasources/backend/payloads.py b/src/govoplan_datasources/backend/payloads.py index 0570549..58a4065 100644 --- a/src/govoplan_datasources/backend/payloads.py +++ b/src/govoplan_datasources/backend/payloads.py @@ -9,6 +9,9 @@ from sqlalchemy import delete, func, insert, select from sqlalchemy.orm import Session from govoplan_core.core.datasources import ( + DatasourceArtifactBackend, + DatasourceArtifactReference, + DatasourceField, DatasourceUnavailableError, DatasourceValidationError, ) @@ -116,6 +119,91 @@ class DatabaseRowsPayloadBackend: ) +class ExternalArtifactPayloadBackend: + """Adapts a Core artifact backend to persisted Datasources payload rows.""" + + def __init__(self, backend: DatasourceArtifactBackend) -> None: + self._backend = backend + self.backend = backend.backend + + def read_rows( + self, + session: Session, + payload: DatasourcePayloadRecord, + *, + offset: int, + limit: int, + ) -> Sequence[Mapping[str, object]]: + return self._backend.read_rows( + session, + tenant_id=payload.tenant_id, + artifact=_artifact_from_payload(payload), + offset=offset, + limit=limit, + ) + + def verify( + self, + session: Session, + payload: DatasourcePayloadRecord, + ) -> None: + self._backend.verify( + session, + tenant_id=payload.tenant_id, + artifact=_artifact_from_payload(payload), + ) + + def delete( + self, + session: Session, + payload: DatasourcePayloadRecord, + ) -> None: + self._backend.delete( + session, + tenant_id=payload.tenant_id, + artifact=_artifact_from_payload(payload), + ) + + +def _artifact_from_payload( + payload: DatasourcePayloadRecord, +) -> DatasourceArtifactReference: + metadata = dict(payload.metadata_) + raw_schema = metadata.get("artifact_schema") + schema = tuple( + DatasourceField( + name=str(item.get("name") or ""), + data_type=str(item.get("data_type") or "unknown"), + nullable=bool(item.get("nullable", True)), + ) + for item in raw_schema + if isinstance(item, Mapping) + ) if isinstance(raw_schema, Sequence) and not isinstance( + raw_schema, (str, bytes) + ) else () + return DatasourceArtifactReference( + backend=payload.backend, + locator=str(payload.locator or ""), + checksum=payload.checksum, + row_count=payload.row_count, + byte_count=payload.byte_count, + schema=schema, + fingerprint=str(metadata.get("publication_fingerprint") or ""), + media_type=payload.media_type, + checkpoint=dict(payload.checkpoint_), + metadata={ + str(key): value + for key, value in metadata.items() + if key not in {"artifact_schema", "artifact_validation"} + }, + validation=( + dict(metadata.get("artifact_validation")) + if isinstance(metadata.get("artifact_validation"), Mapping) + else {} + ), + ) + + class PayloadBackendRegistry: def __init__( self, @@ -363,6 +451,7 @@ __all__ = [ "DATABASE_ROWS_BACKEND", "DatabaseRowsPayloadBackend", "DatasourcePayloadBackend", + "ExternalArtifactPayloadBackend", "PayloadBackendRegistry", "create_database_rows_payload", "create_external_payload_reference", diff --git a/src/govoplan_datasources/backend/service.py b/src/govoplan_datasources/backend/service.py index ff9854c..3f6170a 100644 --- a/src/govoplan_datasources/backend/service.py +++ b/src/govoplan_datasources/backend/service.py @@ -14,6 +14,7 @@ from govoplan_core.audit.logging import audit_event from govoplan_core.auth import ApiPrincipal, has_scope from govoplan_core.core.datasources import ( DatasourceAccessError, + DatasourceArtifactReference, DatasourceDescriptor, DatasourceError, DatasourceField, @@ -25,6 +26,7 @@ from govoplan_core.core.datasources import ( DatasourceOriginReadRequest, DatasourcePublicationRequest, DatasourcePublicationResult, + DatasourcePublicationStatus, DatasourceReadRequest, DatasourceReadResult, DatasourceStage, @@ -46,6 +48,7 @@ from govoplan_datasources.backend.payloads import ( DatasourcePayloadBackend, PayloadBackendRegistry, create_database_rows_payload, + create_external_payload_reference, payload_for_materialization, validate_payload_size, ) @@ -72,9 +75,11 @@ class _PreparedPublication: producer_module: str producer_run_ref: str idempotency_key: str - rows: tuple[dict[str, Any], ...] + rows: tuple[dict[str, Any], ...] | None + artifact: DatasourceArtifactReference | None schema: tuple[DatasourceField, ...] fingerprint: str + row_count: int byte_count: int request_hash: str @@ -120,9 +125,8 @@ class SqlDatasourceProvider: tenant_id=api_principal.tenant_id, request=request, ) - validation = validate_stage( - rows=prepared.rows, - schema=prepared.schema, + validation, publication_status = _validate_publication( + prepared=prepared, quality_policy=governance.quality_policy, baseline_schema=baseline_schema, ) @@ -138,6 +142,12 @@ class SqlDatasourceProvider: "Published output failed governed quality or schema validation" f"{suffix}." ) + reusable_payload = self._publication_payload( + db, + tenant_id=api_principal.tenant_id, + actor_id=actor_id, + prepared=prepared, + ) datasource = _publication_target( db, tenant_id=api_principal.tenant_id, @@ -146,11 +156,14 @@ class SqlDatasourceProvider: prepared=prepared, target=target, governance=governance, + publish_as_current=( + request.set_current and publication_status != "review_required" + ), ) materialization = _append_materialization( db, datasource=datasource, - rows=prepared.rows, + rows=prepared.rows or (), schema=[field_payload(field) for field in prepared.schema], fingerprint=prepared.fingerprint, byte_count=prepared.byte_count, @@ -163,7 +176,15 @@ class SqlDatasourceProvider: "publication_validation": validation, }, metadata=dict(request.metadata), - set_current=request.set_current, + set_current=( + request.set_current and publication_status != "review_required" + ), + reusable_payload=reusable_payload, + state=( + "review_required" + if publication_status == "review_required" + else "published" + ), ) publication = _create_publication_record( db, @@ -173,6 +194,7 @@ class SqlDatasourceProvider: materialization=materialization, request=request, prepared=prepared, + status=publication_status, ) audit_event( db, @@ -180,7 +202,7 @@ class SqlDatasourceProvider: user_id=getattr(api_principal.user, "id", None) or api_principal.account_id, api_key_id=api_principal.api_key_id, - action="datasource.publication.published", + action=f"datasource.publication.{publication_status}", object_type="datasource_publication", object_id=publication.id, details={ @@ -189,7 +211,12 @@ class SqlDatasourceProvider: "datasource_ref": _datasource_ref(datasource.id), "materialization_ref": _materialization_ref(materialization.id), "fingerprint": prepared.fingerprint, - "row_count": len(prepared.rows), + "row_count": prepared.row_count, + "payload_backend": ( + prepared.artifact.backend + if prepared.artifact is not None + else "database_rows" + ), "policy_hash": validation["policy_hash"], "schema_classification": validation["schema_change"][ "classification" @@ -198,12 +225,52 @@ class SqlDatasourceProvider: ) return DatasourcePublicationResult( ref=_publication_ref(publication.id), - status=publication.status, + status=cast(DatasourcePublicationStatus, publication.status), datasource=_datasource_dto(datasource), materialization=_materialization_dto(materialization), replayed=False, ) + def _publication_payload( + self, + session: Session, + *, + tenant_id: str, + actor_id: str | None, + prepared: _PreparedPublication, + ) -> DatasourcePayloadRecord | None: + artifact = prepared.artifact + if artifact is None: + return None + backend = self._payload_backends.require(artifact.backend) + payload = create_external_payload_reference( + session, + tenant_id=tenant_id, + backend=artifact.backend, + locator=artifact.locator, + checksum=artifact.checksum, + row_count=artifact.row_count, + byte_count=artifact.byte_count, + actor_id=actor_id, + media_type=artifact.media_type, + checkpoint=artifact.checkpoint, + metadata={ + **dict(artifact.metadata), + "publication_fingerprint": artifact.fingerprint, + "artifact_schema": [ + field_payload(field) for field in artifact.schema + ], + "artifact_validation": dict(artifact.validation), + }, + ) + try: + backend.verify(session, payload) + except Exception: + session.delete(payload) + session.flush() + raise + return payload + def list_datasources( self, session: object, @@ -1056,6 +1123,7 @@ def _append_materialization( metadata: Mapping[str, object] | None = None, set_current: bool, reusable_payload: DatasourcePayloadRecord | None = None, + state: str = "published", ) -> DatasourceMaterializationRecord: datasource = _lock_datasource_for_materialization(session, datasource) revision = _allocate_materialization_revision(session, datasource) @@ -1082,7 +1150,7 @@ def _append_materialization( tenant_id=datasource.tenant_id, datasource_id=datasource.id, revision=revision, - state="published", + state=state, schema_version=max(1, int(schema_version or 1)), schema_=schema_payload, payload_id=payload.id, @@ -1618,17 +1686,47 @@ def _prepare_publication( producer_run_ref=producer_run_ref, idempotency_key=idempotency_key, ) - normalized = normalize_rows(request.rows) - schema = infer_schema(normalized) - fingerprint = fingerprint_rows(normalized, schema) + if (request.rows is None) == (request.artifact is None): + raise DatasourceValidationError( + "A publication requires exactly one inline row payload or durable " + "artifact reference." + ) + artifact = request.artifact + if artifact is None: + normalized: tuple[dict[str, Any], ...] | None = normalize_rows( + request.rows or () + ) + schema = infer_schema(normalized) + fingerprint = fingerprint_rows(normalized, schema) + row_count = len(normalized) + byte_count = encoded_size(normalized) + else: + normalized = None + schema = _validated_artifact_schema(artifact) + fingerprint = _validated_sha256( + artifact.fingerprint, + label="Artifact publication fingerprints", + ) + _validated_sha256( + artifact.checksum, + label="Artifact publication checksums", + ) + if artifact.row_count < 0 or artifact.byte_count < 0: + raise DatasourceValidationError( + "Artifact publication sizes cannot be negative." + ) + row_count = artifact.row_count + byte_count = artifact.byte_count return _PreparedPublication( producer_module=producer_module, producer_run_ref=producer_run_ref, idempotency_key=idempotency_key, rows=normalized, + artifact=artifact, schema=schema, fingerprint=fingerprint, - byte_count=encoded_size(normalized), + row_count=row_count, + byte_count=byte_count, request_hash=_publication_request_hash( request, normalized=normalized, @@ -1637,6 +1735,147 @@ def _prepare_publication( ) +def _validated_artifact_schema( + artifact: DatasourceArtifactReference, +) -> tuple[DatasourceField, ...]: + schema = tuple(artifact.schema) + names = [field.name.strip() for field in schema] + if not schema or any(not name for name in names): + raise DatasourceValidationError( + "Artifact publications require a non-empty schema." + ) + if len(names) != len(set(names)): + raise DatasourceValidationError( + "Artifact publication schema field names must be unique." + ) + return schema + + +def _validated_sha256(value: str, *, label: str) -> str: + cleaned = value.strip().casefold() + try: + valid = len(cleaned) == 64 and int(cleaned, 16) >= 0 + except ValueError: + valid = False + if not valid: + raise DatasourceValidationError(f"{label} must be SHA-256 values.") + return cleaned + + +def _validate_publication( + *, + prepared: _PreparedPublication, + quality_policy: Mapping[str, object], + baseline_schema: Sequence[DatasourceField] | None, +) -> tuple[dict[str, object], DatasourcePublicationStatus]: + if prepared.artifact is None: + validation = validate_stage( + rows=prepared.rows or (), + schema=prepared.schema, + quality_policy=quality_policy, + baseline_schema=baseline_schema, + ) + status: DatasourcePublicationStatus = ( + "published_with_warnings" + if validation.get("warnings") + else "published" + ) + return validation, status + + raw_rules = quality_policy.get("rules", []) + rules = ( + [item for item in raw_rules if isinstance(item, Mapping)] + if isinstance(raw_rules, Sequence) + and not isinstance(raw_rules, (str, bytes)) + else [] + ) + metadata_rule_types = {"required_fields", "field"} + metadata_policy = { + **dict(quality_policy), + "rules": [ + dict(item) + for item in rules + if str(item.get("type") or "") in metadata_rule_types + ], + } + validation = validate_stage( + rows=(), + schema=prepared.schema, + quality_policy=metadata_policy, + baseline_schema=baseline_schema, + ) + if validation["valid"] is not True: + return validation, "review_required" + + evidence = dict(prepared.artifact.validation) + content_rule_ids = { + str(item.get("id") or f"rule-{index + 1}") + for index, item in enumerate(rules) + if str(item.get("type") or "") not in metadata_rule_types + } + policy_probe = validate_stage( + rows=(), + schema=prepared.schema, + quality_policy=quality_policy, + baseline_schema=baseline_schema, + ) + evidence_rule_ids = { + str(item) + for item in evidence.get("rules_evaluated", []) + if str(item).strip() + } if isinstance(evidence.get("rules_evaluated"), Sequence) and not isinstance( + evidence.get("rules_evaluated"), (str, bytes) + ) else set() + evidence_verified = ( + not content_rule_ids + or ( + evidence.get("valid") is True + and str(evidence.get("policy_hash") or "") + == str(policy_probe["policy_hash"]) + and str(evidence.get("payload_checksum") or "").casefold() + == prepared.artifact.checksum.casefold() + and content_rule_ids.issubset(evidence_rule_ids) + ) + ) + requested_status = str(evidence.get("status") or "verified") + if evidence.get("valid") is False or requested_status == "failed": + errors = evidence.get("errors") + validation["valid"] = False + validation["errors"] = ( + list(errors) + if isinstance(errors, Sequence) + and not isinstance(errors, (str, bytes)) + else [ + { + "severity": "error", + "code": "quality.artifact_validation", + "message": "Artifact validation evidence reports failure.", + } + ] + ) + return validation, "review_required" + + evidence_warnings = evidence.get("warnings") + if isinstance(evidence_warnings, Sequence) and not isinstance( + evidence_warnings, (str, bytes) + ): + validation["warnings"] = [ + *list(validation.get("warnings", [])), + *list(evidence_warnings), + ] + validation["artifact_evidence"] = { + "verified": evidence_verified, + "rules_required": sorted(content_rule_ids), + "rules_evaluated": sorted(evidence_rule_ids), + "payload_checksum": prepared.artifact.checksum.casefold(), + } + if not evidence_verified or requested_status == "review_required": + return validation, "review_required" + if validation.get("warnings") or requested_status == "warning": + return validation, "published_with_warnings" + return validation, "published" + + def _validate_publication_identity( *, producer_module: str, @@ -1730,7 +1969,7 @@ def _existing_publication_result( ) return DatasourcePublicationResult( ref=_publication_ref(publication.id), - status=publication.status, + status=cast(DatasourcePublicationStatus, publication.status), datasource=_datasource_dto(datasource), materialization=_materialization_dto(materialization), replayed=True, @@ -1746,6 +1985,7 @@ def _publication_target( prepared: _PreparedPublication, target: DatasourceRecord | None, governance: DatasourceGovernance, + publish_as_current: bool, ) -> DatasourceRecord: if target is not None: datasource = target @@ -1782,9 +2022,9 @@ def _publication_target( provider_ref=prepared.producer_run_ref, schema_version=1, schema_=[field_payload(field) for field in prepared.schema], - fingerprint=prepared.fingerprint, - row_count=len(prepared.rows), - byte_count=prepared.byte_count, + fingerprint=prepared.fingerprint if publish_as_current else "", + row_count=prepared.row_count if publish_as_current else None, + byte_count=prepared.byte_count if publish_as_current else None, provenance_={ **dict(request.provenance), "producer_module": prepared.producer_module, @@ -1859,6 +2099,7 @@ def _create_publication_record( materialization: DatasourceMaterializationRecord, request: DatasourcePublicationRequest, prepared: _PreparedPublication, + status: DatasourcePublicationStatus, ) -> DatasourcePublicationRecord: publication = DatasourcePublicationRecord( tenant_id=tenant_id, @@ -1868,12 +2109,19 @@ def _create_publication_record( request_hash=prepared.request_hash, datasource_id=datasource.id, materialization_id=materialization.id, - status="published", + status=status, details_={ "fingerprint": prepared.fingerprint, - "row_count": len(prepared.rows), - "set_current": request.set_current, + "row_count": prepared.row_count, + "set_current": ( + request.set_current and status != "review_required" + ), "frozen": request.freeze, + "payload_backend": ( + prepared.artifact.backend + if prepared.artifact is not None + else "database_rows" + ), "validation": materialization.provenance_.get( "publication_validation", {}, @@ -1889,7 +2137,7 @@ def _create_publication_record( def _publication_request_hash( request: DatasourcePublicationRequest, *, - normalized: Sequence[Mapping[str, object]], + normalized: Sequence[Mapping[str, object]] | None, fingerprint: str, ) -> str: payload = { @@ -1899,7 +2147,26 @@ def _publication_request_hash( "name": request.name, "source_name": request.source_name, "description": request.description, - "rows": [dict(row) for row in normalized], + "rows": [dict(row) for row in (normalized or ())], + "artifact": ( + { + "backend": request.artifact.backend, + "locator": request.artifact.locator, + "checksum": request.artifact.checksum.casefold(), + "row_count": request.artifact.row_count, + "byte_count": request.artifact.byte_count, + "schema": [ + field_payload(field) for field in request.artifact.schema + ], + "fingerprint": request.artifact.fingerprint.casefold(), + "media_type": request.artifact.media_type, + "checkpoint": dict(request.artifact.checkpoint), + "metadata": dict(request.artifact.metadata), + "validation": dict(request.artifact.validation), + } + if request.artifact is not None + else None + ), "fingerprint": fingerprint, "freeze": request.freeze, "frozen_label": request.frozen_label, diff --git a/tests/test_lifecycle.py b/tests/test_lifecycle.py index e40ec06..2da7779 100644 --- a/tests/test_lifecycle.py +++ b/tests/test_lifecycle.py @@ -11,6 +11,7 @@ from govoplan_core.core.change_sequence import ChangeSequenceEntry from govoplan_core.core.datasources import ( CAPABILITY_DATASOURCE_ORIGINS, DatasourceAccessError, + DatasourceArtifactReference, DatasourceField, DatasourceGovernance, DatasourceOrigin, @@ -44,6 +45,7 @@ from govoplan_datasources.backend.service import ( SqlDatasourceProvider, ) from govoplan_datasources.backend.payloads import ( + ExternalArtifactPayloadBackend, create_database_rows_payload, finalize_payload_deletion, mark_unreferenced_payload_for_deletion, @@ -166,6 +168,56 @@ class FakeRegistry: return self.origin_provider +class FakeArtifactBackend: + backend = "test_artifact" + + def __init__(self) -> None: + self.verified: list[str] = [] + self.deleted: list[str] = [] + + def read_rows( + self, + _session, + *, + tenant_id: str, + artifact: DatasourceArtifactReference, + offset: int, + limit: int, + ): + self.assert_tenant(tenant_id) + stop = min(artifact.row_count, offset + limit) + return tuple( + {"id": index, "result": "match"} + for index in range(offset, stop) + ) + + def verify( + self, + _session, + *, + tenant_id: str, + artifact: DatasourceArtifactReference, + ) -> None: + self.assert_tenant(tenant_id) + if not artifact.locator.startswith("artifact:"): + raise DatasourceUnavailableError("Unknown test artifact.") + self.verified.append(artifact.locator) + + def delete( + self, + _session, + *, + tenant_id: str, + artifact: DatasourceArtifactReference, + ) -> None: + self.assert_tenant(tenant_id) + self.deleted.append(artifact.locator) + + def assert_tenant(self, tenant_id: str) -> None: + if tenant_id != "tenant-1": + raise AssertionError("Artifact backend crossed a tenant boundary.") + + class DatasourceLifecycleTests(unittest.TestCase): def setUp(self) -> None: self.engine = create_engine("sqlite:///:memory:") @@ -185,8 +237,10 @@ class DatasourceLifecycleTests(unittest.TestCase): self.Session = sessionmaker(bind=self.engine) self.session = self.Session() self.origins = FakeOriginProvider() + self.artifacts = FakeArtifactBackend() self.provider = SqlDatasourceProvider( registry=FakeRegistry(self.origins), + payload_backends=(ExternalArtifactPayloadBackend(self.artifacts),), ) def tearDown(self) -> None: @@ -667,6 +721,147 @@ class DatasourceLifecycleTests(unittest.TestCase): self.session.query(DatasourcePublicationRecord).count(), ) + def test_artifact_publication_pins_large_payload_and_supports_bounded_reads( + self, + ) -> None: + artifact = DatasourceArtifactReference( + backend="test_artifact", + locator="artifact:monthly-output", + checksum="a" * 64, + row_count=25_000, + byte_count=12_000_000, + schema=( + DatasourceField("id", "integer", nullable=False), + DatasourceField("result", "string", nullable=False), + ), + fingerprint="b" * 64, + ) + + published = self.provider.publish_rows( + self.session, + principal(scopes=(SOURCE_WRITE_SCOPE, CATALOGUE_READ_SCOPE)), + request=DatasourcePublicationRequest( + producer_module="dataflow", + producer_run_ref="dataflow-run:large-output", + idempotency_key="large-output", + name="Large output", + source_name="large_output", + artifact=artifact, + ), + ) + preview = self.provider.read_datasource( + self.session, + principal(scopes=(CATALOGUE_READ_SCOPE,)), + request=DatasourceReadRequest( + datasource_ref=published.datasource.ref, + offset=10, + limit=3, + ), + ) + + self.assertEqual("published", published.status) + self.assertEqual(25_000, published.materialization.row_count) + self.assertEqual( + [{"id": 10, "result": "match"}, + {"id": 11, "result": "match"}, + {"id": 12, "result": "match"}], + list(preview.rows), + ) + self.assertEqual( + ["artifact:monthly-output", "artifact:monthly-output"], + self.artifacts.verified, + ) + + def test_unattested_artifact_quality_rules_require_review_without_becoming_current( + self, + ) -> None: + published = self.provider.publish_rows( + self.session, + principal(scopes=(SOURCE_WRITE_SCOPE,)), + request=DatasourcePublicationRequest( + producer_module="reporting", + producer_run_ref="report-run:review", + idempotency_key="review-output", + name="Review output", + source_name="review_output", + artifact=DatasourceArtifactReference( + backend="test_artifact", + locator="artifact:review-output", + checksum="c" * 64, + row_count=2, + byte_count=128, + schema=(DatasourceField("id", "integer", False),), + fingerprint="d" * 64, + ), + governance=DatasourceGovernance( + quality_policy={ + "version": "unique-id-v1", + "rules": [ + { + "id": "unique-id", + "type": "unique", + "fields": ["id"], + } + ], + } + ), + ), + ) + record = self.session.get( + DatasourceRecord, + published.datasource.ref.removeprefix("datasource:"), + ) + + self.assertEqual("review_required", published.status) + self.assertEqual("review_required", published.materialization.state) + self.assertIsNotNone(record) + assert record is not None + self.assertIsNone(record.current_materialization_id) + + def test_artifact_warning_is_a_notification_ready_terminal_state(self) -> None: + published = self.provider.publish_rows( + self.session, + principal(scopes=(SOURCE_WRITE_SCOPE,)), + request=DatasourcePublicationRequest( + producer_module="dataflow", + producer_run_ref="dataflow-run:warning", + idempotency_key="warning-output", + name="Warning output", + source_name="warning_output", + artifact=DatasourceArtifactReference( + backend="test_artifact", + locator="artifact:warning-output", + checksum="e" * 64, + row_count=1, + byte_count=64, + schema=(DatasourceField("id", "integer", False),), + fingerprint="f" * 64, + validation={ + "status": "warning", + "warnings": [ + { + "severity": "warning", + "code": "producer.partial_match", + "message": "One source used a fallback match.", + } + ], + }, + ), + ), + ) + record = self.session.get( + DatasourcePublicationRecord, + published.ref.removeprefix("publication:"), + ) + + self.assertEqual("published_with_warnings", published.status) + self.assertIsNotNone(record) + assert record is not None + self.assertEqual( + "published_with_warnings", + record.status, + ) + def test_publication_idempotency_key_rejects_different_output(self) -> None: producer = principal(scopes=(SOURCE_WRITE_SCOPE,)) base = DatasourcePublicationRequest( diff --git a/webui/package.json b/webui/package.json index 311bc90..3481c0c 100644 --- a/webui/package.json +++ b/webui/package.json @@ -1,6 +1,6 @@ { "name": "@govoplan/datasources-webui", - "version": "0.1.18", + "version": "0.1.19", "private": true, "type": "module", "main": "src/index.ts", diff --git a/webui/src/features/datasources/DatasourcesPage.tsx b/webui/src/features/datasources/DatasourcesPage.tsx index 926afd2..6326248 100644 --- a/webui/src/features/datasources/DatasourcesPage.tsx +++ b/webui/src/features/datasources/DatasourcesPage.tsx @@ -1017,10 +1017,10 @@ function GovernanceDialog({ setValue("privacy_profile_ref", event.target.value || null)} /> - + setValue("retention_policy_ref", event.target.value || null)} /> - + setValue("transfer_agreement_ref", event.target.value || null)} />