feat(core): define durable datasource artifact contract

This commit is contained in:
2026-08-21 17:36:05 +02:00
parent 925dc33696
commit fbea74a74b
2 changed files with 100 additions and 2 deletions
+91 -2
View File
@@ -23,6 +23,7 @@ CAPABILITY_DATASOURCE_CATALOGUE = "datasources.catalogue"
CAPABILITY_DATASOURCE_LIFECYCLE = "datasources.lifecycle" CAPABILITY_DATASOURCE_LIFECYCLE = "datasources.lifecycle"
CAPABILITY_DATASOURCE_PUBLICATION = "datasources.publication" CAPABILITY_DATASOURCE_PUBLICATION = "datasources.publication"
CAPABILITY_DATASOURCE_ORIGINS = "connectors.datasourceOrigins" CAPABILITY_DATASOURCE_ORIGINS = "connectors.datasourceOrigins"
CAPABILITY_DATASOURCE_ARTIFACT_BACKENDS = "datasources.artifactBackends"
DatasourceMode = Literal["live", "cached", "static"] DatasourceMode = Literal["live", "cached", "static"]
DatasourceKind = Literal[ DatasourceKind = Literal[
@@ -37,6 +38,11 @@ DatasourceKind = Literal[
] ]
DatasourceShape = Literal["tabular", "document", "binary", "directory", "stream"] DatasourceShape = Literal["tabular", "document", "binary", "directory", "stream"]
DatasourceConsistency = Literal["current", "live", "frozen"] DatasourceConsistency = Literal["current", "live", "frozen"]
DatasourcePublicationStatus = Literal[
"published",
"published_with_warnings",
"review_required",
]
class DatasourceError(ValueError): class DatasourceError(ValueError):
@@ -309,12 +315,73 @@ class DatasourceStageInput:
governance: DatasourceGovernance | None = None 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) @dataclass(frozen=True, slots=True)
class DatasourcePublicationRequest: class DatasourcePublicationRequest:
producer_module: str producer_module: str
producer_run_ref: str producer_run_ref: str
idempotency_key: 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 target_datasource_ref: str | None = None
name: str | None = None name: str | None = None
source_name: str | None = None source_name: str | None = None
@@ -331,7 +398,7 @@ class DatasourcePublicationRequest:
@dataclass(frozen=True, slots=True) @dataclass(frozen=True, slots=True)
class DatasourcePublicationResult: class DatasourcePublicationResult:
ref: str ref: str
status: str status: DatasourcePublicationStatus
datasource: DatasourceDescriptor datasource: DatasourceDescriptor
materialization: DatasourceMaterialization materialization: DatasourceMaterialization
replayed: bool = False replayed: bool = False
@@ -564,6 +631,17 @@ def datasource_catalogue(registry: object | None) -> DatasourceCatalogueProvider
return capability if isinstance(capability, DatasourceCatalogueProvider) else None 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: def datasource_lifecycle(registry: object | None) -> DatasourceLifecycleProvider | None:
capability = _capability(registry, CAPABILITY_DATASOURCE_LIFECYCLE) capability = _capability(registry, CAPABILITY_DATASOURCE_LIFECYCLE)
return capability if isinstance(capability, DatasourceLifecycleProvider) else None return capability if isinstance(capability, DatasourceLifecycleProvider) else None
@@ -613,9 +691,14 @@ def _governance_mapping(value: object) -> Mapping[str, object]:
__all__ = [ __all__ = [
"CAPABILITY_DATASOURCE_CATALOGUE", "CAPABILITY_DATASOURCE_CATALOGUE",
"CAPABILITY_DATASOURCE_ARTIFACT_BACKENDS",
"CAPABILITY_DATASOURCE_LIFECYCLE", "CAPABILITY_DATASOURCE_LIFECYCLE",
"CAPABILITY_DATASOURCE_ORIGINS", "CAPABILITY_DATASOURCE_ORIGINS",
"CAPABILITY_DATASOURCE_PUBLICATION",
"DatasourceAccessError", "DatasourceAccessError",
"DatasourceArtifactReference",
"DatasourceArtifactBackend",
"DatasourceArtifactBackendProvider",
"DatasourceCatalogueProvider", "DatasourceCatalogueProvider",
"DatasourceConsistency", "DatasourceConsistency",
"DatasourceDescriptor", "DatasourceDescriptor",
@@ -631,6 +714,10 @@ __all__ = [
"DatasourceOriginProvider", "DatasourceOriginProvider",
"DatasourceOriginReadRequest", "DatasourceOriginReadRequest",
"DatasourceOriginReadResult", "DatasourceOriginReadResult",
"DatasourcePublicationProvider",
"DatasourcePublicationRequest",
"DatasourcePublicationResult",
"DatasourcePublicationStatus",
"DatasourceReadRequest", "DatasourceReadRequest",
"DatasourceReadResult", "DatasourceReadResult",
"DatasourceShape", "DatasourceShape",
@@ -639,6 +726,8 @@ __all__ = [
"DatasourceUnavailableError", "DatasourceUnavailableError",
"DatasourceValidationError", "DatasourceValidationError",
"datasource_catalogue", "datasource_catalogue",
"datasource_artifact_backend_provider",
"datasource_lifecycle", "datasource_lifecycle",
"datasource_origins", "datasource_origins",
"datasource_publication",
] ]
+9
View File
@@ -3,11 +3,13 @@ from __future__ import annotations
import unittest import unittest
from govoplan_core.core.datasources import ( from govoplan_core.core.datasources import (
CAPABILITY_DATASOURCE_ARTIFACT_BACKENDS,
CAPABILITY_DATASOURCE_CATALOGUE, CAPABILITY_DATASOURCE_CATALOGUE,
CAPABILITY_DATASOURCE_LIFECYCLE, CAPABILITY_DATASOURCE_LIFECYCLE,
CAPABILITY_DATASOURCE_ORIGINS, CAPABILITY_DATASOURCE_ORIGINS,
CAPABILITY_DATASOURCE_PUBLICATION, CAPABILITY_DATASOURCE_PUBLICATION,
DatasourceCatalogueProvider, DatasourceCatalogueProvider,
DatasourceArtifactBackendProvider,
DatasourceDescriptor, DatasourceDescriptor,
DatasourceField, DatasourceField,
DatasourceLifecycleProvider, DatasourceLifecycleProvider,
@@ -20,6 +22,7 @@ from govoplan_core.core.datasources import (
DatasourceReadResult, DatasourceReadResult,
DatasourceStage, DatasourceStage,
datasource_catalogue, datasource_catalogue,
datasource_artifact_backend_provider,
datasource_lifecycle, datasource_lifecycle,
datasource_origins, datasource_origins,
datasource_publication, datasource_publication,
@@ -155,6 +158,9 @@ class _Provider:
truncated=False, truncated=False,
) )
def artifact_backends(self):
return ()
class DatasourceContractTests(unittest.TestCase): class DatasourceContractTests(unittest.TestCase):
def test_capabilities_are_runtime_checkable_and_resolved_without_modules(self) -> None: 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, DatasourceLifecycleProvider)
self.assertIsInstance(provider, DatasourcePublicationProvider) self.assertIsInstance(provider, DatasourcePublicationProvider)
self.assertIsInstance(provider, DatasourceOriginProvider) self.assertIsInstance(provider, DatasourceOriginProvider)
self.assertIsInstance(provider, DatasourceArtifactBackendProvider)
registry = PlatformRegistry() registry = PlatformRegistry()
registry.register( registry.register(
ModuleManifest( ModuleManifest(
@@ -174,6 +181,7 @@ class DatasourceContractTests(unittest.TestCase):
CAPABILITY_DATASOURCE_LIFECYCLE: lambda context: provider, CAPABILITY_DATASOURCE_LIFECYCLE: lambda context: provider,
CAPABILITY_DATASOURCE_PUBLICATION: lambda context: provider, CAPABILITY_DATASOURCE_PUBLICATION: lambda context: provider,
CAPABILITY_DATASOURCE_ORIGINS: 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_lifecycle(registry))
self.assertIs(provider, datasource_publication(registry)) self.assertIs(provider, datasource_publication(registry))
self.assertIs(provider, datasource_origins(registry)) self.assertIs(provider, datasource_origins(registry))
self.assertIs(provider, datasource_artifact_backend_provider(registry))
self.assertIsNone(datasource_catalogue(PlatformRegistry())) self.assertIsNone(datasource_catalogue(PlatformRegistry()))
def test_descriptor_distinguishes_mode_kind_shape_and_materialization(self) -> None: def test_descriptor_distinguishes_mode_kind_shape_and_materialization(self) -> None: