Preserve connector source runtime evidence

This commit is contained in:
2026-08-04 12:05:35 +02:00
parent 2d93ad07d2
commit 067c8273b6
8 changed files with 202 additions and 4 deletions
+3
View File
@@ -19,6 +19,9 @@ Datasource contracts rather than connector implementations.
The first executable slice supports tabular static uploads, connector-backed The first executable slice supports tabular static uploads, connector-backed
live and cached sources, staging and promotion, refresh, immutable snapshots, live and cached sources, staging and promotion, refresh, immutable snapshots,
explicit frozen states, previews, retirement, and atomic producer publication. 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 Producer modules can append a bounded tabular result or create a new static
datasource through an idempotent capability. The publication ledger retains the datasource through an idempotent capability. The publication ledger retains the
producer run, output materialization, provenance, and replay identity. producer run, output materialization, provenance, and replay identity.
+3 -1
View File
@@ -347,7 +347,9 @@ manifest = ModuleManifest(
"Datasources owns data identity, provenance, lifecycle, read semantics, " "Datasources owns data identity, provenance, lifecycle, read semantics, "
"and the typed governance catalogue for authority, purpose, quality, " "and the typed governance catalogue for authority, purpose, quality, "
"freshness, classification, correction, and dependent services, flows, " "freshness, classification, correction, and dependent services, flows, "
"reports, controls, and decisions. Governance metadata visibility does not " "reports, controls, and decisions. It preserves origin source mode, "
"structured health, declared pushdown, and effective row, byte, and time "
"limits for live previews. Governance metadata visibility does not "
"grant access to protected rows." "grant access to protected rows."
), ),
layer="available", layer="available",
@@ -31,9 +31,12 @@ from govoplan_datasources.backend.schemas import (
DatasourceListResponse, DatasourceListResponse,
DatasourceMaterializationListResponse, DatasourceMaterializationListResponse,
DatasourceMaterializationResponse, DatasourceMaterializationResponse,
DatasourceOriginHealthResponse,
DatasourceOriginListResponse, DatasourceOriginListResponse,
DatasourceOriginPushdownResponse,
DatasourceOriginRegisterRequest, DatasourceOriginRegisterRequest,
DatasourceOriginResponse, DatasourceOriginResponse,
DatasourcePreviewDiagnosticResponse,
DatasourcePreviewResponse, DatasourcePreviewResponse,
DatasourceResponse, DatasourceResponse,
DatasourceRetireResponse, DatasourceRetireResponse,
@@ -383,6 +386,8 @@ def api_preview_datasource(
offset: int = Query(default=0, ge=0), offset: int = Query(default=0, ge=0),
consistency: str = Query(default="current", pattern="^(current|live|frozen)$"), consistency: str = Query(default="current", pattern="^(current|live|frozen)$"),
materialization_ref: str | None = Query(default=None, max_length=120), materialization_ref: str | None = Query(default=None, max_length=120),
max_bytes: int = Query(default=1_000_000, ge=1_024, le=20_000_000),
timeout_ms: int = Query(default=2_000, ge=100, le=10_000),
session: Session = Depends(get_session), session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(get_api_principal), principal: ApiPrincipal = Depends(get_api_principal),
) -> DatasourcePreviewResponse: ) -> DatasourcePreviewResponse:
@@ -397,6 +402,8 @@ def api_preview_datasource(
consistency=consistency, # type: ignore[arg-type] consistency=consistency, # type: ignore[arg-type]
limit=limit, limit=limit,
offset=offset, offset=offset,
max_bytes=max_bytes,
timeout_ms=timeout_ms,
), ),
) )
except DatasourceError as exc: except DatasourceError as exc:
@@ -411,6 +418,20 @@ def api_preview_datasource(
if result.materialization if result.materialization
else None else None
), ),
returned_bytes=result.returned_bytes,
elapsed_ms=result.elapsed_ms,
effective_row_limit=result.effective_row_limit,
effective_byte_limit=result.effective_byte_limit,
effective_timeout_ms=result.effective_timeout_ms,
diagnostics=[
DatasourcePreviewDiagnosticResponse(
severity=item.severity,
code=item.code,
message=item.message,
details=dict(item.details),
)
for item in result.diagnostics
],
) )
@@ -662,6 +683,23 @@ def _origin_response(item: DatasourceOrigin) -> DatasourceOriginResponse:
updated_at=item.updated_at.isoformat() if item.updated_at else None, updated_at=item.updated_at.isoformat() if item.updated_at else None,
capabilities=list(item.capabilities), capabilities=list(item.capabilities),
metadata=dict(item.metadata), metadata=dict(item.metadata),
source_mode=item.source_mode,
pushdown=DatasourceOriginPushdownResponse(
projections=item.pushdown.projections,
pagination=item.pushdown.pagination,
filters=list(item.pushdown.filters),
aggregations=list(item.pushdown.aggregations),
sorting=list(item.pushdown.sorting),
),
health=DatasourceOriginHealthResponse(
status=item.health.status,
code=item.health.code,
summary=item.health.summary,
checked_at=(
item.health.checked_at.isoformat() if item.health.checked_at else None
),
details=dict(item.health.details),
),
) )
@@ -17,6 +17,7 @@ DatasourceKindValue = Literal[
"custom", "custom",
] ]
DatasourceShapeValue = Literal["tabular", "document", "binary", "directory", "stream"] DatasourceShapeValue = Literal["tabular", "document", "binary", "directory", "stream"]
TabularSourceModeValue = Literal["live", "cached", "file_backed", "static"]
SourceAuthorityModeValue = Literal[ SourceAuthorityModeValue = Literal[
"native_authoritative", "native_authoritative",
"external_authoritative", "external_authoritative",
@@ -224,6 +225,29 @@ class DatasourceStagePromoteResponse(BaseModel):
materialization: DatasourceMaterializationResponse materialization: DatasourceMaterializationResponse
class DatasourceOriginPushdownResponse(BaseModel):
projections: bool
pagination: bool
filters: list[str]
aggregations: list[str]
sorting: list[str]
class DatasourceOriginHealthResponse(BaseModel):
status: Literal["healthy", "warning", "error", "unknown"]
code: str
summary: str
checked_at: str | None
details: dict[str, Any]
class DatasourcePreviewDiagnosticResponse(BaseModel):
severity: Literal["info", "warning", "error"]
code: str
message: str
details: dict[str, Any]
class DatasourceOriginResponse(BaseModel): class DatasourceOriginResponse(BaseModel):
ref: str ref: str
source_name: str source_name: str
@@ -241,6 +265,9 @@ class DatasourceOriginResponse(BaseModel):
updated_at: str | None updated_at: str | None
capabilities: list[str] capabilities: list[str]
metadata: dict[str, Any] metadata: dict[str, Any]
source_mode: TabularSourceModeValue
pushdown: DatasourceOriginPushdownResponse
health: DatasourceOriginHealthResponse
class DatasourceOriginListResponse(BaseModel): class DatasourceOriginListResponse(BaseModel):
@@ -267,6 +294,12 @@ class DatasourcePreviewResponse(BaseModel):
total_rows: int total_rows: int
truncated: bool truncated: bool
materialization: DatasourceMaterializationResponse | None materialization: DatasourceMaterializationResponse | None
returned_bytes: int
elapsed_ms: int
effective_row_limit: int
effective_byte_limit: int
effective_timeout_ms: int
diagnostics: list[DatasourcePreviewDiagnosticResponse]
class DatasourceFreezeRequest(BaseModel): class DatasourceFreezeRequest(BaseModel):
+43 -1
View File
@@ -665,7 +665,7 @@ class SqlDatasourceProvider:
"origin_provider": origin.provider, "origin_provider": origin.provider,
"registered_at": utcnow().isoformat(), "registered_at": utcnow().isoformat(),
}, },
metadata_=dict(origin.metadata), metadata_=_origin_metadata(origin),
created_by=_actor_id(api_principal), created_by=_actor_id(api_principal),
updated_by=_actor_id(api_principal), updated_by=_actor_id(api_principal),
) )
@@ -856,6 +856,8 @@ class SqlDatasourceProvider:
offset=offset, offset=offset,
columns=columns, columns=columns,
expected_fingerprint=request.expected_fingerprint, expected_fingerprint=request.expected_fingerprint,
max_bytes=request.max_bytes,
timeout_ms=request.timeout_ms,
), ),
) )
except DatasourceError: except DatasourceError:
@@ -872,12 +874,19 @@ class SqlDatasourceProvider:
row_count=result.origin.row_count, row_count=result.origin.row_count,
byte_count=result.origin.byte_count, byte_count=result.origin.byte_count,
updated_at=result.origin.updated_at, updated_at=result.origin.updated_at,
metadata=_origin_metadata(result.origin, current=item.metadata_),
) )
return DatasourceReadResult( return DatasourceReadResult(
datasource=descriptor, datasource=descriptor,
rows=result.rows, rows=result.rows,
total_rows=result.total_rows, total_rows=result.total_rows,
truncated=result.truncated, truncated=result.truncated,
returned_bytes=result.returned_bytes,
elapsed_ms=result.elapsed_ms,
effective_row_limit=result.effective_row_limit,
effective_byte_limit=result.effective_byte_limit,
effective_timeout_ms=result.effective_timeout_ms,
diagnostics=result.diagnostics,
) )
def _required_origin( def _required_origin(
@@ -964,6 +973,7 @@ class SqlDatasourceProvider:
frozen_label: str | None = None, frozen_label: str | None = None,
set_current: bool, set_current: bool,
) -> DatasourceMaterializationRecord: ) -> DatasourceMaterializationRecord:
item.metadata_ = _origin_metadata(origin, current=item.metadata_)
normalized = normalize_rows(rows) normalized = normalize_rows(rows)
schema = infer_schema(normalized) or origin.schema schema = infer_schema(normalized) or origin.schema
fingerprint = fingerprint_rows(normalized, schema) fingerprint = fingerprint_rows(normalized, schema)
@@ -992,6 +1002,38 @@ class SqlDatasourceProvider:
) )
def _origin_metadata(
origin: DatasourceOrigin,
*,
current: Mapping[str, object] | None = None,
) -> dict[str, object]:
return {
**dict(current or {}),
**dict(origin.metadata),
"source_contract": {
"source_mode": origin.source_mode,
"pushdown": {
"projections": origin.pushdown.projections,
"pagination": origin.pushdown.pagination,
"filters": list(origin.pushdown.filters),
"aggregations": list(origin.pushdown.aggregations),
"sorting": list(origin.pushdown.sorting),
},
"health": {
"status": origin.health.status,
"code": origin.health.code,
"summary": origin.health.summary,
"checked_at": (
origin.health.checked_at.isoformat()
if origin.health.checked_at
else None
),
"details": dict(origin.health.details),
},
},
}
def _append_materialization( def _append_materialization(
session: Session, session: Session,
*, *,
+39
View File
@@ -22,6 +22,11 @@ from govoplan_core.core.datasources import (
DatasourceUnavailableError, DatasourceUnavailableError,
DatasourceValidationError, DatasourceValidationError,
) )
from govoplan_core.core.tabular_sources import (
TabularPreviewDiagnostic,
TabularPushdown,
TabularSourceHealth,
)
from govoplan_core.db.base import Base, utcnow from govoplan_core.db.base import Base, utcnow
from govoplan_datasources.backend.db.models import ( from govoplan_datasources.backend.db.models import (
DatasourceGovernanceReferenceRecord, DatasourceGovernanceReferenceRecord,
@@ -88,6 +93,13 @@ class FakeOriginProvider:
fingerprint=f"version-{len(self.rows)}-{self.rows[-1]['name']}", fingerprint=f"version-{len(self.rows)}-{self.rows[-1]['name']}",
row_count=len(self.rows), row_count=len(self.rows),
updated_at=utcnow(), updated_at=utcnow(),
source_mode="cached",
pushdown=TabularPushdown(projections=True, pagination=True),
health=TabularSourceHealth(
status="healthy",
code="snapshot.ready",
summary="The immutable snapshot is ready.",
),
) )
def list_origins( def list_origins(
@@ -126,6 +138,18 @@ class FakeOriginProvider:
rows=tuple(dict(row) for row in rows), rows=tuple(dict(row) for row in rows),
total_rows=len(self.rows), total_rows=len(self.rows),
truncated=request.offset + len(rows) < len(self.rows), truncated=request.offset + len(rows) < len(self.rows),
returned_bytes=64,
elapsed_ms=4,
effective_row_limit=request.limit,
effective_byte_limit=request.max_bytes,
effective_timeout_ms=request.timeout_ms,
diagnostics=(
TabularPreviewDiagnostic(
severity="info",
code="preview.complete",
message="The bounded preview completed.",
),
),
) )
@@ -513,10 +537,25 @@ class DatasourceLifecycleTests(unittest.TestCase):
) )
self.assertEqual(2, live_read.total_rows) self.assertEqual(2, live_read.total_rows)
self.assertEqual(64, live_read.returned_bytes)
self.assertEqual("preview.complete", live_read.diagnostics[0].code)
self.assertEqual(
"cached",
live_read.datasource.metadata["source_contract"]["source_mode"],
)
self.assertTrue(
live_read.datasource.metadata["source_contract"]["pushdown"][
"projections"
]
)
self.assertEqual(1, cached_before.total_rows) self.assertEqual(1, cached_before.total_rows)
self.assertEqual(2, cached_live.total_rows) self.assertEqual(2, cached_live.total_rows)
self.assertEqual(2, cached_after.total_rows) self.assertEqual(2, cached_after.total_rows)
self.assertEqual(materialization.ref, refreshed.current_materialization_ref) self.assertEqual(materialization.ref, refreshed.current_materialization_ref)
self.assertEqual(
"healthy",
refreshed.metadata["source_contract"]["health"]["status"],
)
def test_tenant_and_scope_isolation(self) -> None: def test_tenant_and_scope_isolation(self) -> None:
stage = self.provider.create_stage( stage = self.provider.create_stage(
+26
View File
@@ -164,6 +164,21 @@ export type DatasourceOrigin = {
updated_at?: string | null; updated_at?: string | null;
capabilities: string[]; capabilities: string[];
metadata: Record<string, unknown>; metadata: Record<string, unknown>;
source_mode: "live" | "cached" | "file_backed" | "static";
pushdown: {
projections: boolean;
pagination: boolean;
filters: string[];
aggregations: string[];
sorting: string[];
};
health: {
status: "healthy" | "warning" | "error" | "unknown";
code: string;
summary: string;
checked_at?: string | null;
details: Record<string, unknown>;
};
}; };
export type DatasourcePreview = { export type DatasourcePreview = {
@@ -172,6 +187,17 @@ export type DatasourcePreview = {
total_rows: number; total_rows: number;
truncated: boolean; truncated: boolean;
materialization?: DatasourceMaterialization | null; materialization?: DatasourceMaterialization | null;
returned_bytes: number;
elapsed_ms: number;
effective_row_limit: number;
effective_byte_limit: number;
effective_timeout_ms: number;
diagnostics: Array<{
severity: "info" | "warning" | "error";
code: string;
message: string;
details: Record<string, unknown>;
}>;
}; };
export async function listDatasources( export async function listDatasources(
@@ -834,7 +834,8 @@ function OriginDetail({ origin }: { origin: DatasourceOrigin }) {
<> <>
<div className="datasources-metrics"> <div className="datasources-metrics">
<Metric label="Provider" value={origin.provider} /> <Metric label="Provider" value={origin.provider} />
<Metric label="Kind" value={origin.kind} /> <Metric label="Mode" value={readableToken(origin.source_mode)} />
<Metric label="Health" value={readableToken(origin.health.status)} />
<Metric label="Rows" value={formatNumber(origin.row_count)} /> <Metric label="Rows" value={formatNumber(origin.row_count)} />
<Metric label="Fields" value={String(origin.schema.length)} /> <Metric label="Fields" value={String(origin.schema.length)} />
<Metric label="Updated" value={formatDate(origin.updated_at)} /> <Metric label="Updated" value={formatDate(origin.updated_at)} />
@@ -849,8 +850,11 @@ function OriginDetail({ origin }: { origin: DatasourceOrigin }) {
</div> </div>
<div className="datasources-key-values"> <div className="datasources-key-values">
<span><small>Origin reference</small><strong>{origin.ref}</strong></span> <span><small>Origin reference</small><strong>{origin.ref}</strong></span>
<span><small>Shape</small><strong>{origin.shape}</strong></span> <span><small>Kind</small><strong>{readableToken(origin.kind)}</strong></span>
<span><small>Shape</small><strong>{readableToken(origin.shape)}</strong></span>
<span><small>Fingerprint</small><strong>{shortFingerprint(origin.fingerprint)}</strong></span> <span><small>Fingerprint</small><strong>{shortFingerprint(origin.fingerprint)}</strong></span>
<span><small>Health</small><strong>{origin.health.summary}</strong></span>
<span><small>Pushdown</small><strong>{pushdownSummary(origin)}</strong></span>
</div> </div>
</section> </section>
<section className="datasources-detail-section"> <section className="datasources-detail-section">
@@ -864,6 +868,17 @@ function OriginDetail({ origin }: { origin: DatasourceOrigin }) {
); );
} }
function pushdownSummary(origin: DatasourceOrigin): string {
const capabilities = [
origin.pushdown.projections ? "projections" : "",
origin.pushdown.pagination ? "pagination" : "",
...origin.pushdown.filters.map((item) => `filter:${item}`),
...origin.pushdown.aggregations.map((item) => `aggregate:${item}`),
...origin.pushdown.sorting.map((item) => `sort:${item}`)
].filter(Boolean);
return capabilities.join(", ") || "None declared";
}
function GovernanceDialog({ function GovernanceDialog({
open, open,
settings, settings,