diff --git a/src/govoplan_core/core/datasources.py b/src/govoplan_core/core/datasources.py index c514826..a8ec014 100644 --- a/src/govoplan_core/core/datasources.py +++ b/src/govoplan_core/core/datasources.py @@ -23,6 +23,7 @@ CAPABILITY_DATASOURCE_CATALOGUE = "datasources.catalogue" CAPABILITY_DATASOURCE_LIFECYCLE = "datasources.lifecycle" CAPABILITY_DATASOURCE_PUBLICATION = "datasources.publication" CAPABILITY_DATASOURCE_ORIGINS = "connectors.datasourceOrigins" +CAPABILITY_DATASOURCE_ARTIFACT_BACKENDS = "datasources.artifactBackends" DatasourceMode = Literal["live", "cached", "static"] DatasourceKind = Literal[ @@ -37,6 +38,11 @@ DatasourceKind = Literal[ ] DatasourceShape = Literal["tabular", "document", "binary", "directory", "stream"] DatasourceConsistency = Literal["current", "live", "frozen"] +DatasourcePublicationStatus = Literal[ + "published", + "published_with_warnings", + "review_required", +] class DatasourceError(ValueError): @@ -309,12 +315,73 @@ class DatasourceStageInput: governance: DatasourceGovernance | None = None +@dataclass(frozen=True, slots=True) +class DatasourceArtifactReference: + """Immutable provider-neutral reference to a durable tabular payload. + + The producer owns creation of the payload. Datasources pins its locator, + checksum and declared shape without importing the artifact-owning module; + a configured payload backend verifies integrity and provides bounded reads. + """ + + backend: str + locator: str + checksum: str + row_count: int + byte_count: int + schema: tuple[DatasourceField, ...] + fingerprint: str + media_type: str = "application/x-ndjson" + checkpoint: Mapping[str, object] = field(default_factory=dict) + metadata: Mapping[str, object] = field(default_factory=dict) + validation: Mapping[str, object] = field(default_factory=dict) + + +@runtime_checkable +class DatasourceArtifactBackend(Protocol): + """Storage-module boundary for immutable artifact-backed tabular data.""" + + backend: str + + def verify( + self, + session: object, + *, + tenant_id: str, + artifact: DatasourceArtifactReference, + ) -> None: ... + + def read_rows( + self, + session: object, + *, + tenant_id: str, + artifact: DatasourceArtifactReference, + offset: int, + limit: int, + ) -> Sequence[Mapping[str, object]]: ... + + def delete( + self, + session: object, + *, + tenant_id: str, + artifact: DatasourceArtifactReference, + ) -> None: ... + + +@runtime_checkable +class DatasourceArtifactBackendProvider(Protocol): + def artifact_backends(self) -> Sequence[DatasourceArtifactBackend]: ... + + @dataclass(frozen=True, slots=True) class DatasourcePublicationRequest: producer_module: str producer_run_ref: str idempotency_key: str - rows: tuple[Mapping[str, object], ...] + rows: tuple[Mapping[str, object], ...] | None = None + artifact: DatasourceArtifactReference | None = None target_datasource_ref: str | None = None name: str | None = None source_name: str | None = None @@ -331,7 +398,7 @@ class DatasourcePublicationRequest: @dataclass(frozen=True, slots=True) class DatasourcePublicationResult: ref: str - status: str + status: DatasourcePublicationStatus datasource: DatasourceDescriptor materialization: DatasourceMaterialization replayed: bool = False @@ -564,6 +631,17 @@ def datasource_catalogue(registry: object | None) -> DatasourceCatalogueProvider return capability if isinstance(capability, DatasourceCatalogueProvider) else None +def datasource_artifact_backend_provider( + registry: object | None, +) -> DatasourceArtifactBackendProvider | None: + capability = _capability(registry, CAPABILITY_DATASOURCE_ARTIFACT_BACKENDS) + return ( + capability + if isinstance(capability, DatasourceArtifactBackendProvider) + else None + ) + + def datasource_lifecycle(registry: object | None) -> DatasourceLifecycleProvider | None: capability = _capability(registry, CAPABILITY_DATASOURCE_LIFECYCLE) return capability if isinstance(capability, DatasourceLifecycleProvider) else None @@ -613,9 +691,14 @@ def _governance_mapping(value: object) -> Mapping[str, object]: __all__ = [ "CAPABILITY_DATASOURCE_CATALOGUE", + "CAPABILITY_DATASOURCE_ARTIFACT_BACKENDS", "CAPABILITY_DATASOURCE_LIFECYCLE", "CAPABILITY_DATASOURCE_ORIGINS", + "CAPABILITY_DATASOURCE_PUBLICATION", "DatasourceAccessError", + "DatasourceArtifactReference", + "DatasourceArtifactBackend", + "DatasourceArtifactBackendProvider", "DatasourceCatalogueProvider", "DatasourceConsistency", "DatasourceDescriptor", @@ -631,6 +714,10 @@ __all__ = [ "DatasourceOriginProvider", "DatasourceOriginReadRequest", "DatasourceOriginReadResult", + "DatasourcePublicationProvider", + "DatasourcePublicationRequest", + "DatasourcePublicationResult", + "DatasourcePublicationStatus", "DatasourceReadRequest", "DatasourceReadResult", "DatasourceShape", @@ -639,6 +726,8 @@ __all__ = [ "DatasourceUnavailableError", "DatasourceValidationError", "datasource_catalogue", + "datasource_artifact_backend_provider", "datasource_lifecycle", "datasource_origins", + "datasource_publication", ] diff --git a/tests/test_datasource_contract.py b/tests/test_datasource_contract.py index d15ccef..a90dcd1 100644 --- a/tests/test_datasource_contract.py +++ b/tests/test_datasource_contract.py @@ -3,11 +3,13 @@ from __future__ import annotations import unittest from govoplan_core.core.datasources import ( + CAPABILITY_DATASOURCE_ARTIFACT_BACKENDS, CAPABILITY_DATASOURCE_CATALOGUE, CAPABILITY_DATASOURCE_LIFECYCLE, CAPABILITY_DATASOURCE_ORIGINS, CAPABILITY_DATASOURCE_PUBLICATION, DatasourceCatalogueProvider, + DatasourceArtifactBackendProvider, DatasourceDescriptor, DatasourceField, DatasourceLifecycleProvider, @@ -20,6 +22,7 @@ from govoplan_core.core.datasources import ( DatasourceReadResult, DatasourceStage, datasource_catalogue, + datasource_artifact_backend_provider, datasource_lifecycle, datasource_origins, datasource_publication, @@ -155,6 +158,9 @@ class _Provider: truncated=False, ) + def artifact_backends(self): + return () + class DatasourceContractTests(unittest.TestCase): def test_capabilities_are_runtime_checkable_and_resolved_without_modules(self) -> None: @@ -163,6 +169,7 @@ class DatasourceContractTests(unittest.TestCase): self.assertIsInstance(provider, DatasourceLifecycleProvider) self.assertIsInstance(provider, DatasourcePublicationProvider) self.assertIsInstance(provider, DatasourceOriginProvider) + self.assertIsInstance(provider, DatasourceArtifactBackendProvider) registry = PlatformRegistry() registry.register( ModuleManifest( @@ -174,6 +181,7 @@ class DatasourceContractTests(unittest.TestCase): CAPABILITY_DATASOURCE_LIFECYCLE: lambda context: provider, CAPABILITY_DATASOURCE_PUBLICATION: lambda context: provider, CAPABILITY_DATASOURCE_ORIGINS: lambda context: provider, + CAPABILITY_DATASOURCE_ARTIFACT_BACKENDS: lambda context: provider, }, ) ) @@ -183,6 +191,7 @@ class DatasourceContractTests(unittest.TestCase): self.assertIs(provider, datasource_lifecycle(registry)) self.assertIs(provider, datasource_publication(registry)) self.assertIs(provider, datasource_origins(registry)) + self.assertIs(provider, datasource_artifact_backend_provider(registry)) self.assertIsNone(datasource_catalogue(PlatformRegistry())) def test_descriptor_distinguishes_mode_kind_shape_and_materialization(self) -> None: