From 5d7ef60c4ae74b63eed1be58e4b2649ec8734720 Mon Sep 17 00:00:00 2001 From: Ning Zhou Date: Mon, 17 Aug 2026 21:32:52 +0400 Subject: [PATCH 1/2] feat: add public lightweight dataops run --- .github/workflows/cd-staging.yml | 2 + .github/workflows/ci.yml | 2 + .../099_synchronous_success_profile.sql | 226 +++++ data_agent/platform_gateway.py | 41 +- data_agent/public_dataops_run.py | 823 ++++++++++++++++++ data_agent/test_public_dataops_run.py | 294 +++++++ .../test_public_dataops_run_postgres.py | 308 +++++++ .../adr-082-public-lightweight-dataops-run.md | 83 ++ .../public-dataops-run-2026-08-17.json | 97 +++ docs/roadmap.md | 7 +- docs/system-of-record-matrix-2026-07-24.md | 12 +- 11 files changed, 1882 insertions(+), 13 deletions(-) create mode 100644 data_agent/migrations/099_synchronous_success_profile.sql create mode 100644 data_agent/public_dataops_run.py create mode 100644 data_agent/test_public_dataops_run.py create mode 100644 data_agent/test_public_dataops_run_postgres.py create mode 100644 docs/architecture-decisions/adr-082-public-lightweight-dataops-run.md create mode 100644 docs/evidence/public-dataops-run-2026-08-17.json diff --git a/.github/workflows/cd-staging.yml b/.github/workflows/cd-staging.yml index 824bb200..d02ef110 100644 --- a/.github/workflows/cd-staging.yml +++ b/.github/workflows/cd-staging.yml @@ -173,6 +173,8 @@ jobs: data_agent/test_chongqing_protected_admission_workflow.py \ data_agent/test_public_source_landing.py \ data_agent/test_public_source_landing_postgres.py \ + data_agent/test_public_dataops_run.py \ + data_agent/test_public_dataops_run_postgres.py \ data_agent/test_dolphinscheduler_adapter.py \ data_agent/test_dolphinscheduler_command_consumer.py \ data_agent/test_dolphinscheduler_command_worker.py \ diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 003d9ce4..0f868407 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -182,6 +182,8 @@ jobs: data_agent/test_chongqing_protected_admission_workflow.py \ data_agent/test_public_source_landing.py \ data_agent/test_public_source_landing_postgres.py \ + data_agent/test_public_dataops_run.py \ + data_agent/test_public_dataops_run_postgres.py \ data_agent/test_dolphinscheduler_adapter.py \ data_agent/test_dolphinscheduler_command_consumer.py \ data_agent/test_dolphinscheduler_command_worker.py \ diff --git a/data_agent/migrations/099_synchronous_success_profile.sql b/data_agent/migrations/099_synchronous_success_profile.sql new file mode 100644 index 00000000..904de5e0 --- /dev/null +++ b/data_agent/migrations/099_synchronous_success_profile.sql @@ -0,0 +1,226 @@ +-- 099: Evidence-gated success for the synchronous public DataOps profile. +-- +-- Migration 096 remains checksum-frozen and continues to own the +-- DolphinScheduler verdict. This separate entry point accepts only the local +-- inline observation profile while preserving the same output, quality, +-- lineage, actor, tenant, state, fingerprint, and replay checks. + +CREATE OR REPLACE FUNCTION gda_control.finalize_synchronous_platform_run_success( + p_tenant_id TEXT, + p_run_id UUID, + p_expected_state_version INTEGER, + p_actor_subject TEXT, + p_reason TEXT, + p_details JSONB +) +RETURNS INTEGER +LANGUAGE plpgsql +SECURITY DEFINER +SET search_path = pg_catalog, gda_control +SET row_security = on +AS $$ +DECLARE + v_run gda_control.platform_run%ROWTYPE; + v_event gda_control.platform_run_event%ROWTYPE; + v_observation_id UUID; + v_output_artifact_id UUID; + v_quality_result_id UUID; + v_lineage_event_id UUID; + v_expected_evidence_sha256 TEXT; + v_output gda_control.artifact%ROWTYPE; + v_quality gda_control.quality_result%ROWTYPE; +BEGIN + IF gda_control.current_tenant() IS NULL + OR p_tenant_id IS DISTINCT FROM gda_control.current_tenant() THEN + RAISE EXCEPTION 'platform run tenant context is missing or mismatched' + USING ERRCODE = '42501'; + END IF; + IF NULLIF(btrim(p_actor_subject), '') IS NULL + OR NULLIF(btrim(p_reason), '') IS NULL THEN + RAISE EXCEPTION 'finalization actor and reason are required' + USING ERRCODE = '22023'; + END IF; + IF jsonb_typeof(p_details) <> 'object' + OR p_details->>'schema' <> 'gda.run_success_evidence.v1' + OR p_details->>'tenant_id' IS DISTINCT FROM p_tenant_id + OR p_details->>'run_id' IS DISTINCT FROM p_run_id::text + OR COALESCE(p_details->>'evidence_sha256', '') + !~ '^[0-9a-f]{64}$' THEN + RAISE EXCEPTION 'success evidence envelope is invalid' + USING ERRCODE = '22023'; + END IF; + + SELECT * INTO v_run + FROM gda_control.platform_run + WHERE tenant_id = p_tenant_id AND run_id = p_run_id + FOR UPDATE; + IF NOT FOUND THEN + RAISE EXCEPTION 'platform run % not found', p_run_id + USING ERRCODE = 'P0002'; + END IF; + + IF v_run.orchestration_class <> 'synchronous' THEN + RAISE EXCEPTION 'synchronous success profile requires synchronous Run' + USING ERRCODE = '23514'; + END IF; + + IF v_run.status = 'succeeded' THEN + SELECT * INTO v_event + FROM gda_control.platform_run_event + WHERE tenant_id = p_tenant_id + AND run_id = p_run_id + AND sequence_no = v_run.state_version; + IF FOUND + AND v_event.to_status = 'succeeded' + AND v_event.actor_subject = p_actor_subject + AND v_event.reason = p_reason + AND v_event.details = p_details THEN + RETURN v_run.state_version; + END IF; + RAISE EXCEPTION 'successful Run has a different terminal verdict' + USING ERRCODE = '40001'; + END IF; + + IF v_run.state_version <> p_expected_state_version THEN + RAISE EXCEPTION 'platform run state version conflict: expected %, actual %', + p_expected_state_version, v_run.state_version + USING ERRCODE = '40001'; + END IF; + IF v_run.status <> 'running' THEN + RAISE EXCEPTION 'synchronous success finalization requires running Run' + USING ERRCODE = '23514'; + END IF; + IF p_actor_subject IS DISTINCT FROM concat( + v_run.subject_context->>'subject_type', + ':', + v_run.subject_context->>'subject_id' + ) THEN + RAISE EXCEPTION 'finalization actor does not match Run workload' + USING ERRCODE = '42501'; + END IF; + + BEGIN + v_observation_id := (p_details->>'attempt_observation_id')::uuid; + v_output_artifact_id := (p_details->>'output_artifact_id')::uuid; + v_quality_result_id := (p_details->>'quality_result_id')::uuid; + v_lineage_event_id := (p_details->>'lineage_event_id')::uuid; + EXCEPTION WHEN invalid_text_representation OR null_value_not_allowed THEN + RAISE EXCEPTION 'success evidence identifiers must be UUIDs' + USING ERRCODE = '22023'; + END; + + v_expected_evidence_sha256 := encode( + sha256( + convert_to( + '{"attempt_observation_id":' + || to_json(v_observation_id::text)::text + || ',"lineage_event_id":' + || to_json(v_lineage_event_id::text)::text + || ',"output_artifact_id":' + || to_json(v_output_artifact_id::text)::text + || ',"quality_result_id":' + || to_json(v_quality_result_id::text)::text + || ',"run_id":' + || to_json(p_run_id::text)::text + || ',"tenant_id":' + || to_json(p_tenant_id)::text + || '}', + 'UTF8' + ) + ), + 'hex' + ); + IF p_details->>'evidence_sha256' + IS DISTINCT FROM v_expected_evidence_sha256 THEN + RAISE EXCEPTION 'success evidence fingerprint does not match its bindings' + USING ERRCODE = '22023'; + END IF; + + PERFORM 1 + FROM gda_control.framework_attempt_observation + WHERE tenant_id = p_tenant_id + AND observation_id = v_observation_id + AND run_id = p_run_id + AND framework_kind = 'legacy' + AND lower(observed_state) = 'success' + AND evidence->>'schema' = 'gda.public_dataops_attempt.v1' + AND evidence->>'execution_mode' = 'local_inline'; + IF NOT FOUND THEN + RAISE EXCEPTION 'local inline success observation was not found' + USING ERRCODE = '23514'; + END IF; + + SELECT artifact.* INTO v_output + FROM gda_control.artifact AS artifact + JOIN gda_control.resource_version AS version + ON version.tenant_id = artifact.tenant_id + AND version.resource_version_id = artifact.resource_version_id + AND version.content_sha256 = artifact.content_sha256 + WHERE artifact.tenant_id = p_tenant_id + AND artifact.artifact_id = v_output_artifact_id + AND artifact.run_id = p_run_id + AND artifact.artifact_role = 'output'; + IF NOT FOUND THEN + RAISE EXCEPTION 'content-bound output Artifact was not found' + USING ERRCODE = '23514'; + END IF; + + SELECT quality.* INTO v_quality + FROM gda_control.quality_result AS quality + JOIN gda_control.artifact AS evidence + ON evidence.tenant_id = quality.tenant_id + AND evidence.artifact_id = quality.evidence_artifact_id + WHERE quality.tenant_id = p_tenant_id + AND quality.quality_result_id = v_quality_result_id + AND quality.run_id = p_run_id + AND quality.resource_version_id = v_output.resource_version_id + AND quality.verdict = 'passed' + AND quality.evaluated_by <> p_actor_subject + AND quality.evaluated_at >= v_output.created_at + AND evidence.artifact_role = 'evidence' + AND evidence.run_id = p_run_id + AND evidence.resource_version_id = v_output.resource_version_id + AND evidence.created_by = quality.evaluated_by; + IF NOT FOUND THEN + RAISE EXCEPTION 'independent passed QualityResult was not found' + USING ERRCODE = '23514'; + END IF; + + PERFORM 1 + FROM gda_control.lineage_event AS lineage + WHERE lineage.tenant_id = p_tenant_id + AND lineage.lineage_event_id = v_lineage_event_id + AND lineage.run_id = p_run_id + AND lineage.definition_version_id = v_run.definition_version_id + AND lineage.artifact_id = v_output_artifact_id + AND lineage.target_resource_version_id = v_output.resource_version_id + AND EXISTS ( + SELECT 1 + FROM gda_control.platform_run_input_binding AS input + WHERE input.tenant_id = p_tenant_id + AND input.run_id = p_run_id + AND input.resource_version_id = lineage.source_resource_version_id + ); + IF NOT FOUND THEN + RAISE EXCEPTION 'input-to-output LineageEvent was not found' + USING ERRCODE = '23514'; + END IF; + + RETURN gda_control.apply_platform_run_transition( + p_tenant_id, + p_run_id, + p_expected_state_version, + 'succeeded', + p_actor_subject, + p_reason, + p_details + ); +END; +$$; + +REVOKE ALL ON FUNCTION gda_control.finalize_synchronous_platform_run_success( + text, uuid, integer, text, text, jsonb +) FROM PUBLIC; +GRANT EXECUTE ON FUNCTION gda_control.finalize_synchronous_platform_run_success( + text, uuid, integer, text, text, jsonb +) TO gda_control_gateway; diff --git a/data_agent/platform_gateway.py b/data_agent/platform_gateway.py index 1a33212a..be6b2114 100644 --- a/data_agent/platform_gateway.py +++ b/data_agent/platform_gateway.py @@ -79,6 +79,11 @@ / "migrations" / "096_platform_success_verdict.sql" ) +SYNCHRONOUS_SUCCESS_MIGRATION = ( + Path(__file__).resolve().parent + / "migrations" + / "099_synchronous_success_profile.sql" +) METADATA_FABRIC_BINDING_MIGRATION = ( Path(__file__).resolve().parent / "migrations" @@ -1002,10 +1007,20 @@ def finalize_run_success( **evidence.model_dump(mode="json"), } with self._transaction(evidence.tenant_id) as connection: + run = self._load_run( + connection, evidence.tenant_id, evidence.run_id + ) + if run is None: + raise GatewayNotFoundError("PlatformRun was not found") + finalizer = ( + "gda_control.finalize_synchronous_platform_run_success" + if run.orchestration_class.value == "synchronous" + else "gda_control.finalize_platform_run_success" + ) connection.execute( text( - """ - SELECT gda_control.finalize_platform_run_success( + f""" + SELECT {finalizer}( :tenant_id, :run_id, :expected_state_version, :actor_subject, :reason, CAST(:details AS jsonb) ) @@ -1020,12 +1035,12 @@ def finalize_run_success( "details": _json(details), }, ).scalar_one() - run = self._load_run( + finalized = self._load_run( connection, evidence.tenant_id, evidence.run_id ) - if run is None: + if finalized is None: raise GatewayNotFoundError("PlatformRun was not found") - return run + return finalized @classmethod def _reconcile_command( @@ -1953,6 +1968,7 @@ def build_gateway_report( role_migration: Path | None = None, command_migration: Path | None = None, success_migration: Path | None = None, + synchronous_success_migration: Path | None = None, binding_migration: Path | None = None, lineage_migration: Path | None = None, gateway_source: Path | None = None, @@ -1970,6 +1986,9 @@ def build_gateway_report( "success_migration": ( success_migration or SUCCESS_VERDICT_MIGRATION ).resolve(), + "synchronous_success_migration": ( + synchronous_success_migration or SYNCHRONOUS_SUCCESS_MIGRATION + ).resolve(), "binding_migration": ( binding_migration or METADATA_FABRIC_BINDING_MIGRATION ).resolve(), @@ -2037,6 +2056,17 @@ def build_gateway_report( "GRANT SELECT, INSERT ON gda_control.quality_result", "finalize_platform_run_success", ), + "synchronous_success_migration": ( + "finalize_synchronous_platform_run_success", + "synchronous success profile requires synchronous Run", + "framework_kind = 'legacy'", + "gda.public_dataops_attempt.v1", + "evidence->>'execution_mode' = 'local_inline'", + "content-bound output Artifact was not found", + "independent passed QualityResult was not found", + "input-to-output LineageEvent was not found", + "TO gda_control_gateway", + ), "binding_migration": ( "CREATE TABLE IF NOT EXISTS gda_control.metadata_fabric_binding", "FOREIGN KEY (tenant_id, execution_plan_artifact_id)", @@ -2124,6 +2154,7 @@ def build_gateway_report( forbidden in role_sql or forbidden in texts.get("command_migration", "") or forbidden in texts.get("success_migration", "") + or forbidden in texts.get("synchronous_success_migration", "") or forbidden in texts.get("binding_migration", "") or forbidden in texts.get("lineage_migration", "") ): diff --git a/data_agent/public_dataops_run.py b/data_agent/public_dataops_run.py new file mode 100644 index 00000000..da7c8617 --- /dev/null +++ b/data_agent/public_dataops_run.py @@ -0,0 +1,823 @@ +"""Run a small public Landing through the lightweight DataOps profile.""" + +from __future__ import annotations + +import argparse +import hashlib +import json +import re +import stat +import tempfile +import zipfile +from datetime import UTC, datetime +from pathlib import Path, PurePosixPath +from typing import Annotated, Any +from uuid import NAMESPACE_URL, uuid5 + +from pydantic import Field, field_validator, model_validator + +from .platform_contracts import ( + Artifact, + FrameworkAttemptObservation, + FrozenContract, + LineageEvent, + PlatformDefinitionVersion, + PlatformRun, + QualityResult, + Resource, + ResourceBinding, + ResourceVersion, + RunSuccessEvidence, + SubjectContext, + canonical_json_bytes, + canonical_json_fingerprint, + platform_definition_fingerprint, + quality_result_fingerprint, + run_success_evidence_fingerprint, +) +from .platform_gateway import DefinitionRegistration, PlatformGateway +from .public_source_landing import ( + PublicSourceLandingResult, + _install_immutable_bytes, + verify_public_source_landing, +) + +PUBLIC_DATAOPS_SCHEMA = "gda.public_dataops_run.v1" +DEFINITION_SCHEMA = "gda.public_dataops_definition.v1" +QUALITY_SCHEMA = "gda.public_dataops_quality.v1" +ATTEMPT_SCHEMA = "gda.public_dataops_attempt.v1" +LINEAGE_SCHEMA = "gda.public_dataops_lineage.v1" +QUALITY_RULE_VERSION = "gda://public-open/quality-rule/geojson-v1" +DEFINITION_PUBLISHED_BY = "workload:gda-release" +DEFINITION_PUBLISHED_AT = datetime(2026, 8, 17, tzinfo=UTC) +_DATASET_ID_RE = re.compile(r"^[a-z0-9][a-z0-9._-]{0,79}$") + +DatasetId = Annotated[ + str, + Field(min_length=1, max_length=80, pattern=r"^[a-z0-9][a-z0-9._-]{0,79}$"), +] + + +class PublicDataOpsError(RuntimeError): + """The public lightweight DataOps run cannot proceed safely.""" + + +class PublicDataOpsRequest(FrozenContract): + schema_id = "public_dataops_request" + + executor: str + quality_evaluator: str + output_dataset_id: DatasetId + executed_at: datetime + min_feature_count: int = Field(default=1, ge=1, le=10_000_000) + + @field_validator("executor", "quality_evaluator") + @classmethod + def _workload_identity(cls, value: str) -> str: + if not value.startswith("workload:") or not value.removeprefix("workload:"): + raise ValueError("DataOps actors must use workload identities") + return value + + @field_validator("executed_at") + @classmethod + def _utc_executed_at(cls, value: datetime) -> datetime: + if value.tzinfo is None or value.utcoffset() is None: + raise ValueError("executed_at must include a timezone") + return value.astimezone(UTC) + + @model_validator(mode="after") + def _independent_quality(self) -> PublicDataOpsRequest: + if self.executor == self.quality_evaluator: + raise ValueError("quality evaluator must be independent from executor") + return self + + +class PublicDataOpsResult(FrozenContract): + schema_id = PUBLIC_DATAOPS_SCHEMA + + landing: PublicSourceLandingResult + definition_registration: DefinitionRegistration + run: PlatformRun + target_resource: Resource + target_version: ResourceVersion + output_artifact: Artifact + quality_evidence_artifact: Artifact + quality_result: QualityResult + lineage_event: LineageEvent + attempt_observation: FrameworkAttemptObservation + success_evidence: RunSuccessEvidence + output_path: str + quality_path: str + output_created: bool + quality_created: bool + final_run: PlatformRun | None = None + ledger_completed: bool | None = None + + @model_validator(mode="after") + def _consistent_bundle(self) -> PublicDataOpsResult: + tenant = self.landing.registration.resource.tenant_id + if self.run.tenant_id != tenant: + raise ValueError("DataOps bundle tenants must match") + if self.run.orchestration_class.value != "synchronous": + raise ValueError("public lightweight run must use synchronous orchestration") + if self.target_resource.tenant_id != tenant: + raise ValueError("target Resource tenant must match bundle") + if self.target_resource.resource_urn != self.target_version.resource_urn: + raise ValueError("target ResourceVersion must bind target Resource") + if self.target_version.content_sha256 != self.output_artifact.content_sha256: + raise ValueError("output Artifact must bind target ResourceVersion hash") + if self.output_artifact.artifact_role.value != "output": + raise ValueError("DataOps output Artifact must use output role") + if self.output_artifact.run_id != self.run.run_id: + raise ValueError("output Artifact must bind the PlatformRun") + if self.output_artifact.resource_version_id != self.target_version.resource_version_id: + raise ValueError("output Artifact must bind target ResourceVersion") + if self.quality_evidence_artifact.artifact_role.value != "evidence": + raise ValueError("quality Artifact must use evidence role") + if self.quality_evidence_artifact.run_id != self.run.run_id: + raise ValueError("quality Artifact must bind the PlatformRun") + if ( + self.quality_evidence_artifact.resource_version_id + != self.target_version.resource_version_id + ): + raise ValueError("quality Artifact must bind target ResourceVersion") + if self.quality_result.run_id != self.run.run_id: + raise ValueError("QualityResult must bind the PlatformRun") + if self.quality_result.resource_version_id != self.target_version.resource_version_id: + raise ValueError("QualityResult must bind target ResourceVersion") + if self.quality_result.evidence_artifact_id != self.quality_evidence_artifact.artifact_id: + raise ValueError("QualityResult must bind quality Artifact") + if self.lineage_event.run_id != self.run.run_id: + raise ValueError("LineageEvent must bind the PlatformRun") + if self.lineage_event.definition_version_id != self.run.definition_version_id: + raise ValueError("LineageEvent must bind the definition version") + if self.lineage_event.target_resource_version_id != self.target_version.resource_version_id: + raise ValueError("LineageEvent must bind target ResourceVersion") + if self.lineage_event.artifact_id != self.output_artifact.artifact_id: + raise ValueError("LineageEvent must bind output Artifact") + if self.attempt_observation.run_id != self.run.run_id: + raise ValueError("attempt observation must bind the PlatformRun") + if self.attempt_observation.framework_kind.value != "legacy": + raise ValueError("local inline attempt must use legacy framework kind") + if self.attempt_observation.evidence.get("execution_mode") != "local_inline": + raise ValueError("local inline attempt evidence is required") + if self.success_evidence.run_id != self.run.run_id: + raise ValueError("success evidence must bind the PlatformRun") + if self.success_evidence.output_artifact_id != self.output_artifact.artifact_id: + raise ValueError("success evidence must bind output Artifact") + if self.success_evidence.quality_result_id != self.quality_result.quality_result_id: + raise ValueError("success evidence must bind QualityResult") + if self.success_evidence.lineage_event_id != self.lineage_event.lineage_event_id: + raise ValueError("success evidence must bind LineageEvent") + if not Path(self.output_path).is_absolute() or not Path(self.quality_path).is_absolute(): + raise ValueError("serving paths must be absolute") + return self + + +def _safe_json_document(path: Path) -> dict[str, Any]: + try: + value = json.loads(path.read_text(encoding="utf-8")) + except (OSError, ValueError) as exc: + raise PublicDataOpsError("GeoJSON input is unreadable") from exc + if not isinstance(value, dict) or value.get("type") != "FeatureCollection": + raise PublicDataOpsError("GeoJSON input must be a FeatureCollection") + return value + + +def _normalize_features(features: Any) -> list[dict[str, Any]]: + if not isinstance(features, list): + raise PublicDataOpsError("GeoJSON features must be an array") + normalized: list[dict[str, Any]] = [] + for feature in features: + if not isinstance(feature, dict) or feature.get("type") != "Feature": + raise PublicDataOpsError("GeoJSON features must be Feature objects") + properties = feature.get("properties") + if properties is None: + properties = {} + if not isinstance(properties, dict): + raise PublicDataOpsError("GeoJSON feature properties must be objects") + item: dict[str, Any] = { + "geometry": feature.get("geometry"), + "properties": properties, + "type": "Feature", + } + if isinstance(feature.get("id"), (int, str)): + item["id"] = feature["id"] + normalized.append(item) + return normalized + + +def _safe_extract_zip(payload_path: Path, extraction_root: Path) -> list[Path]: + extracted: list[Path] = [] + try: + archive = zipfile.ZipFile(payload_path) + except (OSError, zipfile.BadZipFile) as exc: + raise PublicDataOpsError("Landing ZIP is unreadable") from exc + with archive: + infos = archive.infolist() + if not infos or len(infos) > 256: + raise PublicDataOpsError("Landing ZIP has an unsupported entry count") + total_size = sum(info.file_size for info in infos) + if total_size > 512 * 1024 * 1024: + raise PublicDataOpsError("Landing ZIP exceeds the lightweight size limit") + for info in infos: + name = PurePosixPath(info.filename) + if ( + name.is_absolute() + or not info.filename + or ".." in name.parts + or "\\" in info.filename + or "\x00" in info.filename + or (info.flag_bits & 0x1) + ): + raise PublicDataOpsError("Landing ZIP contains an unsafe entry") + mode = (info.external_attr >> 16) & 0o170000 + if mode == stat.S_IFLNK: + raise PublicDataOpsError("Landing ZIP contains a symbolic link") + destination = extraction_root.joinpath(*name.parts) + if info.is_dir(): + destination.mkdir(parents=True, exist_ok=True, mode=0o750) + continue + destination.parent.mkdir(parents=True, exist_ok=True, mode=0o750) + try: + with archive.open(info, "r") as source, destination.open("xb") as target: + while chunk := source.read(1024 * 1024): + target.write(chunk) + except (OSError, RuntimeError, zipfile.BadZipFile) as exc: + raise PublicDataOpsError("Landing ZIP extraction failed") from exc + extracted.append(destination) + return extracted + + +def _read_landing_features( + payload_path: Path, media_type: str, scratch_root: Path +) -> list[dict[str, Any]]: + if media_type == "application/geo+json" or payload_path.suffix.lower() in {".json", ".geojson"}: + return _normalize_features(_safe_json_document(payload_path).get("features")) + if media_type != "application/zip" and payload_path.suffix.lower() != ".zip": + raise PublicDataOpsError("lightweight DataOps supports GeoJSON or ZIP input") + extracted = _safe_extract_zip(payload_path, scratch_root) + candidates = sorted( + item for item in extracted if item.suffix.lower() in {".geojson", ".json", ".shp", ".gpkg"} + ) + if not candidates: + raise PublicDataOpsError("Landing ZIP contains no supported geospatial dataset") + candidate = next((item for item in candidates if item.suffix.lower() == ".shp"), candidates[0]) + if candidate.suffix.lower() in {".json", ".geojson"}: + return _normalize_features(_safe_json_document(candidate).get("features")) + try: + import geopandas as gpd + + frame = gpd.read_file(candidate) + except Exception as exc: + raise PublicDataOpsError("Landing vector dataset could not be read") from exc + if frame.crs is None: + raise PublicDataOpsError("Landing vector dataset has no declared CRS") + try: + frame = frame.to_crs("EPSG:4326") + return _normalize_features(list(frame.iterfeatures(drop_id=True))) + except Exception as exc: + raise PublicDataOpsError("Landing vector dataset could not be normalized") from exc + + +def _quality_metrics(features: list[dict[str, Any]]) -> dict[str, Any]: + from shapely.geometry import shape + + null_geometry_count = 0 + empty_geometry_count = 0 + invalid_geometry_count = 0 + geometry_types: dict[str, int] = {} + bbox: list[float] | None = None + for feature in features: + geometry = feature.get("geometry") + if geometry is None: + null_geometry_count += 1 + continue + try: + parsed = shape(geometry) + except Exception: + invalid_geometry_count += 1 + continue + geometry_types[parsed.geom_type] = geometry_types.get(parsed.geom_type, 0) + 1 + if parsed.is_empty: + empty_geometry_count += 1 + if not parsed.is_valid: + invalid_geometry_count += 1 + if not parsed.is_empty: + bounds = parsed.bounds + if bbox is None: + bbox = list(bounds) + else: + bbox = [ + min(bbox[0], bounds[0]), + min(bbox[1], bounds[1]), + max(bbox[2], bounds[2]), + max(bbox[3], bounds[3]), + ] + return { + "feature_count": len(features), + "null_geometry_count": null_geometry_count, + "empty_geometry_count": empty_geometry_count, + "invalid_geometry_count": invalid_geometry_count, + "geometry_types": dict(sorted(geometry_types.items())), + "bbox_epsg4326": bbox, + } + + +def _definition_registration( + request: PublicDataOpsRequest, landing: PublicSourceLandingResult +) -> DefinitionRegistration: + tenant = landing.registration.resource.tenant_id + definition_urn = f"gda://{tenant}/definition/public-geojson-materialize" + definition_document = { + "schema": DEFINITION_SCHEMA, + "operation": "materialize_public_geojson", + "execution_mode": "local_inline", + "input_admission_class": "public_open", + "quality_rule_version": QUALITY_RULE_VERSION, + "output_media_type": "application/geo+json", + } + input_contract = { + "binding_name": "source", + "resource_kind": "dataset", + "media_types": ["application/geo+json", "application/zip"], + "content_admission": "public_open", + } + output_contract = { + "resource_kind": "dataset", + "media_type": "application/geo+json", + "serving_profile": "lightweight", + "crs": "EPSG:4326", + } + definition_sha256 = platform_definition_fingerprint( + orchestration_class="synchronous", + capability_id="public.geojson.materialize", + portability_class="portable", + definition_document=definition_document, + input_contract=input_contract, + output_contract=output_contract, + ) + definition_version_id = uuid5(NAMESPACE_URL, f"{definition_urn}@sha256:{definition_sha256}") + return DefinitionRegistration( + resource=Resource( + tenant_id=tenant, + resource_urn=definition_urn, + resource_kind="definition", + authority_system="gda", + authority_locator="definition/public-geojson-materialize", + owner_ref=landing.registration.resource.owner_ref, + governance_ref={"profile": "public_open_lightweight"}, + ), + resource_version=ResourceVersion( + tenant_id=tenant, + resource_urn=definition_urn, + resource_version_id=definition_version_id, + version_key=f"sha256:{definition_sha256[:16]}", + content_sha256=definition_sha256, + authority_version_ref={"schema": DEFINITION_SCHEMA, "revision": 1}, + created_by=DEFINITION_PUBLISHED_BY, + created_at=DEFINITION_PUBLISHED_AT, + ), + definition=PlatformDefinitionVersion( + tenant_id=tenant, + definition_urn=definition_urn, + definition_version_id=definition_version_id, + orchestration_class="synchronous", + capability_id="public.geojson.materialize", + portability_class="portable", + definition_document=definition_document, + input_contract=input_contract, + output_contract=output_contract, + definition_sha256=definition_sha256, + ), + ) + + +def materialize_public_dataops( + landing: PublicSourceLandingResult, + request: PublicDataOpsRequest, + *, + serving_root: Path, +) -> PublicDataOpsResult: + """Materialize, quality-check, and build an idempotent control bundle.""" + verify_public_source_landing(landing) + if landing.registration.resource.tenant_id != landing.registration.artifact.tenant_id: + raise PublicDataOpsError("Landing tenant binding is invalid") + if serving_root.is_symlink(): + raise PublicDataOpsError("serving root cannot be a symbolic link") + serving_root = serving_root.resolve() + serving_root.mkdir(parents=True, exist_ok=True, mode=0o750) + staging_root = serving_root / ".staging" + staging_root.mkdir(parents=True, exist_ok=True, mode=0o700) + with tempfile.TemporaryDirectory(prefix="public-dataops-", dir=staging_root) as scratch: + features = _read_landing_features( + Path(landing.payload_path), landing.registration.artifact.media_type, Path(scratch) + ) + metrics = _quality_metrics(features) + if ( + metrics["feature_count"] < request.min_feature_count + or metrics["null_geometry_count"] + or metrics["empty_geometry_count"] + or metrics["invalid_geometry_count"] + ): + raise PublicDataOpsError(f"GeoJSON quality gate failed: {metrics}") + normalized = {"features": features, "type": "FeatureCollection"} + output_bytes = canonical_json_bytes(normalized) + b"\n" + output_sha256 = hashlib.sha256(output_bytes).hexdigest() + target_urn = ( + f"gda://{landing.registration.resource.tenant_id}/dataset/{request.output_dataset_id}" + ) + target_version_id = uuid5(NAMESPACE_URL, f"{target_urn}@sha256:{output_sha256}") + definition_registration = _definition_registration(request, landing) + run_id = uuid5( + NAMESPACE_URL, + f"{definition_registration.resource_version.resource_version_id}:run:" + f"{landing.registration.resource_version.resource_version_id}:{request.output_dataset_id}", + ) + output_root = ( + serving_root + / landing.registration.resource.tenant_id + / request.output_dataset_id + / "sha256" + / output_sha256 + ) + output_path = output_root / "data.geojson" + output_created = _install_immutable_bytes(output_path, output_bytes) + target_resource = Resource( + tenant_id=landing.registration.resource.tenant_id, + resource_urn=target_urn, + resource_kind="dataset", + authority_system="gda_lightweight_serving", + authority_locator=f"{landing.registration.resource.tenant_id}/{request.output_dataset_id}", + owner_ref=landing.registration.resource.owner_ref, + governance_ref={ + "profile": "public_open_lightweight", + "source_resource_urn": landing.registration.resource.resource_urn, + "production_ready": False, + }, + ) + target_version = ResourceVersion( + tenant_id=target_resource.tenant_id, + resource_urn=target_urn, + resource_version_id=target_version_id, + version_key=f"sha256:{output_sha256[:16]}", + content_sha256=output_sha256, + authority_version_ref={ + "authority_system": "gda_lightweight_serving", + "object_key": output_path.relative_to(serving_root).as_posix(), + "source_resource_version_id": str( + landing.registration.resource_version.resource_version_id + ), + }, + created_by=request.executor, + created_at=request.executed_at, + ) + run = PlatformRun( + tenant_id=target_resource.tenant_id, + run_id=run_id, + definition_version_id=definition_registration.resource_version.resource_version_id, + orchestration_class="synchronous", + subject_context=SubjectContext( + tenant_id=target_resource.tenant_id, + subject_id=request.executor.removeprefix("workload:"), + subject_type="workload", + roles=("platform_operator",), + purpose="materialize public Landing into lightweight GeoJSON serving", + ), + input_bindings=( + ResourceBinding( + binding_name="source", + resource_version_id=landing.registration.resource_version.resource_version_id, + semantic_type="geo.public_source", + ), + ), + idempotency_key=f"public-geojson:{landing.registration.resource_version.resource_version_id}:{request.output_dataset_id}", + config_fingerprint=canonical_json_fingerprint( + { + "min_feature_count": request.min_feature_count, + "quality_rule_version": QUALITY_RULE_VERSION, + } + ), + submitted_at=request.executed_at, + ) + output_manifest = { + "schema": "gda.public_dataops_output.v1", + "run_id": str(run_id), + "source_resource_version_id": str( + landing.registration.resource_version.resource_version_id + ), + "target_resource_version_id": str(target_version_id), + "content_sha256": output_sha256, + "size_bytes": len(output_bytes), + "feature_count": metrics["feature_count"], + "crs": "EPSG:4326", + "profile": "public_open_lightweight", + "production_ready": False, + } + output_artifact_id = uuid5(NAMESPACE_URL, f"{run_id}:output:{output_sha256}") + output_artifact = Artifact( + tenant_id=run.tenant_id, + artifact_id=output_artifact_id, + artifact_key=f"public-geojson-output:{request.output_dataset_id}:{output_sha256[:12]}", + artifact_role="output", + storage_uri=output_path.as_uri(), + media_type="application/geo+json", + content_sha256=output_sha256, + size_bytes=len(output_bytes), + run_id=run_id, + resource_version_id=target_version_id, + manifest=output_manifest, + created_by=request.executor, + created_at=request.executed_at, + ) + quality_metrics = {**metrics, "output_sha256": output_sha256, "crs": "EPSG:4326"} + quality_result_id = uuid5(NAMESPACE_URL, f"{run_id}:quality:{output_sha256}") + quality_document = { + "schema": QUALITY_SCHEMA, + "run_id": str(run_id), + "resource_version_id": str(target_version_id), + "rule_version_ref": QUALITY_RULE_VERSION, + "verdict": "passed", + "metrics": quality_metrics, + "evaluated_by": request.quality_evaluator, + "evaluated_at": request.executed_at.isoformat().replace("+00:00", "Z"), + } + quality_bytes = canonical_json_bytes(quality_document) + b"\n" + quality_sha256 = hashlib.sha256(quality_bytes).hexdigest() + quality_path = ( + serving_root + / landing.registration.resource.tenant_id + / request.output_dataset_id + / "evidence" + / "sha256" + / quality_sha256 + / "quality.json" + ) + quality_created = _install_immutable_bytes(quality_path, quality_bytes) + quality_artifact_id = uuid5(NAMESPACE_URL, f"{run_id}:quality-artifact:{quality_sha256}") + quality_artifact = Artifact( + tenant_id=run.tenant_id, + artifact_id=quality_artifact_id, + artifact_key=f"public-geojson-quality:{request.output_dataset_id}:{quality_sha256[:12]}", + artifact_role="evidence", + storage_uri=quality_path.as_uri(), + media_type="application/vnd.gda.public-dataops-quality+json", + content_sha256=quality_sha256, + size_bytes=len(quality_bytes), + run_id=run_id, + resource_version_id=target_version_id, + manifest=quality_document, + created_by=request.quality_evaluator, + created_at=request.executed_at, + ) + quality_result = QualityResult( + tenant_id=run.tenant_id, + quality_result_id=quality_result_id, + run_id=run_id, + resource_version_id=target_version_id, + rule_version_ref=QUALITY_RULE_VERSION, + verdict="passed", + metrics=quality_metrics, + evidence_artifact_id=quality_artifact_id, + result_sha256=quality_result_fingerprint( + tenant_id=run.tenant_id, + run_id=run_id, + resource_version_id=target_version_id, + rule_version_ref=QUALITY_RULE_VERSION, + verdict="passed", + metrics=quality_metrics, + evidence_artifact_id=quality_artifact_id, + evaluated_by=request.quality_evaluator, + evaluated_at=request.executed_at, + ), + evaluated_by=request.quality_evaluator, + evaluated_at=request.executed_at, + ) + lineage_event_id = uuid5(NAMESPACE_URL, f"{run_id}:lineage:{output_sha256}") + lineage_facets = { + "schema": LINEAGE_SCHEMA, + "operation": "materialize_public_geojson", + "input_media_type": landing.registration.artifact.media_type, + "output_media_type": "application/geo+json", + "output_sha256": output_sha256, + } + lineage_event = LineageEvent( + tenant_id=run.tenant_id, + lineage_event_id=lineage_event_id, + event_type="materialize", + source_resource_version_id=landing.registration.resource_version.resource_version_id, + target_resource_version_id=target_version_id, + producer=request.executor, + event_sha256=canonical_json_fingerprint( + { + "source_resource_version_id": str( + landing.registration.resource_version.resource_version_id + ), + "target_resource_version_id": str(target_version_id), + "run_id": str(run_id), + "facets": lineage_facets, + } + ), + run_id=run_id, + definition_version_id=run.definition_version_id, + artifact_id=output_artifact_id, + facets=lineage_facets, + occurred_at=request.executed_at, + ) + attempt_id = uuid5(NAMESPACE_URL, f"{run_id}:attempt:1:success") + attempt_evidence = { + "schema": ATTEMPT_SCHEMA, + "execution_mode": "local_inline", + "executor": request.executor, + "output_sha256": output_sha256, + "quality_result_id": str(quality_result_id), + } + attempt = FrameworkAttemptObservation( + tenant_id=run.tenant_id, + observation_id=attempt_id, + run_id=run_id, + attempt_no=1, + framework_kind="legacy", + external_namespace="gda-public-lightweight", + external_run_id=str(run_id), + external_attempt_id="1", + observed_state="success", + observation_sha256=canonical_json_fingerprint(attempt_evidence), + evidence=attempt_evidence, + observed_at=request.executed_at, + ) + success_evidence = RunSuccessEvidence( + tenant_id=run.tenant_id, + run_id=run_id, + attempt_observation_id=attempt_id, + output_artifact_id=output_artifact_id, + quality_result_id=quality_result_id, + lineage_event_id=lineage_event_id, + evidence_sha256=run_success_evidence_fingerprint( + tenant_id=run.tenant_id, + run_id=run_id, + attempt_observation_id=attempt_id, + output_artifact_id=output_artifact_id, + quality_result_id=quality_result_id, + lineage_event_id=lineage_event_id, + ), + ) + return PublicDataOpsResult( + landing=landing, + definition_registration=definition_registration, + run=run, + target_resource=target_resource, + target_version=target_version, + output_artifact=output_artifact, + quality_evidence_artifact=quality_artifact, + quality_result=quality_result, + lineage_event=lineage_event, + attempt_observation=attempt, + success_evidence=success_evidence, + output_path=str(output_path), + quality_path=str(quality_path), + output_created=output_created, + quality_created=quality_created, + ) + + +def register_public_dataops( + result: PublicDataOpsResult, gateway: PlatformGateway +) -> PublicDataOpsResult: + """Replayably register the bundle and finalize the synchronous Run.""" + gateway.register_landing(result.landing.registration) + gateway.register_definition(result.definition_registration) + gateway.register_resource(result.target_resource) + gateway.register_resource_version(result.target_version) + gateway.submit_run(result.run, request_dispatch=False) + current = gateway.get_run(result.run.tenant_id, result.run.run_id) + actor = ( + result.run.subject_context.subject_type.value + ":" + result.run.subject_context.subject_id + ) + if current.status.value == "accepted": + current = gateway.transition_run( + result.run.tenant_id, + result.run.run_id, + 0, + "dispatching", + actor, + "local lightweight profile accepted Run", + ) + if current.status.value == "dispatching": + current = gateway.transition_run( + result.run.tenant_id, + result.run.run_id, + 1, + "running", + actor, + "local lightweight executor started Run", + ) + if current.status.value not in {"running", "succeeded"}: + raise PublicDataOpsError(f"cannot finalize Run in status {current.status.value}") + gateway.record_artifact(result.output_artifact) + gateway.record_artifact(result.quality_evidence_artifact) + gateway.record_attempt(result.attempt_observation) + gateway.record_quality_result(result.quality_result) + gateway.record_lineage(result.lineage_event) + final_run = gateway.finalize_run_success( + result.success_evidence, + expected_state_version=2, + actor_subject=actor, + reason="public lightweight GeoJSON materialization passed content and quality gates", + ) + return result.model_copy(update={"final_run": final_run, "ledger_completed": True}) + + +def verify_public_dataops_result(result: PublicDataOpsResult) -> None: + verify_public_source_landing(result.landing) + output = Path(result.output_path) + quality = Path(result.quality_path) + if output.is_symlink() or quality.is_symlink() or not output.is_file() or not quality.is_file(): + raise PublicDataOpsError("serving output or quality evidence is missing") + output_bytes = output.read_bytes() + if hashlib.sha256(output_bytes).hexdigest() != result.target_version.content_sha256: + raise PublicDataOpsError("serving output does not match target ResourceVersion") + if output.as_uri() != result.output_artifact.storage_uri: + raise PublicDataOpsError("output Artifact URI does not match serving output") + quality_bytes = quality.read_bytes() + if hashlib.sha256(quality_bytes).hexdigest() != result.quality_evidence_artifact.content_sha256: + raise PublicDataOpsError("quality evidence does not match its Artifact hash") + quality_document = json.loads(quality_bytes) + if quality_document != result.quality_evidence_artifact.manifest: + raise PublicDataOpsError("quality evidence does not match its Artifact") + if quality.as_uri() != result.quality_evidence_artifact.storage_uri: + raise PublicDataOpsError("quality Artifact URI does not match evidence file") + if result.final_run is not None and result.final_run.status.value != "succeeded": + raise PublicDataOpsError("registered DataOps Run is not succeeded") + + +def _write_result(result: PublicDataOpsResult, output: Path | None) -> None: + rendered = json.dumps( + result.model_dump(mode="json"), ensure_ascii=True, indent=2, sort_keys=True + ) + if output is not None: + output.parent.mkdir(parents=True, exist_ok=True) + output.write_text(rendered + "\n", encoding="utf-8") + print(rendered) + + +def _parse_time(value: str) -> datetime: + parsed = datetime.fromisoformat(value.replace("Z", "+00:00")) + if parsed.tzinfo is None or parsed.utcoffset() is None: + raise argparse.ArgumentTypeError("timestamp must include a timezone") + return parsed.astimezone(UTC) + + +def main(argv: list[str] | None = None) -> int: + parser = argparse.ArgumentParser(description=__doc__) + subparsers = parser.add_subparsers(dest="command", required=True) + run_parser = subparsers.add_parser("run") + run_parser.add_argument("--landing-result", type=Path, required=True) + run_parser.add_argument("--serving-root", type=Path, required=True) + run_parser.add_argument("--output-dataset-id", required=True) + run_parser.add_argument("--executor", default="workload:public-dataops") + run_parser.add_argument("--quality-evaluator", default="workload:public-quality") + run_parser.add_argument("--executed-at", type=_parse_time, required=True) + run_parser.add_argument("--min-feature-count", type=int, default=1) + run_parser.add_argument("--database-url") + run_parser.add_argument("--output", type=Path) + verify_parser = subparsers.add_parser("verify") + verify_parser.add_argument("--input", type=Path, required=True) + args = parser.parse_args(argv) + if args.command == "run": + landing = PublicSourceLandingResult.model_validate_json( + args.landing_result.read_text(encoding="utf-8") + ) + request = PublicDataOpsRequest( + executor=args.executor, + quality_evaluator=args.quality_evaluator, + output_dataset_id=args.output_dataset_id, + executed_at=args.executed_at, + min_feature_count=args.min_feature_count, + ) + result = materialize_public_dataops(landing, request, serving_root=args.serving_root) + if args.database_url: + from sqlalchemy import create_engine + + result = register_public_dataops( + result, PlatformGateway(create_engine(args.database_url)) + ) + verify_public_dataops_result(result) + _write_result(result, args.output) + return 0 + result = PublicDataOpsResult.model_validate_json(args.input.read_text(encoding="utf-8")) + verify_public_dataops_result(result) + print( + json.dumps( + { + "valid": True, + "run_id": str(result.run.run_id), + "output_sha256": result.output_artifact.content_sha256, + } + ) + ) + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/data_agent/test_public_dataops_run.py b/data_agent/test_public_dataops_run.py new file mode 100644 index 00000000..f56c7499 --- /dev/null +++ b/data_agent/test_public_dataops_run.py @@ -0,0 +1,294 @@ +import hashlib +import json +import os +import stat +import zipfile +from datetime import UTC, datetime, timedelta +from pathlib import Path + +import pytest +from pydantic import ValidationError + +from data_agent.public_dataops_run import ( + ATTEMPT_SCHEMA, + DEFINITION_PUBLISHED_AT, + PublicDataOpsError, + PublicDataOpsRequest, + PublicDataOpsResult, + _safe_extract_zip, + main, + materialize_public_dataops, + verify_public_dataops_result, +) +from data_agent.public_source_landing import ( + PublicSourceLandingRequest, + stage_public_source, +) + +NOW = datetime(2026, 8, 17, 14, 0, tzinfo=UTC) +FEATURE_COLLECTION = { + "type": "FeatureCollection", + "features": [ + { + "type": "Feature", + "id": "capital", + "properties": {"name": "Example City"}, + "geometry": {"type": "Point", "coordinates": [106.55, 29.56]}, + }, + { + "type": "Feature", + "properties": {"name": "Example Area"}, + "geometry": { + "type": "Polygon", + "coordinates": [[[106.0, 29.0], [107.0, 29.0], [107.0, 30.0], [106.0, 29.0]]], + }, + }, + ], +} + + +def _request(**updates) -> PublicDataOpsRequest: + values = { + "executor": "workload:public-dataops", + "quality_evaluator": "workload:public-quality", + "output_dataset_id": "countries-serving", + "executed_at": NOW, + "min_feature_count": 1, + } + values.update(updates) + return PublicDataOpsRequest(**values) + + +def _stage_geojson( + tmp_path: Path, + document=FEATURE_COLLECTION, + *, + dataset_id: str = "public-features", + created_at: datetime = NOW, +): + payload = json.dumps(document, sort_keys=True).encode() + b"\n" + source = tmp_path / f"{dataset_id}.geojson" + source.write_bytes(payload) + return stage_public_source( + PublicSourceLandingRequest( + tenant_id="public-demo", + dataset_id=dataset_id, + source_uri=f"https://example.org/open/{dataset_id}.geojson", + license_id="CC0-1.0", + owner_ref="team:data-platform", + expected_sha256=hashlib.sha256(payload).hexdigest(), + media_type="application/geo+json", + created_by="workload:public-source-ingest", + created_at=created_at, + ), + source_path=source, + landing_root=tmp_path / "landing", + ) + + +def _materialize(tmp_path: Path, document=FEATURE_COLLECTION): + return materialize_public_dataops( + _stage_geojson(tmp_path, document), + _request(), + serving_root=tmp_path / "serving", + ) + + +def test_materializes_deterministic_content_addressed_geojson_and_replays(tmp_path): + landing = _stage_geojson(tmp_path) + request = _request() + first = materialize_public_dataops(landing, request, serving_root=tmp_path / "serving") + verify_public_dataops_result(first) + + output = json.loads(Path(first.output_path).read_bytes()) + assert output["type"] == "FeatureCollection" + assert len(output["features"]) == 2 + assert first.output_created is True + assert first.quality_created is True + assert f"sha256/{first.output_artifact.content_sha256}/data.geojson" in first.output_path + assert ( + f"evidence/sha256/{first.quality_evidence_artifact.content_sha256}/quality.json" + in first.quality_path + ) + assert first.quality_result.metrics == { + "bbox_epsg4326": [106.0, 29.0, 107.0, 30.0], + "crs": "EPSG:4326", + "empty_geometry_count": 0, + "feature_count": 2, + "geometry_types": {"Point": 1, "Polygon": 1}, + "invalid_geometry_count": 0, + "null_geometry_count": 0, + "output_sha256": first.output_artifact.content_sha256, + } + + replay = materialize_public_dataops(landing, request, serving_root=tmp_path / "serving") + verify_public_dataops_result(replay) + assert replay.output_created is False + assert replay.quality_created is False + assert replay.model_copy(update={"output_created": True, "quality_created": True}) == first + + +def test_identity_chain_and_portable_definition_are_exact(tmp_path): + landing = _stage_geojson(tmp_path) + result = materialize_public_dataops(landing, _request(), serving_root=tmp_path / "serving") + definition = result.definition_registration + + assert "source_resource_urn" not in definition.definition.definition_document + assert definition.resource_version.created_at == DEFINITION_PUBLISHED_AT + assert result.target_resource.governance_ref["source_resource_urn"] == ( + landing.registration.resource.resource_urn + ) + assert result.run.definition_version_id == definition.definition.definition_version_id + assert result.run.input_bindings[0].resource_version_id == ( + landing.registration.resource_version.resource_version_id + ) + assert result.output_artifact.resource_version_id == result.target_version.resource_version_id + assert result.quality_result.evidence_artifact_id == ( + result.quality_evidence_artifact.artifact_id + ) + assert result.lineage_event.source_resource_version_id == ( + landing.registration.resource_version.resource_version_id + ) + assert result.lineage_event.target_resource_version_id == ( + result.target_version.resource_version_id + ) + assert result.success_evidence.attempt_observation_id == ( + result.attempt_observation.observation_id + ) + assert result.attempt_observation.evidence["schema"] == ATTEMPT_SCHEMA + assert result.attempt_observation.evidence["execution_mode"] == "local_inline" + + later_landing = _stage_geojson( + tmp_path, + dataset_id="other-public-features", + created_at=NOW + timedelta(days=1), + ) + later = materialize_public_dataops( + later_landing, + _request( + executor="workload:alternate-dataops", + output_dataset_id="other-serving", + executed_at=NOW + timedelta(days=1), + ), + serving_root=tmp_path / "serving", + ) + assert later.definition_registration == definition + + +@pytest.mark.parametrize( + "geometry", + [ + None, + {"type": "GeometryCollection", "geometries": []}, + { + "type": "Polygon", + "coordinates": [[[0.0, 0.0], [1.0, 1.0], [1.0, 0.0], [0.0, 1.0], [0.0, 0.0]]], + }, + ], +) +def test_quality_gate_rejects_null_empty_or_invalid_geometry_before_publish(tmp_path, geometry): + document = { + "type": "FeatureCollection", + "features": [{"type": "Feature", "properties": {}, "geometry": geometry}], + } + with pytest.raises(PublicDataOpsError, match="quality gate failed"): + _materialize(tmp_path, document) + assert not tuple((tmp_path / "serving").glob("**/data.geojson")) + + +def test_zip_with_directory_and_geojson_is_supported(tmp_path): + archive = tmp_path / "features.zip" + with zipfile.ZipFile(archive, "w") as handle: + handle.writestr("dataset/", b"") + handle.writestr("dataset/features.geojson", json.dumps(FEATURE_COLLECTION)) + landing = stage_public_source( + PublicSourceLandingRequest( + tenant_id="public-demo", + dataset_id="zipped-features", + source_uri="https://example.org/open/zipped-features.zip", + license_id="CC0-1.0", + owner_ref="team:data-platform", + expected_sha256=hashlib.sha256(archive.read_bytes()).hexdigest(), + media_type="application/zip", + created_by="workload:public-source-ingest", + created_at=NOW, + ), + source_path=archive, + landing_root=tmp_path / "landing", + ) + result = materialize_public_dataops(landing, _request(), serving_root=tmp_path / "serving") + assert result.quality_result.metrics["feature_count"] == 2 + + +def _mark_zip_encrypted(path: Path) -> None: + payload = bytearray(path.read_bytes()) + for signature, flag_offset in ((b"PK\x03\x04", 6), (b"PK\x01\x02", 8)): + position = payload.find(signature) + assert position >= 0 + payload[position + flag_offset] |= 1 + path.write_bytes(payload) + + +@pytest.mark.parametrize("unsafe_kind", ["traversal", "symlink", "encrypted"]) +def test_zip_rejects_unsafe_entries(tmp_path, unsafe_kind): + archive = tmp_path / f"{unsafe_kind}.zip" + with zipfile.ZipFile(archive, "w") as handle: + if unsafe_kind == "traversal": + handle.writestr("../outside.geojson", json.dumps(FEATURE_COLLECTION)) + elif unsafe_kind == "symlink": + info = zipfile.ZipInfo("features.geojson") + info.create_system = 3 + info.external_attr = (stat.S_IFLNK | 0o777) << 16 + handle.writestr(info, "target.geojson") + else: + handle.writestr("features.geojson", json.dumps(FEATURE_COLLECTION)) + if unsafe_kind == "encrypted": + _mark_zip_encrypted(archive) + with pytest.raises(PublicDataOpsError, match="unsafe entry|symbolic link"): + _safe_extract_zip(archive, tmp_path / "extract") + assert not (tmp_path / "outside.geojson").exists() + + +def test_request_requires_independent_quality_evaluator(): + with pytest.raises(ValidationError, match="quality evaluator must be independent"): + _request(quality_evaluator="workload:public-dataops") + + +@pytest.mark.parametrize("target", ["output", "quality"]) +def test_verification_detects_output_and_quality_tampering(tmp_path, target): + result = _materialize(tmp_path) + path = Path(result.output_path if target == "output" else result.quality_path) + os.chmod(path, 0o640) + path.write_bytes(path.read_bytes() + b"tampered") + with pytest.raises(PublicDataOpsError, match="does not match"): + verify_public_dataops_result(result) + + +def test_cli_runs_and_verifies_bundle(tmp_path, capsys): + landing = _stage_geojson(tmp_path) + landing_result = tmp_path / "landing-result.json" + landing_result.write_text(landing.model_dump_json(indent=2) + "\n", encoding="utf-8") + output = tmp_path / "dataops-result.json" + assert ( + main( + [ + "run", + "--landing-result", + str(landing_result), + "--serving-root", + str(tmp_path / "serving"), + "--output-dataset-id", + "countries-serving", + "--executed-at", + "2026-08-17T14:00:00Z", + "--output", + str(output), + ] + ) + == 0 + ) + capsys.readouterr() + result = PublicDataOpsResult.model_validate_json(output.read_text(encoding="utf-8")) + assert result.ledger_completed is None + assert main(["verify", "--input", str(output)]) == 0 + assert json.loads(capsys.readouterr().out)["valid"] is True diff --git a/data_agent/test_public_dataops_run_postgres.py b/data_agent/test_public_dataops_run_postgres.py new file mode 100644 index 00000000..48b2de1f --- /dev/null +++ b/data_agent/test_public_dataops_run_postgres.py @@ -0,0 +1,308 @@ +import hashlib +import json +import os +from datetime import UTC, datetime +from pathlib import Path +from uuid import uuid4 + +import pytest +from sqlalchemy import create_engine, text + +from data_agent.platform_contracts import ( + FrameworkAttemptObservation, + QualityResult, + RunSuccessEvidence, + canonical_json_fingerprint, + quality_result_fingerprint, + run_success_evidence_fingerprint, +) +from data_agent.platform_gateway import GatewayValidationError, PlatformGateway +from data_agent.public_dataops_run import ( + PublicDataOpsRequest, + materialize_public_dataops, + register_public_dataops, +) +from data_agent.public_source_landing import ( + PublicSourceLandingRequest, + stage_public_source, +) + +DATABASE_URL = os.environ.get("DATABASE_URL") +MIGRATIONS = tuple( + Path(__file__).resolve().parent / "migrations" / filename + for filename in ( + "092_platform_control_ledger.sql", + "093_app_user_tenant_context.sql", + "094_platform_control_gateway.sql", + "095_platform_command_outbox.sql", + "096_platform_success_verdict.sql", + "099_synchronous_success_profile.sql", + ) +) +NOW = datetime(2026, 8, 17, 14, 0, tzinfo=UTC) +FEATURE_COLLECTION = { + "type": "FeatureCollection", + "features": [ + { + "type": "Feature", + "properties": {"name": "Example"}, + "geometry": {"type": "Point", "coordinates": [106.55, 29.56]}, + } + ], +} +FINALIZE_REASON = "public lightweight GeoJSON materialization passed content and quality gates" + + +@pytest.fixture(scope="module") +def postgres_engine(): + if not DATABASE_URL: + pytest.skip("DATABASE_URL is not configured") + engine = create_engine(DATABASE_URL) + with engine.begin() as connection: + is_superuser = connection.exec_driver_sql( + "SELECT rolsuper FROM pg_roles WHERE rolname = current_user" + ).scalar_one() + if not is_superuser: + pytest.skip("public DataOps gateway test requires a PostgreSQL superuser") + connection.exec_driver_sql( + """ + CREATE TABLE IF NOT EXISTS agent_app_users ( + id SERIAL PRIMARY KEY, + username VARCHAR(100) UNIQUE NOT NULL + ) + """ + ) + for migration in MIGRATIONS: + connection.execute(text(migration.read_text(encoding="utf-8"))) + try: + yield engine + finally: + engine.dispose() + + +def _bundle(tmp_path: Path, *, tenant: str, suffix: str): + payload = json.dumps(FEATURE_COLLECTION, sort_keys=True).encode() + b"\n" + source = tmp_path / f"source-{suffix}.geojson" + source.write_bytes(payload) + landing = stage_public_source( + PublicSourceLandingRequest( + tenant_id=tenant, + dataset_id=f"source-{suffix}", + source_uri=f"https://example.org/open/source-{suffix}.geojson", + license_id="CC0-1.0", + owner_ref="team:data-platform", + expected_sha256=hashlib.sha256(payload).hexdigest(), + media_type="application/geo+json", + created_by="workload:public-source-ingest", + created_at=NOW, + ), + source_path=source, + landing_root=tmp_path / "landing", + ) + request = PublicDataOpsRequest( + executor="workload:public-dataops", + quality_evaluator="workload:public-quality", + output_dataset_id=f"serving-{suffix}", + executed_at=NOW, + ) + return ( + landing, + request, + materialize_public_dataops(landing, request, serving_root=tmp_path / "serving"), + ) + + +def _actor(result) -> str: + return ( + f"{result.run.subject_context.subject_type.value}:{result.run.subject_context.subject_id}" + ) + + +def _prepare_running_bundle( + gateway: PlatformGateway, + result, + *, + attempt=None, + quality_artifact=None, + quality_result=None, + lineage=None, +): + gateway.register_landing(result.landing.registration) + gateway.register_definition(result.definition_registration) + gateway.register_resource(result.target_resource) + gateway.register_resource_version(result.target_version) + gateway.submit_run(result.run, request_dispatch=False) + actor = _actor(result) + gateway.transition_run( + result.run.tenant_id, + result.run.run_id, + 0, + "dispatching", + actor, + "local lightweight profile accepted Run", + ) + gateway.transition_run( + result.run.tenant_id, + result.run.run_id, + 1, + "running", + actor, + "local lightweight executor started Run", + ) + gateway.record_artifact(result.output_artifact) + gateway.record_artifact(quality_artifact or result.quality_evidence_artifact) + gateway.record_attempt(attempt or result.attempt_observation) + gateway.record_quality_result(quality_result or result.quality_result) + gateway.record_lineage(lineage or result.lineage_event) + + +def test_postgres_public_dataops_succeeds_once_and_replays(tmp_path, postgres_engine): + tenant = f"public-run-{uuid4().hex[:12]}" + landing, request, result = _bundle(tmp_path, tenant=tenant, suffix="countries") + gateway = PlatformGateway(postgres_engine) + + completed = register_public_dataops(result, gateway) + assert completed.final_run is not None + assert completed.final_run.status.value == "succeeded" + assert completed.final_run.state_version == 3 + assert completed.ledger_completed is True + + replay_bundle = materialize_public_dataops(landing, request, serving_root=tmp_path / "serving") + assert replay_bundle.output_created is False + assert replay_bundle.quality_created is False + replay = register_public_dataops(replay_bundle, gateway) + assert replay.final_run == completed.final_run + + with postgres_engine.connect() as connection: + counts = connection.execute( + text( + """ + SELECT + (SELECT count(*) FROM gda_control.platform_run + WHERE tenant_id = :tenant_id), + (SELECT count(*) FROM gda_control.platform_run_event + WHERE tenant_id = :tenant_id), + (SELECT count(*) FROM gda_control.artifact + WHERE tenant_id = :tenant_id AND run_id = :run_id), + (SELECT count(*) FROM gda_control.framework_attempt_observation + WHERE tenant_id = :tenant_id AND run_id = :run_id), + (SELECT count(*) FROM gda_control.quality_result + WHERE tenant_id = :tenant_id AND run_id = :run_id), + (SELECT count(*) FROM gda_control.lineage_event + WHERE tenant_id = :tenant_id AND run_id = :run_id) + """ + ), + {"tenant_id": tenant, "run_id": result.run.run_id}, + ).one() + privilege = connection.execute( + text( + """ + SELECT has_function_privilege( + 'gda_control_gateway', + 'gda_control.finalize_synchronous_platform_run_success(text,uuid,integer,text,text,jsonb)', + 'EXECUTE' + ) + """ + ) + ).scalar_one() + assert counts == (1, 4, 2, 1, 1, 1) + assert privilege is True + + +def _attempt_with(result, **updates) -> FrameworkAttemptObservation: + values = result.attempt_observation.model_dump(mode="python") + evidence = updates.pop("evidence", values["evidence"]) + values.update(updates) + values["evidence"] = evidence + values["observation_sha256"] = canonical_json_fingerprint(evidence) + return FrameworkAttemptObservation.model_validate(values) + + +def _success_with_output(result, artifact_id) -> RunSuccessEvidence: + values = result.success_evidence.model_dump(mode="python") + values["output_artifact_id"] = artifact_id + values["evidence_sha256"] = run_success_evidence_fingerprint( + tenant_id=values["tenant_id"], + run_id=values["run_id"], + attempt_observation_id=values["attempt_observation_id"], + output_artifact_id=values["output_artifact_id"], + quality_result_id=values["quality_result_id"], + lineage_event_id=values["lineage_event_id"], + ) + return RunSuccessEvidence.model_validate(values) + + +@pytest.mark.parametrize( + "rejection", + [ + "dolphinscheduler_observation", + "wrong_observation_schema", + "wrong_execution_mode", + "same_quality_evaluator", + "unbound_output", + "unbound_lineage", + ], +) +def test_postgres_synchronous_finalizer_rejects_unbound_or_wrong_profile_evidence( + tmp_path, postgres_engine, rejection +): + tenant = f"sync-reject-{uuid4().hex[:12]}" + _, _, result = _bundle(tmp_path, tenant=tenant, suffix=rejection[:24]) + gateway = PlatformGateway(postgres_engine) + attempt = result.attempt_observation + quality_artifact = result.quality_evidence_artifact + quality_result = result.quality_result + lineage = result.lineage_event + success = result.success_evidence + + if rejection == "dolphinscheduler_observation": + attempt = _attempt_with(result, framework_kind="dolphinscheduler") + elif rejection == "wrong_observation_schema": + attempt = _attempt_with( + result, + evidence={**attempt.evidence, "schema": "gda.wrong.v1"}, + ) + elif rejection == "wrong_execution_mode": + attempt = _attempt_with( + result, + evidence={**attempt.evidence, "execution_mode": "remote"}, + ) + elif rejection == "same_quality_evaluator": + evaluator = _actor(result) + quality_artifact = quality_artifact.model_copy(update={"created_by": evaluator}) + values = quality_result.model_dump(mode="python") + values["evaluated_by"] = evaluator + values["result_sha256"] = quality_result_fingerprint( + tenant_id=values["tenant_id"], + run_id=values["run_id"], + resource_version_id=values["resource_version_id"], + rule_version_ref=values["rule_version_ref"], + verdict=values["verdict"], + metrics=values["metrics"], + evidence_artifact_id=values["evidence_artifact_id"], + evaluated_by=evaluator, + evaluated_at=values["evaluated_at"], + ) + quality_result = QualityResult.model_validate(values) + elif rejection == "unbound_output": + success = _success_with_output(result, result.quality_evidence_artifact.artifact_id) + else: + lineage = lineage.model_copy( + update={"artifact_id": result.quality_evidence_artifact.artifact_id} + ) + + _prepare_running_bundle( + gateway, + result, + attempt=attempt, + quality_artifact=quality_artifact, + quality_result=quality_result, + lineage=lineage, + ) + with pytest.raises(GatewayValidationError, match="platform contract was rejected"): + gateway.finalize_run_success( + success, + expected_state_version=2, + actor_subject=_actor(result), + reason=FINALIZE_REASON, + ) diff --git a/docs/architecture-decisions/adr-082-public-lightweight-dataops-run.md b/docs/architecture-decisions/adr-082-public-lightweight-dataops-run.md new file mode 100644 index 00000000..16c9c550 --- /dev/null +++ b/docs/architecture-decisions/adr-082-public-lightweight-dataops-run.md @@ -0,0 +1,83 @@ +# ADR-082:公共轻量 DataOps Run 采用同步执行与独立成功证据门 + +**Status**: Accepted + +**Date**: 2026-08-17 + +**Decision owners**: Data Platform, DataOps, Platform Architecture, Governance + +**Related decisions**: ADR-020、ADR-026、ADR-081 + +## Context + +M3-34 已把显式 public/open source 写入不可变 Landing,但尚未证明平台能消费真实空间数据并形成可服务的数据产品版本。继续增加 readiness contract 无法回答核心问题:同一 Landing ResourceVersion 是否能被一个受控 Definition/Run 解包、标准化、质量检查、发布、记录血缘并由数据库裁决成功。 + +这条公共路径规模较小、无需等待或补偿,也不需要分布式调度。强行伪装成 DolphinScheduler attempt 会制造错误运行事实;新增服务、调度器或账本则会重复现有平台控制面。 + +## Options Considered + +| 方案 | 优点 | 缺点 | 结论 | +|---|---|---|---| +| 把本地执行记录为 DolphinScheduler | 可复用 migration 096 | observation 与真实执行器不符,破坏审计可信度 | 拒绝 | +| 为公共路径新增服务、表和调度器 | 可独立演进 | 重复 Resource/Run/Artifact/Quality/Lineage authority | 拒绝 | +| 复用控制账本,增加受限 synchronous success profile | 最小增量,运行事实真实,完整复用证据链 | 只适合短时本地任务,不提供生产调度能力 | **选择** | + +## Decision + +### 1. 复用既有平台对象 + +公共轻量执行复用 `Resource`、`ResourceVersion`、`PlatformDefinitionVersion`、`PlatformRun`、`Artifact`、`QualityResult`、`LineageEvent` 和 `FrameworkAttemptObservation`。逻辑 Definition 使用 `orchestration_class=synchronous`,不绑定具体 source;Run 的 immutable input binding 绑定实际 Landing ResourceVersion。 + +Definition 的发布主体和时间属于 Definition version 本身,不随每次 Run 变化。目标 Resource 只记录稳定 source ResourceURN,目标 ResourceVersion 再记录精确 source ResourceVersion。 + +### 2. 执行器形成真实内容寻址数据面 + +`public_dataops_run.py` 接受 GeoJSON 或安全 ZIP。ZIP extraction 拒绝路径穿越、符号链接、加密 entry、无效名称、超限 entry 和超限解压体积;Shapefile/GeoPackage 通过 GeoPandas 读取并归一化为 EPSG:4326。 + +发布前使用 Shapely 检查 feature count、null/empty/invalid geometry、geometry type 和 bbox。通过后写入不可覆盖的 content-addressed GeoJSON;独立 evaluator 生成另一条 content-addressed quality evidence。相同输入、配置和执行时间重放不得重写任一文件。 + +### 3. 本地 attempt 不伪装为调度器 + +真实 observation 固定为: + +- `framework_kind=legacy`; +- `evidence.schema=gda.public_dataops_attempt.v1`; +- `evidence.execution_mode=local_inline`; +- `observed_state=success`。 + +`legacy` 在此只表示已登记的本地 inline executor,不表示旧 Run 获得迁移权威,也不证明 DolphinScheduler execution。 + +### 4. 数据库继续拥有成功终局 + +已应用 migration 096 保持 checksum-frozen。migration 099 新增独立的 `finalize_synchronous_platform_run_success(...)`,完整保留 tenant、actor、state/CAS、replay、success evidence fingerprint、content-bound output、独立 passed quality 和 input-to-output lineage 检查,并额外强制 synchronous Run 与上述 local inline observation profile。 + +新函数为 `SECURITY DEFINER`,PUBLIC 无执行权,只有 `gda_control_gateway` 可调用。Gateway 先读取 Run;只有 synchronous Run 选择新函数,其他 orchestration class 继续调用 migration 096 的 DolphinScheduler finalizer。 + +## Consequences + +正面影响: + +- 首次形成真实 `Landing -> Run -> GeoJSON -> Quality -> Lineage -> succeeded` 公共数据垂直切片; +- 文件、平台对象和数据库终态均可确定重放,不新增第二套权威; +- 本地执行与调度器执行在 observation 和 finalizer 层明确分离; +- 可移植 Definition 能被同 tenant 的多个公共 source 重用。 + +限制与缓解: + +- 当前 serving 是本地 content-addressed GeoJSON,不是 active service revision;下一可执行里程碑必须增加消费者可访问的版本化发布、切换和回滚; +- 当前为同步短任务,不提供 schedule、补数、队列、checkpoint 或故障恢复;这些仍由 DolphinScheduler/Spark/Flink profile 验收; +- 本地文件系统不证明 object lock、跨节点 durability、backup/RPO/RTO 或生产 identity;`production_ready=false`; +- 本决策不改变重庆 protected admission,`source_content_admitted` 仍为 false。 + +## Verification + +- 13 个聚焦单元测试覆盖确定性输出、重放、质量、ZIP 安全、独立 evaluator、篡改和 CLI; +- 7 个 PostgreSQL 测试覆盖成功终态、无重复重放,以及错误 framework/schema/mode、同主体 evaluator、未绑定 output/lineage 的拒绝; +- Natural Earth 110m public-domain ZIP 实际产生 177 个 EPSG:4326 feature、1,065,275-byte GeoJSON 和完整账本链,重放未创建新文件或账本记录; +- 机器可读实测证据见 `docs/evidence/public-dataops-run-2026-08-17.json`。 + +## Revisit Triggers + +- 同步任务达到必须异步排队、取消、checkpoint 或恢复的规模; +- 版本化 GeoJSON 需要 active revision、HTTP gateway、cache、consumer impact 或 rollback; +- 同一 Definition 接入 PostGIS、STAC、DuckDB 或 lakehouse provider,需要 profile conformance 与 golden equivalence。 diff --git a/docs/evidence/public-dataops-run-2026-08-17.json b/docs/evidence/public-dataops-run-2026-08-17.json new file mode 100644 index 00000000..e7e04e50 --- /dev/null +++ b/docs/evidence/public-dataops-run-2026-08-17.json @@ -0,0 +1,97 @@ +{ + "schema": "gda.public_dataops_run.evidence.v1", + "evidence_date": "2026-08-17", + "profile": "public_open_lightweight", + "source": { + "label": "natural-earth-admin0-countries", + "uri": "https://naturalearth.s3.amazonaws.com/110m_cultural/ne_110m_admin_0_countries.zip", + "license_id": "public-domain", + "media_type": "application/zip", + "size_bytes": 214976, + "content_sha256": "0f243aeac8ac6cf26f0417285b0bd33ac47f1b5bdb719fd3e0df37d03ea37110", + "resource_urn": "gda://public-demo/dataset/natural-earth-admin0-countries", + "resource_version_id": "905bc019-71ef-5207-8392-cd2c46b7b5b7" + }, + "definition": { + "definition_urn": "gda://public-demo/definition/public-geojson-materialize", + "definition_version_id": "d92e3e09-b4ba-5aa7-b6cd-dba5d2e3c12d", + "definition_sha256": "c14d55805c00986fbdf737fc231715ca8c4dbce7512b19a9f7b17fff0d16d8f1", + "orchestration_class": "synchronous", + "capability_id": "public.geojson.materialize" + }, + "run": { + "run_id": "f3ae1ea2-1603-566c-809f-95a8401e295e", + "status": "succeeded", + "state_version": 3, + "attempt_observation_id": "1330a187-658c-5fdc-b186-de6b19602867", + "attempt_framework_kind": "legacy", + "attempt_schema": "gda.public_dataops_attempt.v1", + "execution_mode": "local_inline" + }, + "output": { + "resource_urn": "gda://public-demo/dataset/natural-earth-countries-serving", + "resource_version_id": "9431cf7c-c01b-5f26-a4cb-e05a2c56a6b8", + "artifact_id": "1a9099ce-08dc-5dab-ab62-5a6e218d87c2", + "media_type": "application/geo+json", + "crs": "EPSG:4326", + "content_sha256": "d81aeae51ab6759a692b21702cf159f1b94d0aec8e1bbe60828f47c65cf90f68", + "size_bytes": 1065275, + "feature_count": 177, + "geometry_types": { + "MultiPolygon": 29, + "Polygon": 148 + }, + "bbox_epsg4326": [ + -180.0, + -90.0, + 180.00000000000006, + 83.64513000000001 + ] + }, + "quality": { + "verdict": "passed", + "quality_result_id": "4badc6ad-4081-5fd6-b05e-13900970e230", + "evidence_artifact_id": "3762ba75-55fd-5c98-a149-e719926dbc18", + "evidence_content_sha256": "ce075670c559de543c97007d5e5cb3aed449d21bde82cab60135c6cf3225423f", + "evaluated_by": "workload:public-quality", + "independent_evaluator_verified": true, + "null_geometry_count": 0, + "empty_geometry_count": 0, + "invalid_geometry_count": 0 + }, + "lineage": { + "lineage_event_id": "e705fa4c-57c3-54a4-8c1e-a8bb5f18cf95", + "source_resource_version_id": "905bc019-71ef-5207-8392-cd2c46b7b5b7", + "target_resource_version_id": "9431cf7c-c01b-5f26-a4cb-e05a2c56a6b8", + "output_artifact_id": "1a9099ce-08dc-5dab-ab62-5a6e218d87c2" + }, + "replay": { + "output_created": false, + "quality_created": false, + "ledger_completed": true, + "platform_run_count": 1, + "platform_run_event_count": 4, + "run_artifact_count": 2, + "attempt_observation_count": 1, + "quality_result_count": 1, + "lineage_event_count": 1 + }, + "verification": { + "result_bundle_valid": true, + "output_file_sha256_valid": true, + "quality_file_sha256_valid": true, + "focused_unit_tests_passed": 13, + "postgres_integration_tests_passed": 7, + "local_inline_execution_verified": true, + "geojson_serving_projection_verified": true, + "postgres_ledger_verified": true, + "dolphinscheduler_execution_verified": false, + "postgis_serving_verified": false, + "stac_serving_verified": false, + "active_service_revision_verified": false, + "rollback_verified": false + }, + "production_ready": false, + "protected_chongqing_admission_unchanged": true, + "protected_chongqing_source_content_admitted": false +} diff --git a/docs/roadmap.md b/docs/roadmap.md index 1bae86ff..44d37465 100644 --- a/docs/roadmap.md +++ b/docs/roadmap.md @@ -1,6 +1,6 @@ # GIS Data Agent — 总体架构 Roadmap -**Last updated**: 2026-07-24 +**Last updated**: 2026-08-17 **Status**: Architecture reset, authoritative mainline @@ -446,6 +446,7 @@ AR-0 Architecture / Schema / Runtime Truth - M3-32 protected admission attestation intake:固定外部 attestation 的 M3-31 logical/file fingerprint binding、15 项逐项 SHA-256 证明、受保护 verifier/source binding、24 小时 freshness、七天 validity 上限和八项 no-payload/no-mutation checks;`evaluate` 只产生 fingerprinted readiness report,不创建任何 Landing、ResourceVersion、PlatformRun、scheduler submission 或 provider mutation authority。 - M3-33 protected admission verifier workflow:只允许从 `main` 手工触发受保护 `chongqing-admission` environment 的专用 verifier runner,消费 metadata-only secret bundle,绑定 exact M3-31 fingerprints,执行 M3-32 evaluate/verify,并以 GitHub OIDC provenance attestation 固定 input/report;workflow 不含 source scan、Landing/ResourceVersion/PlatformRun 创建、scheduler submission 或 provider client,环境/runner/15 项真实 attestation 未 provision 前仍 blocked。 - M3-34 public/open-source immutable Landing slice:以显式 HTTPS source、license、owner、expected SHA-256 和 controlled actor 将真实 public-domain bytes 写入 content-addressed local Landing;复用现有 Resource/ResourceVersion/Artifact authority,通过单事务 gateway registration 保证幂等 replay 和冲突回滚;Natural Earth 110m smoke 已通过,public profile 仍非 production-ready,也不改变 Chongqing protected admission。 +- M3-35 public lightweight DataOps run:复用既有 Definition/Run/Artifact/Quality/Lineage control ledger,以同步本地 executor 将 M3-34 Natural Earth ZIP 实际解包并归一化为 EPSG:4326 content-addressed GeoJSON;177 个 feature 的独立质量证据、legacy/local-inline attempt、input-to-output lineage 和 PostgreSQL success verdict 已闭环,文件与账本重放无重复。该路径不证明 DolphinScheduler、STAC/PostGIS active service、rollback 或 production readiness,也不改变 Chongqing protected admission。 - SourceDefinition、CredentialReference、SourceCapability、SyncDefinition/Version、SyncRun、Cursor/Watermark、SchemaDriftEvent 和 Reconciliation。 - 数据库、对象存储/空间文件、HTTP/STAC 三类代表 source 的连接、凭据、连通、发现、preview、profile 和 owner 登记。 - 全量/增量微批的 Append/Overwrite/Merge 策略,以及至少一个真实 CDC 或事件流 source 通过 Flink 写入版本化 Bronze;覆盖 watermark/offset、checkpoint、迟到/乱序、源端删除、幂等、对账、重放和失败恢复。 @@ -693,7 +694,7 @@ Golden checks 至少覆盖: 3. 冻结 ResourceURN、ResourceVersion、PlatformDefinition/PlatformRun/FrameworkAttemptObservation/Artifact/LineageEvent、SubjectContext 与 storage/table/compute provider 最小合同。 4. 分阶段实现 `gda-metadata-fabric-bridge`:M1 只读 mapping/reconciliation、M2a 本地 foundation/重启连续性、M2b-1 本地三存储恢复、M2b-2 隔离 versioned/Object-Locked repository round-trip、M2b-3 本机双集群 + Kubernetes 外 COMPLIANCE repository + 独立 writer/reader、M2c-1 provider-native metrics、M2c-2 临时 OTel Collector + JSON Exporter 的双周期本地 pipeline、M2c-3 本地单 job scrape 故障检测/配置恢复/完整清理 evidence、M2c-4 绑定 source revision 的 production observability readiness contract,以及 M2d-1 本地 kindnet 跨节点 NetworkPolicy enforcement 已验证;M2c-4 当前仍有 20 项 blockers,M2d-1 也未验证生产 provider policy 或 tenant isolation。下一步完成 source host/cluster 外的生产 bucket、KMS/TLS/workload identity、source-loss recovery 与 RPO/RTO,并批准 metrics backend、retention、OTel/TLS、tenant、alert/SLO/owner 后在受保护环境验证持续采集、存储、查询、真实告警投递、runbook 响应和 provider NetworkPolicy;再推进 OIDC、upgrade/rollback、registry provenance 和 owner/runbook;之后才进入 M3 ingestion/OpenLineage/conformance。 5. 实现 `gda-orchestration-gateway`、DolphinScheduler process/task/schedule/complement/worker-group、Spark/Flink provider task adapter 和故障注入;不再开发新的 lease/queue/scheduler。 -6. M3-28 已冻结全量重庆真实源的 path-free physical/metadata admission baseline,M3-29 已建立 metadata-only extraction provenance gap baseline,M3-30 已选择 `bishan_land_use_dltb_local` 并把八项治理输入固化为 fail-closed pending records,M3-31 已将三层证据统一成 protected admission readiness contract,M3-32 已固定 protected attestation intake/evaluate/verify boundary,M3-33 已固定受保护 workflow execution/provenance boundary,M3-34 已用 Natural Earth public-domain bytes 验证 immutable Landing 与 ledger registration path;下一步用 public Landing 绑定最小 DataOps PlatformRun 并物化轻量 serving slice,同时 provision 专用 environment/runner 并取得 Chongqing 的 15 项外部 attestation,未获批准前不得 content admission。 +6. M3-28 至 M3-33 已冻结重庆 protected admission 的 metadata-only fail-closed 边界,M3-34 已验证 Natural Earth immutable Landing,M3-35 已实际完成 public Landing -> synchronous PlatformRun -> EPSG:4326 GeoJSON -> independent QualityResult -> LineageEvent -> succeeded ledger 的轻量垂直切片及无重复重放;下一步直接实现该 public ResourceVersion 的版本化服务 revision、消费者读取、原子切换与 rollback,不再增加与执行无关的 readiness-only 合同。同时 provision 重庆专用 environment/runner 并取得 15 项外部 attestation,未获批准前不得 content admission。 7. 冻结 Default Lakehouse、Cloud Managed、Lightweight Integrated profiles;以统一 Run 完成默认 MinIO/Iceberg/Spark/Flink、轻量 PostGIS/DuckDB 和 Azure 代表 adapter 的 conformance smoke。 8. 实现跨 profile 的 Raw -> ODS -> DIM/DWD -> DWS -> ADS 通用生产、质量、发布、回滚和 golden equivalence。 9. 建立 DataProductBlueprint、模型版本和 Visual/SQL/Notebook 共用 definition 的 Build 工作台,打通 preview、test、publish、approval 和 rollback。 @@ -734,7 +735,7 @@ AR-4 parity/control gate 退出前暂停以下主线扩张: |---|---|---| | AR-0 Architecture/Schema/Runtime Truth Freeze | `in_progress` | 全环境 schema/config fingerprint、迁移 fail-closed、事实清单、storage/compute/GIS serving provider profile/capability、ADR-017 benchmark、owner/SLO 和首条数据/服务验收集冻结 | | AR-1 Unified Metadata + Orchestration Control Planes | `in_progress` | controlled gateway、DolphinScheduler adapter、Metadata Fabric M1/M2a/M2b、M2c-1 provider metrics、M2c-2 本地临时 OTel pipeline、M2c-3 本地 scrape failure/recovery、M2c-4 production observability readiness contract 与 M2d-1 本地跨节点 NetworkPolicy enforcement 已验证,生产观测和生产 policy/tenant isolation 仍 blocked;下一证据是 source host/cluster 外的生产 recovery、持久 metrics backend/TLS/tenant/真实 alert delivery/SLO、OIDC、受保护 provider NetworkPolicy、升级回滚/registry provenance,以及受控 ingestion/replay 与无双写验收 | -| AR-2 Source/Ingestion + Geospatial Lakehouse Vertical Slice | `in_progress` | M3-28 admission baseline、M3-29 extraction provenance gap baseline、M3-30 first-candidate governance gate、M3-31 protected admission readiness、M3-32 attestation intake、M3-33 protected verifier workflow contract 与 M3-34 public immutable Landing slice 已完成;public profile 已有真实 bytes/manifest/ledger path,但 DataOps PlatformRun、Bronze/Silver/Gold、serving 和 rollback 尚未完成;checked Chongqing baseline 仍为 `admission_eligible=false`,下一证据是 public Landing 的最小 DataOps run,以及专用 environment/runner provisioning、15 项外部 attestation 和首次 protected verifier run | +| AR-2 Source/Ingestion + Geospatial Lakehouse Vertical Slice | `in_progress` | M3-28 至 M3-33 的重庆 protected admission 边界、M3-34 public immutable Landing 与 M3-35 public lightweight DataOps run 已完成;Natural Earth 已形成真实 synchronous PlatformRun、EPSG:4326 GeoJSON、独立质量、血缘、PostgreSQL succeeded verdict 和无重复 replay,但 Bronze/Silver/Gold、active service revision、消费者访问、rollback 与 DolphinScheduler execution 尚未完成。下一证据是 public GeoJSON 的版本化服务激活/切换/回滚实操,以及重庆专用 environment/runner、15 项外部 attestation 和首次 protected verifier run | | AR-3 Data Product Engineering + Governance Workbench | `planned` | Blueprint、模型、Visual/SQL/Notebook、DataOps CI/CD、质量/安全/审批共用 definition 和产品生命周期 | | AR-4 Asset/GIS Service/Spatial Experience Operations | `planned` | Service Control Plane、Features/Tiles/MVT/COG/STAC/export 及条件 legacy OGC/3D/EDR provider、Gateway/权限/缓存、原子切换/回滚、Discover/Operate/Govern 和无 LLM 多入口通过 conformance/parity/control gate | | AR-5 AgentOps Runtime + UX Uplift | `planned` | DataOps parity/control 通过;Agent bundle eval、deployment、online observation、incident/rollback 和 uplift gate | diff --git a/docs/system-of-record-matrix-2026-07-24.md b/docs/system-of-record-matrix-2026-07-24.md index 81f5c696..851b6a16 100644 --- a/docs/system-of-record-matrix-2026-07-24.md +++ b/docs/system-of-record-matrix-2026-07-24.md @@ -2,7 +2,7 @@ 日期:2026-08-17 -阶段:AR-0 `in_progress`;AR-1 gateway、成功终局 evidence gate、DolphinScheduler adapter sandbox POC、Metadata Fabric M1/M2、M2c-4/M2d-2 production readiness contracts、M3-1/M3-2、M3-3 local binding ledger、M3-4 local OpenLineage wire delivery、M3-5 local OpenMetadata bounded identity、M3-6 local Gravitino Basic bounded identity 与 M3-7 production identity readiness contract 已验证;AR-2 M3-28 重庆真实源 metadata-only admission baseline、M3-29 extraction provenance gap baseline、M3-30 first-candidate governance gate、M3-31 protected admission readiness、M3-32 attestation intake、M3-33 protected verifier workflow contract 与 M3-34 public immutable Landing slice 已建立,public profile 的 DataOps run/serving、重庆内容准入、生产 provider ingestion、生产观测、生产 policy/tenant isolation、生产 identity attestation 和生产切换仍 `in_progress` +阶段:AR-0 `in_progress`;AR-1 gateway、成功终局 evidence gate、DolphinScheduler adapter sandbox POC、Metadata Fabric M1/M2、M2c-4/M2d-2 production readiness contracts、M3-1/M3-2、M3-3 local binding ledger、M3-4 local OpenLineage wire delivery、M3-5 local OpenMetadata bounded identity、M3-6 local Gravitino Basic bounded identity 与 M3-7 production identity readiness contract 已验证;AR-2 M3-28 至 M3-33 的重庆 protected admission 边界、M3-34 public immutable Landing 与 M3-35 public lightweight DataOps run 已建立,public profile 已形成真实 synchronous Run/GeoJSON/quality/lineage/success replay,但 active service revision/rollback、重庆内容准入、生产 provider ingestion、生产观测、生产 policy/tenant isolation、生产 identity attestation 和生产切换仍 `in_progress` 适用分支:`main` @@ -21,7 +21,7 @@ | 部署配置策略 | Compose/K8s/进程环境;`platform_truth.CONFIG_SPECS` 定义关键类型与策略;DolphinScheduler worker 有默认零副本、外部 ConfigMap/Secret 驱动的 Kustomize 模板、静态 validator、staging activation preflight 和受保护的单副本 activation admission/workflow | `.env` 仅补默认;脱敏 snapshot、Secret key attestation、未扩容 Deployment、`ready_for_activation` 和未执行的 activation workflow 都是观测/变更能力 | 版本化 DeploymentProfile + secret reference;部署环境始终优先;模板、preflight 或 admission 通过都不等于环境已启用 | Platform/SRE/Security | AR-0,部分实现;worker 激活合同已验证、真实运行待审批 | | 环境发布与晋级 | publisher `31862363442`、protected verifier `31862984294` 与 staging deploy/observe `31863077257` 已将 `main@5fffc85`、GHCR digest、attested release、cluster/namespace identity 和 live revision 绑定;schema/config/runtime/health/rollout 通过,golden slice 缺失使 promotion fail closed | 旧 mainline、feature branch、CI artifact、JSON、离线 report、单独的 staging deployment 或人工批准都不能成为 production 发布权威 | 由受保护 environment 的 DeploymentRevision 绑定 OCI、provenance artifact、release manifest、golden slice 与全部 live verdict | Platform/SRE/Security/Repository Owner | AR-1 真实 staging 已部署 -> golden slice/production exit gates 待完成 | | 后台运行时清单 | `platform_truth.RUNTIME_INVENTORY` 是代码层登记;`gda_control` 已有受控 PlatformRun 写入口;DolphinScheduler adapter、tenant-scoped managed command worker 与受保护的单副本激活边界已有代码和测试,但 worker 尚未在 staging 运行;M2b recovery runner、M2c-1 provider probe、M2c-2 `_OtelPortForward`、M2c-3 failure rehearsal 与 M2d-1 NetworkPolicy rehearsal 均登记为 `local_verification_only`,不是 scheduler、worker、持续监控、生产 policy controller 或状态权威 | AST primitive report、worker status JSON、FrameworkAttemptObservation、DolphinScheduler instance state、本地 recovery/metrics/network-policy evidence | PlatformRun ledger 唯一登记最终状态;framework/provider attempt 只能回报观测;worker status 仅为进程健康投影;本地演练进程与 evidence 不得变成生产控制器、监控后端或 tenant-isolation 权威 | Platform Architecture | AR-1 adapter/worker/activation 合同已验证 -> staging 运行待审批;metadata recovery/metrics/policy runner 仅本地验证 | -| 原始文件/对象 | M3-28 已以 path-free metadata manifest 记录重庆真实源 archive/extracted fingerprints、source groups、asset profiles 与治理 blockers;M3-29/M3-30 固定派生证明缺口和首条候选的八项 pending 治理决策;M3-31/M3-32 将其绑定为 15 项 pending requirements 与 protected intake contract;M3-33 已定义 protected verifier workflow,但未 provision 或执行;M3-34 已用 Natural Earth public-domain ZIP 验证 content-addressed local Landing、manifest、ResourceVersion 和 input Artifact 的实际路径;重庆 `source_content_admitted=false`,public Landing 也尚未绑定 DataOps Run 或 production serving | public-source Landing bytes/manifest、Resource/ResourceVersion/Artifact、admission/provenance/governance/readiness/evaluation/workflow evidence、archive/extracted comparison、临时上传、下载缓存、预览文件 | public/open profile 可由受控 DataOps definition/Run 消费;重庆仍必须在专用 environment/runner 中由 15 项外部 attestation 产生受保护 verifier report,并经独立 admission decision 后才可建立其 immutable Landing authority;本地 scratch、checked contract 与 workflow 文件均不可替代生产 Landing | Data Platform | AR-2 `in_progress` | +| 原始文件/对象 | M3-28 至 M3-33 固定重庆 metadata-only protected admission 边界但未 provision/执行;M3-34 已用 Natural Earth public-domain ZIP 验证 content-addressed Landing,M3-35 已将其 ResourceVersion 绑定 synchronous DataOps Run 并生成 content-addressed EPSG:4326 GeoJSON 与独立 quality evidence;重庆 `source_content_admitted=false`,public 输出仍是本地轻量 serving projection,不是 production object/service authority | public-source Landing bytes/manifest、Resource/ResourceVersion/Artifact、public GeoJSON/quality projection、admission/provenance/governance/readiness/evaluation/workflow evidence、临时上传、下载缓存、预览文件 | public/open profile 只允许受控 Definition/Run 和内容寻址 projection;active service 必须另经版本化发布/切换/回滚。重庆仍必须由 15 项外部 attestation 和独立 admission decision 建立 immutable Landing authority;本地 scratch、checked contract 与 workflow 文件均不可替代生产 Landing | Data Platform | AR-2 `in_progress` | | 湖仓表与 snapshot | Iceberg/STAC/S3A 有局部实现,尚无通用发布权威 | STAC item、GeoParquet export | Iceberg catalog snapshot 是分析表版本权威;对象是物理内容,STAC 是发现投影 | Data Platform | AR-2 | | 在线空间数据 | PostGIS 业务表是当前编辑/查询事实,部分临时表混入 | Martin MVT、API JSON、导出文件 | 已批准 DataProductVersion 物化到 PostGIS;不能由瓦片或临时表反向定义产品版本 | GIS/Data Platform | AR-2 -> AR-4 | | 数据资产身份与版本 | `gda_control.resource/resource_version` 已实现 identity、hash、predecessor、tenant FK 和幂等 gateway 写入;`agent_data_assets`、`agent_asset_versions` 仍是兼容写路径 | UI catalog、search index、STAC | GDA ledger 管身份与版本绑定;旧行只有在 tenant、authority identity、checksum 和 version evidence 完整时才可形成 eligible plan;OpenMetadata 管治理目录,Gravitino 管技术对象映射 | Metadata Platform | AR-1 gateway 已验证 -> 生产切换待验收 | @@ -29,10 +29,10 @@ | 治理目录 | M1 已冻结 OpenMetadata table ref/reconciliation;M2a 已运行 OpenMetadata `1.13.1` + 独立 PostgreSQL/OpenSearch;M2b 已完成本地新 PVC、锁定 repository 与独立 kind cluster 恢复;M2c-1 已验证原生指标,M2c-2 已验证 OpenMetadata -> OTel 本地双周期抓取,M2c-3 已证明 Gravitino scrape 故障时 OpenMetadata 保持 `up=1` 并可恢复;M2c-4 profile 合同有效但 production observability gate 仍 blocked;M2d-1 已验证本地 kindnet 跨节点合成流量,不是 OpenMetadata provider policy;M2d-2 NetworkPolicy profile 合同有效但 production provider/tenant gate 仍 blocked;M3-1 已生成 owner/domain/tag/generic-lineage intent 与 OpenLineage COMPLETE candidate;M3-2 已用 bootstrap admin 本地创建目标并回读真实 UUID,第二次 replay 为 `no_op/0 mutations`,OpenLineage 未发送;M3-3 已将该 UUID 经 evidence gate 追加到本地 GDA binding ledger;M3-4 已将精确 OpenLineage candidate 经 tenant-scoped outbox 投递到本地 loopback HTTP receiver,并验证 at-least-once + receiver idempotency | 搜索/页面视图、合成 OpenMetadata response、本地 sandbox/recovery/repository/metrics/network-policy/ingestion observation、projection/provider evidence、binding ledger/OpenLineage candidate、lineage outbox/receipt、receiver event 与 readiness report | OpenMetadata 为 owner/glossary/classification/quality discoverability 权威;GDA ledger 保留审批与 provider identity 关系;outbox 只拥有投递状态,receiver 拥有接收状态;本地 admin apply、临时 ledger、loopback receiver 和 OpenLineage event 均不能改变 ResourceVersion、Run 或审批权威 | Governance | AR-1 M1 + 本地 M2a/M2b/M2c/M2d-1/M2d-2 + M3-2/M3-3/M3-4 local evidence 已验证 -> 外部生产恢复/持续 metrics/受保护 policy/最小权限 ingestion/生产持久 binding/受保护 production receiver/live OpenLineage 待执行 | | 血缘 | `gda_control.lineage_event` 已实现 immutable version edge 和幂等 gateway ingest;`agent_asset_lineage` 旧记录仍是可变 asset edge | OpenMetadata lineage graph、UI DAG | 只有 source/target ResourceVersion 与 event checksum 证据完整的旧记录可形成 eligible plan;目录图只作可重建投影 | Data Platform | AR-1 gateway 已验证 -> adapter 待接入 | | Definition | `gda_control.platform_definition_version` 已绑定 definition ResourceVersion、完整逻辑 hash 和原子 gateway registration;3.4.2 adapter 可编译、创建并上线 provider DAG;binding 已以 append-only `execution_plan` Artifact 持久化并可按 tenant + artifact UUID 读取,旧 workflow/template/YAML 仍在写入 | 编辑器状态、DolphinScheduler DAG/definition | 旧 workflow 必须规范化并完整 hash 后才可形成 PlatformDefinitionVersion;provider binding 作为 ExecutionPlanArtifact/evidence,不可反写 definition | DataOps | AR-1 binding persistence 代码已验证 -> staging 调用链待验收 | -| Run 最终状态 | `gda_control.platform_run/event` 已实现受控 submit/read/CAS;通用 transition 已禁止 `succeeded`,专用数据库 finalizer 只接受精确 workload、DolphinScheduler success observation、内容匹配 output、独立 passed QualityResult/evidence 和 input-to-output lineage;adapter standalone API path 已验证,但端到端 staging 尚未完成,legacy 路径继续运行 | Redis progress、日志、DolphinScheduler state、attempt observation | 旧 run 到 PlatformRun 永久 prohibited;已有 PlatformRun correlation 时才可转为 observation;provider 终态只进入 `reconciling`,ledger 经证据门唯一裁决成功 | DataOps/AgentOps | AR-1 success authority 本地/PostgreSQL 已验证 -> staging/生产切换待验收 | +| Run 最终状态 | `gda_control.platform_run/event` 已实现受控 submit/read/CAS,通用 transition 禁止 `succeeded`;migration 096 只接受 DolphinScheduler success profile,migration 099 只接受 synchronous Run + `legacy/gda.public_dataops_attempt.v1/local_inline` profile,两者都强制内容匹配 output、独立 passed QualityResult/evidence 和 input-to-output lineage。M3-35 Natural Earth Run 已在 PostgreSQL 达到 `succeeded`/state 3 并幂等 replay;DolphinScheduler 端到端 staging 尚未完成 | Redis progress、日志、DolphinScheduler state、attempt observation | 旧 run 到 PlatformRun 永久 prohibited;已有 PlatformRun correlation 时才可转为 observation;provider/profile observation 必须走各自受限 finalizer,ledger 经证据门唯一裁决成功 | DataOps/AgentOps | AR-1/AR-2 本地 PostgreSQL 已验证 -> staging/生产切换待验收 | | 调度与补数 | APScheduler、自进化 scheduler 和调用方定时逻辑并存;DolphinScheduler POC 只验证 manual start/list/variables/STOP | UI schedule 列表 | DolphinScheduler 管 DataOps schedule/complement;Temporal 只管需要 durable signal/compensation 的 Agent/GWM workflow | DataOps/AgentOps | AR-1 manual correlation 已验证;schedule/complement/failover 待验收 | | 事件交付 | Standards outbox 已数据库耐久;`gda_control.platform_command_outbox` 已为 DolphinScheduler dispatch/reconcile 提供 tenant RLS、lease claim、幂等 callback、薄 consumer library、managed worker process 和受保护的单副本 staging activation boundary;其他 WebSocket/bot/feedback 多为 best effort | command delivery status、消费者 claim、worker status JSON、activation/readiness evidence、WebSocket 消息 | command/event 与源事实同事务入 outbox,幂等 consumer 交付;worker status、outbox 状态、activation evidence、缓存或 socket 都不是 Run/业务权威 | Platform/Integrations | AR-1 command delivery/activation 代码已验证 -> worker/callback live evidence 待完成 | -| 质量结果 | `gda_control.quality_result` 已提供 tenant RLS、append-only gateway 写入,绑定 Run、output ResourceVersion、rule version、verdict、metrics、evidence Artifact 和独立 evaluator;standards、QC、MMFE 专项结果仍未迁移 | dashboard、OpenMetadata quality summary | GDA ledger 保存产品终局所需的不可变 verdict/evidence;OpenMetadata 与 UI 只作可重建发现投影;旧结果缺稳定版本和证据时不得升级为终局依据 | Governance/DataOps | AR-1 最小成功证据已验证 -> 真实规则/staging 待接入 | +| 质量结果 | `gda_control.quality_result` 已提供 tenant RLS、append-only gateway 写入,绑定 Run、output ResourceVersion、rule version、verdict、metrics、evidence Artifact 和独立 evaluator;M3-35 已对 Natural Earth 177 个 feature 实测 geometry/null/empty/invalid/type/bbox 并形成独立 passed evidence,standards、QC、MMFE 专项结果仍未迁移 | dashboard、OpenMetadata quality summary | GDA ledger 保存产品终局所需的不可变 verdict/evidence;OpenMetadata 与 UI 只作可重建发现投影;旧结果缺稳定版本和证据时不得升级为终局依据 | Governance/DataOps | AR-2 public lightweight 实测 -> 领域规则/staging 待接入 | | 标准与语义定义 | `std_*`、semantic registry 和 YAML 共同存在,生命周期未统一 | prompt/context、搜索索引 | 版本化 Standard/SemanticDefinition 经审批后为权威;Agent context 只消费批准版本 | Governance | AR-1 -> AR-3 | | 身份与权限 | Chainlit user 可显式绑定 tenant;versioned API 从认证 principal 派生 SubjectContext;`gda_control_gateway` 是 non-login/non-bypass 最小权限角色;Run 可引用强类型 PolicyDecision/Approval Artifact;M3-5/M3-6 分别验证本地 provider scoped grant、越权拒绝和 credential rotation/revocation;M3-7 已冻结生产 OIDC/workload/tenant binding、TLS、持久 catalog 与 attestation contract,但 40 个外部输入仍 blocked | session/cache、前端菜单权限、本地 provider identity evidence、pending profile 与合成 readiness report | IdP/workload identity 提供真实 service identity;PolicyDecision/Approval 继续绑定不可变资源与 execution plan;只有 fresh protected attestation 可派生双 provider production identity claims,profile、Basic/JWT evidence 或人工批准均不可替代 | Security | AR-1 local identities + production readiness contract 已验证 -> protected 双 provider IAM/attestation 待执行 | | GIS 服务定义与 active revision | Martin、REST/MVT/STAC endpoints 和配置直接暴露 | Ingress、tile cache、客户端图层 | GIS Service Control Plane 管 Service/Layer/Style/TMS/DeploymentRevision;provider/Gateway 仅执行 | GIS Platform | AR-4 | @@ -69,6 +69,7 @@ 23. M3-32 只允许受保护 verifier 消费绑定 M3-31 logical/file fingerprints 的外部 attestation;15 项 requirement、八项 no-payload/no-mutation checks、freshness、expiry 和 verifier/source binding 任一缺失即 blocked。完整 attestation 只能使 fingerprinted report 的 `admission_eligible=true`,不创建 Landing object、ResourceVersion、PlatformRun、scheduler submission、provider mutation、content admission 或 production readiness 权威;合成测试 attestation 不计入生产证据。 24. M3-33 只允许 `main` 上手工触发的受保护 `chongqing-admission` workflow 在专用 runner 中消费 metadata-only secret bundle,绑定 M3-31 exact logical/file fingerprints,执行 M3-32 evaluate/verify 并 attested/upload input/report;workflow 不读取 source payload、不调用 connector/provider/scheduler、不创建 Landing object、ResourceVersion、PlatformRun 或任何 ingestion authority。environment、runner、reviewer policy、secret rotation 与真实 15 项 attestation 未 provision 前,workflow contract 不计入 protected admission evidence。 25. M3-34 public/open Landing 只允许显式 HTTPS source、license、owner、expected SHA-256 和 controlled actor;payload 与 manifest 以 content-addressed no-overwrite 文件写入,Resource、ResourceVersion、input Artifact 必须由同一 gateway transaction 登记并可幂等 replay。该 public profile 不创建 PlatformRun、不代表 Chongqing protected admission、不证明 object-lock、provider identity、production serving 或 production readiness。 +26. M3-35 public lightweight DataOps 只允许 synchronous Run 消费已验证的 public/open Landing;本地 attempt 必须明确登记为 `legacy` + `gda.public_dataops_attempt.v1` + `local_inline`,不得伪装成 DolphinScheduler。GeoJSON output、quality evidence、QualityResult 和 LineageEvent 必须内容寻址并由 migration 099 的独立终态函数裁决;该 slice 不形成 active service revision、rollback、DolphinScheduler execution、Chongqing content admission 或 production readiness。 ## 已建立的 AR-0/AR-1 entry 证据 @@ -107,6 +108,7 @@ - AR-2 M3-32 已建立 protected admission attestation intake/evaluate/verify contract:外部输入必须绑定 M3-31 logical/file fingerprints,逐项提供 15 个 SHA-256 attestation、八项 protected checks、verifier/source binding 和受控 freshness/expiry;checked baseline 的 `readiness_valid=true`、`attestation_valid=false`、`admission_eligible=false`,所有 content/authority/production claims 继续为 `false`。合成完整 attestation 仅覆盖 evaluator 单元测试,不计入真实准入证据。 - AR-2 M3-33 已建立 protected admission verifier workflow contract:workflow 仅从 `main` 手工触发、绑定专用 environment/runner、以 restrictive umask 解码 metadata-only attestation secret、执行 M3-32 evaluate/verify 并通过 GitHub OIDC provenance attestation 固定 input/report;checked workflow 未执行,专用环境、runner、reviewer policy、secret rotation 与真实 15 项 attestation 均未 provision,不形成 content admission 或生产 authority。 - AR-2 M3-34 已用 Natural Earth 110m public-domain ZIP 完成 public/open immutable Landing smoke:214,976 bytes 的 content SHA、content-addressed payload、只读 manifest、ResourceVersion、input Artifact、replay 和 verify 均通过;该 smoke 的 `ledger_registered=false` 只表示 CLI 未连接生产数据库,独立 PostgreSQL 集成测试已验证同一原子 registration/replay/conflict rollback contract,public profile 仍 `production_ready=false`。 +- AR-2 M3-35 已将同一 Natural Earth ZIP 实际解包并归一化为 EPSG:4326 content-addressed GeoJSON:177 features、148 Polygon、29 MultiPolygon、1,065,275 bytes,output SHA 为 `d81aeae51ab6759a692b21702cf159f1b94d0aec8e1bbe60828f47c65cf90f68`;Run `f3ae1ea2-1603-566c-809f-95a8401e295e` 在 PostgreSQL 达到 `succeeded`/state 3,独立 QualityResult、quality Artifact、LineageEvent 和 local-inline attempt 均为单条,重放 `output_created=false`、`quality_created=false` 且账本无重复。该实测不证明 DolphinScheduler、active service revision、STAC/PostGIS serving、rollback 或 `production_ready`,重庆 protected admission 保持不变。 ## 下一验收证据 @@ -115,5 +117,5 @@ - 为受保护 activation 提供真实 ConfigMap snapshot、Secret key attestation、provider identity 和 reviewer approval,随后完成 managed outbox worker/provider callback 单副本实际部署、唯一 worker ID、status/lease 故障恢复和无双写证据; - 首条真实图斑链对 golden slice 的 output hash、独立质量结果/evidence、血缘、发布 revision 和 rollback 演练; - OpenMetadata/Gravitino 的 source host/cluster 外生产 backup account/bucket、KMS/TLS/workload identity、PITR/source-loss recovery、RPO/RTO、OIDC、受保护环境 provider NetworkPolicy/tenant isolation、upgrade/rollback、registry provenance、持续 metrics backend/retention/query、真实 alert delivery/SLO owner/runbook,以及受保护 PolicyDecision/Approval、双 provider 最小权限 ingestion、生产持久 binding、受保护 production OpenLineage receiver、无双写 read-back 和 conformance;M1 fixture、M2 本地 evidence/readiness contracts、M3-1 projection candidate、M3-2 local replay、M3-3 临时 binding ledger、M3-4 loopback delivery、M3-5/M3-6 本地临时 provider identity 与 M3-7 pending profile/合成 attestation 均不计入生产退出门; -- M3-34 的下一证据是将 public Landing ResourceVersion 绑定最小 DataOps definition/PlatformRun,完成一次轻量 profile 的解包、质量、GeoJSON/STAC 或 PostGIS serving、lineage、replay 和 rollback;M3-29/M3-30/M3-31/M3-32/M3-33 的重庆下一证据仍必须 provision 专用 `chongqing-admission` environment/runner/reviewer/rotation policy,补齐 operator/tool/command、modified/additional manifest、archive-to-working-set attestation、owner/license/retention/access/privacy-sensitivity/standard-version/DataSLO/golden-result 八项签名决策与 fresh protected attestation,并由 M3-33 workflow 按 M3-32 intake contract 重新计算 admission eligibility;metadata-only admission/provenance/governance/readiness/evaluation/workflow evidence 是不可变检查结果,不能被编辑或重解释为内容批准、Landing authority 或 ingestion authorization; +- M3-35 的下一证据是将 public GeoJSON ResourceVersion 绑定版本化 service revision,完成真实消费者读取、原子 active pointer 切换、旧 revision 回滚和无重复 replay;不得以新增 readiness-only 合同替代执行。M3-29/M3-30/M3-31/M3-32/M3-33 的重庆下一证据仍必须 provision 专用 `chongqing-admission` environment/runner/reviewer/rotation policy,补齐 operator/tool/command、modified/additional manifest、archive-to-working-set attestation、owner/license/retention/access/privacy-sensitivity/standard-version/DataSLO/golden-result 八项签名决策与 fresh protected attestation,并由 M3-33 workflow 按 M3-32 intake contract 重新计算 admission eligibility;metadata-only admission/provenance/governance/readiness/evaluation/workflow evidence 是不可变检查结果,不能被编辑或重解释为内容批准、Landing authority 或 ingestion authorization; - DolphinScheduler/Temporal sandbox 的独立数据库、备份恢复、身份、版本和升级责任证明;DolphinScheduler standalone/H2 不计入此退出门。 From ec056cd932597000dacb4817c5e090cbc99e917f Mon Sep 17 00:00:00 2001 From: Ning Zhou Date: Mon, 17 Aug 2026 21:45:36 +0400 Subject: [PATCH 2/2] test: recognize synchronous success migration --- data_agent/test_platform_contracts.py | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/data_agent/test_platform_contracts.py b/data_agent/test_platform_contracts.py index 495fbd9e..a0dbd704 100644 --- a/data_agent/test_platform_contracts.py +++ b/data_agent/test_platform_contracts.py @@ -471,9 +471,7 @@ def test_control_ledger_contract_and_migration_catalog_are_valid(): assert report["status"] == "valid" assert report["contract_count"] == 16 assert report["migration"]["sha256"] == migration["checksum"] - assert migrations[-1]["migration_id"] == ( - "098_metadata_fabric_openlineage_delivery" - ) + assert migrations[-1]["migration_id"] == "099_synchronous_success_profile" def test_sql_contract_has_tenant_fks_rls_append_only_and_no_legacy_backfill():