From 6bfe16c944390992367fef69188943c8a3cf62f5 Mon Sep 17 00:00:00 2001 From: Ning Zhou Date: Thu, 30 Jul 2026 10:01:46 +0800 Subject: [PATCH 1/3] feat(platform): add active metadata change outbox --- .github/workflows/ci.yml | 11 + data_agent/active_metadata_change_contract.py | 355 ++++++++++++ .../metadata_fabric_active_metadata_outbox.py | 535 ++++++++++++++++++ .../099_active_metadata_change_outbox.sql | 353 ++++++++++++ data_agent/platform_gateway.py | 317 +++++++++++ data_agent/platform_truth.py | 20 + .../test_active_metadata_change_contract.py | 139 +++++ ..._metadata_fabric_active_metadata_outbox.py | 45 ++ ..._fabric_active_metadata_outbox_postgres.py | 111 ++++ data_agent/test_platform_gateway.py | 15 + data_agent/test_platform_truth.py | 5 + ...sactional-active-metadata-change-outbox.md | 62 ++ ...ric-active-metadata-outbox-2026-07-30.json | 33 ++ docs/roadmap-ar0-platform-truth-2026-07-24.md | 5 +- docs/system-of-record-matrix-2026-07-24.md | 13 +- .../metadata-fabric-active-metadata-outbox.sh | 4 + 16 files changed, 2016 insertions(+), 7 deletions(-) create mode 100644 data_agent/active_metadata_change_contract.py create mode 100644 data_agent/metadata_fabric_active_metadata_outbox.py create mode 100644 data_agent/migrations/099_active_metadata_change_outbox.sql create mode 100644 data_agent/test_active_metadata_change_contract.py create mode 100644 data_agent/test_metadata_fabric_active_metadata_outbox.py create mode 100644 data_agent/test_metadata_fabric_active_metadata_outbox_postgres.py create mode 100644 docs/architecture-decisions/adr-060-transactional-active-metadata-change-outbox.md create mode 100644 docs/evidence/metadata-fabric-active-metadata-outbox-2026-07-30.json create mode 100755 scripts/metadata-fabric-active-metadata-outbox.sh diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 39aa0778..1c04c3b7 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -29,6 +29,7 @@ on: - feat/ar1-metadata-fabric-object-store-readiness-gate - feat/ar1-metadata-fabric-spark-commit-failure-recovery - feat/ar1-metadata-fabric-spark-uncertain-commit-reconciliation + - feat/ar1-metadata-fabric-active-metadata-outbox env: PYTHON_VERSION: "3.13" @@ -165,6 +166,9 @@ jobs: - name: Validate metadata fabric Spark uncertain-commit reconciliation evidence run: python -m data_agent.metadata_fabric_spark_uncertain_commit_reconciliation validate + - name: Validate metadata fabric Active Metadata outbox evidence + run: python -m data_agent.metadata_fabric_active_metadata_outbox validate + - name: Validate DolphinScheduler adapter boundary run: python -m data_agent.dolphinscheduler_adapter validate @@ -191,6 +195,11 @@ jobs: DATABASE_URL: postgresql://postgres:postgres@localhost:5432/gis_agent_test run: python -m pytest data_agent/test_metadata_fabric_lineage_delivery_postgres.py -q + - name: Verify metadata fabric Active Metadata outbox on PostgreSQL + env: + DATABASE_URL: postgresql://postgres:postgres@localhost:5432/gis_agent_test + run: python -m pytest data_agent/test_metadata_fabric_active_metadata_outbox_postgres.py -q + - name: Run required platform tests env: DATABASE_URL: postgresql://postgres:postgres@localhost:5432/gis_agent_test @@ -219,6 +228,8 @@ jobs: data_agent/test_metadata_fabric_ingestion.py \ data_agent/test_metadata_fabric_ingestion_replay.py \ data_agent/test_metadata_fabric_binding_contract.py \ + data_agent/test_active_metadata_change_contract.py \ + data_agent/test_metadata_fabric_active_metadata_outbox.py \ data_agent/test_metadata_fabric_lineage_delivery.py \ data_agent/test_metadata_fabric_provider_identity.py \ data_agent/test_metadata_fabric_gravitino_identity.py \ diff --git a/data_agent/active_metadata_change_contract.py b/data_agent/active_metadata_change_contract.py new file mode 100644 index 00000000..2cee2389 --- /dev/null +++ b/data_agent/active_metadata_change_contract.py @@ -0,0 +1,355 @@ +"""Content-bound contracts for the Active Metadata change outbox.""" + +from __future__ import annotations + +import re +from datetime import UTC, datetime +from enum import Enum +from typing import Annotated, Any, Literal, Self +from uuid import UUID, uuid5 + +from pydantic import ( + BaseModel, + ConfigDict, + Field, + StringConstraints, + field_validator, + model_validator, +) + +from .platform_contracts import ( + ResourceVersion, + Sha256, + TenantId, + canonical_json_fingerprint, +) + + +EVENT_SCHEMA = "gda.metadata_change_event.v1" +DELIVERY_SCHEMA = "gda.metadata_change_delivery.v1" +REGISTRATION_SCHEMA = "gda.active_metadata_registration.v1" +ACTIVATION_INTENT_SCHEMA = "gda.metadata_activation_intent.v1" +RESOURCE_VERSION_REGISTERED = "resource_version.registered" +METADATA_PROJECTION_ROUTE = "metadata_fabric.projection_plan" + +ActorSubject = Annotated[ + str, + StringConstraints( + strip_whitespace=True, + min_length=3, + max_length=256, + pattern=r"^(human|workload|agent):[a-zA-Z0-9][a-zA-Z0-9._:@/-]{0,247}$", + ), +] +WorkloadSubject = Annotated[ + str, + StringConstraints( + strip_whitespace=True, + min_length=10, + max_length=256, + pattern=r"^workload:[a-zA-Z0-9][a-zA-Z0-9._:@/-]{0,247}$", + ), +] +ErrorCode = Annotated[ + str, + StringConstraints(pattern=r"^[a-z0-9_]{1,64}$"), +] + + +class ActiveMetadataContractError(ValueError): + """An Active Metadata change or activation intent is not content-bound.""" + + +class _FrozenModel(BaseModel): + model_config = ConfigDict(extra="forbid", frozen=True) + + +class MetadataChangeDeliveryStatus(str, Enum): + PENDING = "pending" + IN_FLIGHT = "in_flight" + PROCESSED = "processed" + FAILED = "failed" + + +def _utc(value: datetime) -> datetime: + if value.tzinfo is None or value.utcoffset() is None: + raise ValueError("timestamp must include a timezone") + return value.astimezone(UTC) + + +def _event_stable(values: dict[str, Any]) -> dict[str, Any]: + return { + "schema": EVENT_SCHEMA, + "event_id": str(values["event_id"]), + "event_type": RESOURCE_VERSION_REGISTERED, + "tenant_id": values["tenant_id"], + "resource_urn": values["resource_urn"], + "resource_version_id": str(values["resource_version_id"]), + "version_key": values["version_key"], + "predecessor_version_id": ( + str(values["predecessor_version_id"]) + if values["predecessor_version_id"] is not None + else None + ), + "content_sha256": values["content_sha256"], + "producer_subject": values["producer_subject"], + "consumer_subject": values["consumer_subject"], + "occurred_at": _utc(values["occurred_at"]) + .isoformat() + .replace("+00:00", "Z"), + } + + +class MetadataChangeEvent(_FrozenModel): + event_schema: Literal["gda.metadata_change_event.v1"] = Field( + default=EVENT_SCHEMA, + alias="schema", + ) + event_id: UUID + event_type: Literal["resource_version.registered"] = ( + RESOURCE_VERSION_REGISTERED + ) + tenant_id: TenantId + resource_urn: str + resource_version_id: UUID + version_key: str + predecessor_version_id: UUID | None = None + content_sha256: Sha256 + producer_subject: ActorSubject + consumer_subject: WorkloadSubject + occurred_at: datetime + event_sha256: Sha256 + + @field_validator("occurred_at") + @classmethod + def _aware_time(cls, value: datetime) -> datetime: + return _utc(value) + + @model_validator(mode="after") + def _content_bound(self) -> Self: + expected_id = uuid5( + self.resource_version_id, + f"active-metadata:{self.event_type}", + ) + if self.event_id != expected_id: + raise ValueError("MetadataChangeEvent ID does not match the version") + expected_sha = canonical_json_fingerprint( + _event_stable( + self.model_dump( + mode="python", + by_alias=True, + exclude={"event_sha256"}, + ) + ) + ) + if self.event_sha256 != expected_sha: + raise ValueError("MetadataChangeEvent SHA-256 does not match") + return self + + +def build_metadata_change_event( + version: ResourceVersion, + *, + consumer_subject: str, +) -> MetadataChangeEvent: + if not re.fullmatch( + r"^(human|workload|agent):[a-zA-Z0-9][a-zA-Z0-9._:@/-]{0,247}$", + version.created_by, + ): + raise ActiveMetadataContractError( + "ResourceVersion creator must use an authenticated subject" + ) + event_id = uuid5( + version.resource_version_id, + f"active-metadata:{RESOURCE_VERSION_REGISTERED}", + ) + values: dict[str, Any] = { + "event_id": event_id, + "tenant_id": version.tenant_id, + "resource_urn": version.resource_urn, + "resource_version_id": version.resource_version_id, + "version_key": version.version_key, + "predecessor_version_id": version.predecessor_version_id, + "content_sha256": version.content_sha256, + "producer_subject": version.created_by, + "consumer_subject": consumer_subject, + "occurred_at": version.created_at, + } + return MetadataChangeEvent( + **values, + event_sha256=canonical_json_fingerprint(_event_stable(values)), + ) + + +class ActiveMetadataRegistration(_FrozenModel): + registration_schema: Literal["gda.active_metadata_registration.v1"] = Field( + default=REGISTRATION_SCHEMA, + alias="schema", + ) + resource_version: ResourceVersion + event: MetadataChangeEvent + + @model_validator(mode="after") + def _same_change(self) -> Self: + expected = build_metadata_change_event( + self.resource_version, + consumer_subject=self.event.consumer_subject, + ) + if self.event != expected: + raise ValueError( + "MetadataChangeEvent does not match its ResourceVersion" + ) + return self + + +def build_active_metadata_registration( + version: ResourceVersion, + *, + consumer_subject: str, +) -> ActiveMetadataRegistration: + return ActiveMetadataRegistration( + resource_version=version, + event=build_metadata_change_event( + version, + consumer_subject=consumer_subject, + ), + ) + + +class MetadataActivationIntent(_FrozenModel): + intent_schema: Literal["gda.metadata_activation_intent.v1"] = Field( + default=ACTIVATION_INTENT_SCHEMA, + alias="schema", + ) + event_id: UUID + event_sha256: Sha256 + tenant_id: TenantId + resource_urn: str + resource_version_id: UUID + content_sha256: Sha256 + route: Literal["metadata_fabric.projection_plan"] = METADATA_PROJECTION_ROUTE + routed_by: WorkloadSubject + provider_apply_authorized: Literal[False] = False + provider_mutations_executed: Literal[False] = False + production_ingestion_verified: Literal[False] = False + intent_sha256: Sha256 + + @model_validator(mode="after") + def _content_bound(self) -> Self: + stable = self.model_dump( + mode="json", + by_alias=True, + exclude={"intent_sha256"}, + ) + if self.intent_sha256 != canonical_json_fingerprint(stable): + raise ValueError("Metadata activation intent SHA-256 does not match") + return self + + +def build_metadata_activation_intent( + event: MetadataChangeEvent, + *, + routed_by: str, +) -> MetadataActivationIntent: + if routed_by != event.consumer_subject: + raise ActiveMetadataContractError( + "activation router must match the event consumer" + ) + values: dict[str, Any] = { + "event_id": event.event_id, + "event_sha256": event.event_sha256, + "tenant_id": event.tenant_id, + "resource_urn": event.resource_urn, + "resource_version_id": event.resource_version_id, + "content_sha256": event.content_sha256, + "routed_by": routed_by, + } + stable = { + "schema": ACTIVATION_INTENT_SCHEMA, + "event_id": str(event.event_id), + "event_sha256": event.event_sha256, + "tenant_id": event.tenant_id, + "resource_urn": event.resource_urn, + "resource_version_id": str(event.resource_version_id), + "content_sha256": event.content_sha256, + "route": METADATA_PROJECTION_ROUTE, + "routed_by": routed_by, + "provider_apply_authorized": False, + "provider_mutations_executed": False, + "production_ingestion_verified": False, + } + return MetadataActivationIntent( + **values, + intent_sha256=canonical_json_fingerprint(stable), + ) + + +class MetadataChangeDelivery(_FrozenModel): + delivery_schema: Literal["gda.metadata_change_delivery.v1"] = Field( + default=DELIVERY_SCHEMA, + alias="schema", + ) + event: MetadataChangeEvent + status: MetadataChangeDeliveryStatus = MetadataChangeDeliveryStatus.PENDING + attempt_count: Annotated[int, Field(ge=0)] = 0 + max_attempts: Annotated[int, Field(ge=1, le=20)] = 5 + available_at: datetime + claimed_by: str | None = None + claimed_until: datetime | None = None + last_error_code: ErrorCode | None = None + activation_intent_sha256: Sha256 | None = None + completed_at: datetime | None = None + + @field_validator("available_at", "claimed_until", "completed_at") + @classmethod + def _aware_delivery_time(cls, value: datetime | None) -> datetime | None: + return _utc(value) if value is not None else None + + @model_validator(mode="after") + def _consistent_state(self) -> Self: + claimed = self.claimed_by is not None and self.claimed_until is not None + if (self.claimed_by is None) != (self.claimed_until is None): + raise ValueError("metadata change claim fields must be set together") + if self.status == MetadataChangeDeliveryStatus.PENDING: + if claimed or self.completed_at or self.activation_intent_sha256: + raise ValueError("pending metadata change has invalid state") + elif self.status == MetadataChangeDeliveryStatus.IN_FLIGHT: + if not claimed or self.completed_at or self.activation_intent_sha256: + raise ValueError("in-flight metadata change has invalid state") + elif self.status == MetadataChangeDeliveryStatus.PROCESSED: + if ( + claimed + or self.completed_at is None + or self.activation_intent_sha256 is None + or self.last_error_code is not None + ): + raise ValueError("processed metadata change has invalid state") + elif ( + claimed + or self.completed_at is None + or self.last_error_code is None + or self.activation_intent_sha256 is not None + ): + raise ValueError("failed metadata change has invalid state") + return self + + +def build_metadata_change_delivery( + event: MetadataChangeEvent, + *, + max_attempts: int = 5, +) -> MetadataChangeDelivery: + return MetadataChangeDelivery( + event=event, + max_attempts=max_attempts, + available_at=event.occurred_at, + ) + + +def metadata_change_binding_payload( + delivery: MetadataChangeDelivery, +) -> dict[str, Any]: + return { + "event": delivery.event.model_dump(mode="json", by_alias=True), + "max_attempts": delivery.max_attempts, + } diff --git a/data_agent/metadata_fabric_active_metadata_outbox.py b/data_agent/metadata_fabric_active_metadata_outbox.py new file mode 100644 index 00000000..e0a8fef2 --- /dev/null +++ b/data_agent/metadata_fabric_active_metadata_outbox.py @@ -0,0 +1,535 @@ +"""Validate and rehearse the local Active Metadata transactional outbox.""" + +from __future__ import annotations + +import argparse +import hashlib +import json +from datetime import datetime, timezone +from pathlib import Path +from typing import Any +from uuid import UUID + +from pydantic import BaseModel, ConfigDict +from sqlalchemy import create_engine, text + +from .active_metadata_change_contract import ( + ActiveMetadataRegistration, + build_active_metadata_registration, + build_metadata_activation_intent, +) +from .platform_contracts import ( + Resource, + ResourceVersion, + canonical_json_fingerprint, +) +from .platform_gateway import ( + GatewayConflictError, + GatewayNotFoundError, + PlatformGateway, +) + + +CONTRACT_SCHEMA = "gda.active_metadata_outbox_contract.v1" +EVIDENCE_SCHEMA = "gda.active_metadata_outbox_evidence.v1" +TENANT = "active-metadata-local" +ISOLATED_TENANT = "active-metadata-isolated" +RESOURCE_URN = f"gda://{TENANT}/dataset/parcels" +RESOURCE_VERSION_ID = UUID("a4000000-0000-4000-8000-000000000001") +LEGACY_VERSION_ID = UUID("a4000000-0000-4000-8000-000000000002") +OCCURRED_AT = datetime(2026, 7, 30, 0, 0, tzinfo=timezone.utc) +PRODUCER_SUBJECT = "workload:metadata-registrar" +CONSUMER_SUBJECT = "workload:metadata-router" +WORKER_1 = "worker:metadata-router-1" +WORKER_2 = "worker:metadata-router-2" +WORKER_3 = "worker:metadata-router-3" + +REPO_ROOT = Path(__file__).resolve().parent.parent +DEFAULT_EVIDENCE_PATH = ( + REPO_ROOT / "docs/evidence/metadata-fabric-active-metadata-outbox-2026-07-30.json" +) +DEFAULT_WRAPPER_PATH = REPO_ROOT / "scripts/metadata-fabric-active-metadata-outbox.sh" +CONTRACT_PATH = Path(__file__).resolve().parent / "active_metadata_change_contract.py" +GATEWAY_PATH = Path(__file__).resolve().parent / "platform_gateway.py" +MIGRATION_PATH = ( + Path(__file__).resolve().parent + / "migrations" + / "099_active_metadata_change_outbox.sql" +) +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", + "099_active_metadata_change_outbox.sql", + ) +) + + +class ActiveMetadataOutboxError(RuntimeError): + """The Active Metadata contract or local rehearsal failed closed.""" + + +class _FrozenModel(BaseModel): + model_config = ConfigDict(extra="forbid", frozen=True) + + +class ActiveMetadataBundle(_FrozenModel): + resource: Resource + registration: ActiveMetadataRegistration + legacy_version: ResourceVersion + + +def _load_json_object(path: Path) -> dict[str, Any]: + value = json.loads(path.read_text(encoding="utf-8")) + if not isinstance(value, dict): + raise ActiveMetadataOutboxError(f"{path.name} must contain an object") + return value + + +def build_active_metadata_bundle() -> ActiveMetadataBundle: + resource = Resource( + tenant_id=TENANT, + resource_urn=RESOURCE_URN, + resource_kind="dataset", + authority_system="iceberg", + authority_locator="geo.parcels", + owner_ref="team:data-platform", + ) + version = ResourceVersion( + tenant_id=TENANT, + resource_urn=RESOURCE_URN, + resource_version_id=RESOURCE_VERSION_ID, + version_key="snapshot-1", + content_sha256="a" * 64, + authority_version_ref={"snapshot_id": 1}, + created_by=PRODUCER_SUBJECT, + created_at=OCCURRED_AT, + ) + legacy_version = ResourceVersion( + tenant_id=TENANT, + resource_urn=RESOURCE_URN, + resource_version_id=LEGACY_VERSION_ID, + version_key="legacy-snapshot", + content_sha256="b" * 64, + authority_version_ref={"snapshot_id": 0}, + created_by=PRODUCER_SUBJECT, + created_at=OCCURRED_AT, + ) + return ActiveMetadataBundle( + resource=resource, + registration=build_active_metadata_registration( + version, + consumer_subject=CONSUMER_SUBJECT, + ), + legacy_version=legacy_version, + ) + + +def build_contract_report() -> dict[str, Any]: + errors: list[str] = [] + paths = { + "contract": CONTRACT_PATH, + "migration": MIGRATION_PATH, + "gateway": GATEWAY_PATH, + "wrapper": DEFAULT_WRAPPER_PATH, + } + required = { + "contract": ( + "class MetadataChangeEvent", + "class MetadataActivationIntent", + "provider_apply_authorized: Literal[False]", + "provider_mutations_executed: Literal[False]", + "production_ingestion_verified: Literal[False]", + ), + "migration": ( + "CREATE TABLE IF NOT EXISTS gda_control.metadata_change_outbox", + "FOR UPDATE SKIP LOCKED", + "FORCE ROW LEVEL SECURITY", + "claim_metadata_changes", + "complete_metadata_change", + "fail_metadata_change", + "GRANT SELECT, INSERT ON gda_control.metadata_change_outbox", + ), + "gateway": ( + "def register_resource_version_with_metadata_event(", + "version_result.created != event_result.created", + "def claim_metadata_changes(", + "def complete_metadata_change(", + "activation_intent != expected", + ), + "wrapper": ( + "data_agent.metadata_fabric_active_metadata_outbox", + '"$@"', + ), + } + files: dict[str, dict[str, Any]] = {} + for name, path in paths.items(): + if not path.is_file(): + errors.append(f"{name} is missing") + files[name] = {"path": path.resolve().as_posix(), "sha256": None} + continue + raw = path.read_bytes() + source = raw.decode("utf-8") + files[name] = { + "path": path.resolve().as_posix(), + "sha256": hashlib.sha256(raw).hexdigest(), + } + missing = [marker for marker in required[name] if marker not in source] + if missing: + errors.append(f"{name} is missing required Active Metadata markers") + + bundle = build_active_metadata_bundle() + stable = { + "schema": CONTRACT_SCHEMA, + "event_id": str(bundle.registration.event.event_id), + "event_sha256": bundle.registration.event.event_sha256, + "consumer_subject": bundle.registration.event.consumer_subject, + "activation_route": "metadata_fabric.projection_plan", + "files": files, + "errors": errors, + } + return { + **stable, + "status": "valid" if not errors else "invalid", + "contract_sha256": canonical_json_fingerprint(stable), + "local_postgresql_active_metadata_loop_verified": False, + "production_ready": False, + } + + +def _apply_migrations(engine) -> None: + 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: + raise ActiveMetadataOutboxError( + "local Active Metadata rehearsal requires a fresh superuser database" + ) + 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"))) + + +def run_local_rehearsal(database_url: str) -> dict[str, Any]: + if not database_url.startswith( + ("postgresql://", "postgresql+psycopg://", "postgresql+psycopg2://") + ): + raise ActiveMetadataOutboxError( + "Active Metadata rehearsal requires a PostgreSQL database URL" + ) + bundle = build_active_metadata_bundle() + engine = create_engine(database_url) + try: + _apply_migrations(engine) + gateway = PlatformGateway(engine) + gateway.register_resource(bundle.resource) + first = gateway.register_resource_version_with_metadata_event( + bundle.registration, + max_attempts=3, + ) + replay = gateway.register_resource_version_with_metadata_event( + bundle.registration, + max_attempts=3, + ) + wrong_consumer_claim = gateway.claim_metadata_changes( + TENANT, + WORKER_1, + consumer_subject="workload:other-router", + ) + first_claim = gateway.claim_metadata_changes( + TENANT, + WORKER_1, + consumer_subject=CONSUMER_SUBJECT, + lease_seconds=60, + ) + intent = build_metadata_activation_intent( + first_claim[0].event, + routed_by=CONSUMER_SUBJECT, + ) + wrong_worker_blocked = False + try: + gateway.complete_metadata_change( + TENANT, + bundle.registration.event.event_id, + worker_id=WORKER_2, + activation_intent=intent, + ) + except GatewayConflictError: + wrong_worker_blocked = True + after_retry = gateway.fail_metadata_change( + TENANT, + bundle.registration.event.event_id, + worker_id=WORKER_1, + error_code="router_unavailable", + retryable=True, + retry_delay_seconds=0, + ) + second_claim = gateway.claim_metadata_changes( + TENANT, + WORKER_2, + consumer_subject=CONSUMER_SUBJECT, + lease_seconds=60, + ) + with engine.begin() as connection: + connection.execute( + text( + """ + UPDATE gda_control.metadata_change_outbox + SET claimed_until = clock_timestamp() - interval '1 second' + WHERE tenant_id = :tenant_id AND event_id = :event_id + """ + ), + { + "tenant_id": TENANT, + "event_id": bundle.registration.event.event_id, + }, + ) + third_claim = gateway.claim_metadata_changes( + TENANT, + WORKER_3, + consumer_subject=CONSUMER_SUBJECT, + lease_seconds=60, + ) + completed = gateway.complete_metadata_change( + TENANT, + bundle.registration.event.event_id, + worker_id=WORKER_3, + activation_intent=intent, + ) + final_claim = gateway.claim_metadata_changes( + TENANT, + WORKER_3, + consumer_subject=CONSUMER_SUBJECT, + ) + final_replay = gateway.register_resource_version_with_metadata_event( + bundle.registration, + max_attempts=3, + ) + + cross_tenant_blocked = False + try: + gateway.get_metadata_change_delivery( + ISOLATED_TENANT, + bundle.registration.event.event_id, + ) + except GatewayNotFoundError: + cross_tenant_blocked = True + + gateway.register_resource_version(bundle.legacy_version) + legacy_registration = build_active_metadata_registration( + bundle.legacy_version, + consumer_subject=CONSUMER_SUBJECT, + ) + legacy_backfill_blocked = False + try: + gateway.register_resource_version_with_metadata_event( + legacy_registration + ) + except GatewayConflictError: + legacy_backfill_blocked = True + legacy_event_rolled_back = False + try: + gateway.get_metadata_change_delivery( + TENANT, + legacy_registration.event.event_id, + ) + except GatewayNotFoundError: + legacy_event_rolled_back = True + + with engine.connect() as connection: + privileges = connection.exec_driver_sql( + """ + SELECT + has_table_privilege( + 'gda_control_gateway', + 'gda_control.metadata_change_outbox', 'SELECT,INSERT' + ), + NOT has_table_privilege( + 'gda_control_gateway', + 'gda_control.metadata_change_outbox', 'UPDATE' + ), + NOT has_table_privilege( + 'gda_control_gateway', + 'gda_control.metadata_change_outbox', 'DELETE' + ), + has_function_privilege( + 'gda_control_gateway', + 'gda_control.claim_metadata_changes(text,text,text,integer,integer)', + 'EXECUTE' + ) + """ + ).one() + force_rls = connection.exec_driver_sql( + """ + SELECT relforcerowsecurity + FROM pg_class + WHERE oid = 'gda_control.metadata_change_outbox'::regclass + """ + ).scalar_one() + event_count = connection.execute( + text( + """ + SELECT count(*) + FROM gda_control.metadata_change_outbox + WHERE tenant_id = :tenant_id + """ + ), + {"tenant_id": TENANT}, + ).scalar_one() + connection.rollback() + finally: + engine.dispose() + + verified = ( + first.created + and not replay.created + and not wrong_consumer_claim + and len(first_claim) == len(second_claim) == len(third_claim) == 1 + and first_claim[0].attempt_count == 1 + and after_retry.status.value == "pending" + and after_retry.attempt_count == 1 + and second_claim[0].attempt_count == 2 + and third_claim[0].attempt_count == 3 + and wrong_worker_blocked + and completed.status.value == "processed" + and completed.activation_intent_sha256 == intent.intent_sha256 + and not final_claim + and not final_replay.created + and cross_tenant_blocked + and legacy_backfill_blocked + and legacy_event_rolled_back + and privileges == (True, True, True, True) + and force_rls + and event_count == 1 + ) + contract = build_contract_report() + stable = { + "schema": EVIDENCE_SCHEMA, + "status": ( + "local_postgresql_active_metadata_loop_verified" + if verified + else "blocked" + ), + "contract_sha256": contract["contract_sha256"], + "event_id": str(bundle.registration.event.event_id), + "event_sha256": bundle.registration.event.event_sha256, + "activation_intent_sha256": intent.intent_sha256, + "activation_route": intent.route, + "first_registration_created": first.created, + "exact_replay_created": replay.created, + "processed_replay_created": final_replay.created, + "wrong_consumer_claim_blocked": not wrong_consumer_claim, + "wrong_worker_completion_blocked": wrong_worker_blocked, + "retry_pending_verified": after_retry.status.value == "pending", + "lease_expiry_reclaim_verified": ( + len(third_claim) == 1 and third_claim[0].attempt_count == 3 + ), + "final_attempt_count": completed.attempt_count, + "processed_delivery_not_reclaimed": not final_claim, + "legacy_backfill_blocked": legacy_backfill_blocked, + "legacy_event_transaction_rolled_back": legacy_event_rolled_back, + "cross_tenant_read_blocked": cross_tenant_blocked, + "gateway_select_insert_only_verified": privileges == (True, True, True, True), + "force_rls_verified": bool(force_rls), + "authoritative_event_count": event_count, + "local_postgresql_active_metadata_loop_verified": verified, + "transactional_outbox_verified": verified, + "provider_apply_authorized": False, + "provider_mutations_executed": False, + "production_ingestion_verified": False, + "production_scheduler_submission_verified": False, + "production_ready": False, + "errors": [] if verified else ["local Active Metadata loop did not verify"], + } + return {**stable, "evidence_sha256": canonical_json_fingerprint(stable)} + + +def validate_rehearsal_evidence(evidence: dict[str, Any]) -> list[str]: + errors: list[str] = [] + stable = {key: value for key, value in evidence.items() if key != "evidence_sha256"} + if evidence.get("schema") != EVIDENCE_SCHEMA: + errors.append("Active Metadata evidence schema does not match") + if evidence.get("evidence_sha256") != canonical_json_fingerprint(stable): + errors.append("Active Metadata evidence SHA-256 does not match") + contract = build_contract_report() + if evidence.get("contract_sha256") != contract.get("contract_sha256"): + errors.append("Active Metadata evidence contract fingerprint is stale") + for claim in ( + "provider_apply_authorized", + "provider_mutations_executed", + "production_ingestion_verified", + "production_scheduler_submission_verified", + "production_ready", + ): + if evidence.get(claim) is not False: + errors.append(f"local Active Metadata evidence may not claim {claim}") + for claim in ( + "wrong_consumer_claim_blocked", + "wrong_worker_completion_blocked", + "retry_pending_verified", + "lease_expiry_reclaim_verified", + "processed_delivery_not_reclaimed", + "legacy_backfill_blocked", + "legacy_event_transaction_rolled_back", + "cross_tenant_read_blocked", + "gateway_select_insert_only_verified", + "force_rls_verified", + "local_postgresql_active_metadata_loop_verified", + "transactional_outbox_verified", + ): + if evidence.get(claim) is not True: + errors.append(f"local Active Metadata evidence did not verify {claim}") + if evidence.get("activation_route") != "metadata_fabric.projection_plan": + errors.append("Active Metadata activation route is invalid") + if evidence.get("authoritative_event_count") != 1: + errors.append("Active Metadata evidence must contain exactly one event") + return errors + + +def main(argv: list[str] | None = None) -> int: + parser = argparse.ArgumentParser(description=__doc__) + subparsers = parser.add_subparsers(dest="command", required=True) + validate = subparsers.add_parser("validate") + validate.add_argument("--evidence", type=Path, default=DEFAULT_EVIDENCE_PATH) + rehearse = subparsers.add_parser("rehearse") + rehearse.add_argument("--database-url", required=True) + rehearse.add_argument("--evidence-out", type=Path, required=True) + args = parser.parse_args(argv) + + if args.command == "validate": + report = build_contract_report() + try: + evidence = _load_json_object(args.evidence) + report["errors"].extend(validate_rehearsal_evidence(evidence)) + except (OSError, ValueError) as exc: + report["errors"].append( + f"Active Metadata evidence is invalid: {type(exc).__name__}" + ) + report["status"] = "valid" if not report["errors"] else "invalid" + report["local_postgresql_active_metadata_loop_verified"] = not report[ + "errors" + ] + print(json.dumps(report, ensure_ascii=True, indent=2, sort_keys=True)) + return 0 if not report["errors"] else 1 + + evidence = run_local_rehearsal(args.database_url) + args.evidence_out.write_text( + json.dumps(evidence, ensure_ascii=True, indent=2, sort_keys=True) + "\n", + encoding="utf-8", + ) + print(json.dumps(evidence, ensure_ascii=True, indent=2, sort_keys=True)) + return 0 if not evidence["errors"] else 1 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/data_agent/migrations/099_active_metadata_change_outbox.sql b/data_agent/migrations/099_active_metadata_change_outbox.sql new file mode 100644 index 00000000..266cbf7e --- /dev/null +++ b/data_agent/migrations/099_active_metadata_change_outbox.sql @@ -0,0 +1,353 @@ +-- 099: Transactional outbox for authoritative ResourceVersion changes. +-- +-- The ResourceVersion remains authoritative. This table owns only delivery of +-- its content-bound MetadataChangeEvent to the Active Metadata router. + +CREATE TABLE IF NOT EXISTS gda_control.metadata_change_outbox ( + tenant_id TEXT NOT NULL, + event_id UUID PRIMARY KEY, + event_type TEXT NOT NULL, + resource_urn TEXT NOT NULL, + resource_version_id UUID NOT NULL, + version_key TEXT NOT NULL, + predecessor_version_id UUID, + content_sha256 CHAR(64) NOT NULL, + producer_subject TEXT NOT NULL, + consumer_subject TEXT NOT NULL, + occurred_at TIMESTAMPTZ NOT NULL, + event JSONB NOT NULL, + event_sha256 CHAR(64) NOT NULL, + status TEXT NOT NULL DEFAULT 'pending', + attempt_count INTEGER NOT NULL DEFAULT 0, + max_attempts INTEGER NOT NULL DEFAULT 5, + available_at TIMESTAMPTZ NOT NULL, + claimed_by TEXT, + claimed_until TIMESTAMPTZ, + last_error_code TEXT, + activation_intent_sha256 CHAR(64), + created_at TIMESTAMPTZ NOT NULL DEFAULT now(), + completed_at TIMESTAMPTZ, + CONSTRAINT uq_gda_metadata_change_tenant_event + UNIQUE (tenant_id, event_id), + CONSTRAINT uq_gda_metadata_change_version_type + UNIQUE (tenant_id, resource_version_id, event_type), + CONSTRAINT uq_gda_metadata_change_sha + UNIQUE (tenant_id, event_sha256), + CONSTRAINT fk_gda_metadata_change_version + FOREIGN KEY ( + tenant_id, resource_urn, resource_version_id, content_sha256 + ) REFERENCES gda_control.resource_version( + tenant_id, resource_urn, resource_version_id, content_sha256 + ), + CONSTRAINT ck_gda_metadata_change_type CHECK ( + event_type = 'resource_version.registered' + ), + CONSTRAINT ck_gda_metadata_change_subjects CHECK ( + producer_subject ~ '^(human|workload|agent):.+' + AND consumer_subject ~ '^workload:.+' + ), + CONSTRAINT ck_gda_metadata_change_event_document CHECK ( + jsonb_typeof(event) = 'object' + AND event ?& ARRAY[ + 'schema', 'event_id', 'event_type', 'tenant_id', + 'resource_urn', 'resource_version_id', 'version_key', + 'predecessor_version_id', 'content_sha256', + 'producer_subject', 'consumer_subject', 'occurred_at', + 'event_sha256' + ] + AND event - ARRAY[ + 'schema', 'event_id', 'event_type', 'tenant_id', + 'resource_urn', 'resource_version_id', 'version_key', + 'predecessor_version_id', 'content_sha256', + 'producer_subject', 'consumer_subject', 'occurred_at', + 'event_sha256' + ] = '{}'::jsonb + AND event->>'schema' = 'gda.metadata_change_event.v1' + AND event->>'event_id' = event_id::text + AND event->>'event_type' = event_type + AND event->>'tenant_id' = tenant_id + AND event->>'resource_urn' = resource_urn + AND event->>'resource_version_id' = resource_version_id::text + AND event->>'version_key' = version_key + AND event->>'predecessor_version_id' + IS NOT DISTINCT FROM predecessor_version_id::text + AND event->>'content_sha256' = content_sha256 + AND event->>'producer_subject' = producer_subject + AND event->>'consumer_subject' = consumer_subject + AND (event->>'occurred_at')::timestamptz = occurred_at + AND event->>'event_sha256' = event_sha256 + ), + CONSTRAINT ck_gda_metadata_change_sha256 CHECK ( + content_sha256 ~ '^[0-9a-f]{64}$' + AND event_sha256 ~ '^[0-9a-f]{64}$' + AND ( + activation_intent_sha256 IS NULL + OR activation_intent_sha256 ~ '^[0-9a-f]{64}$' + ) + ), + CONSTRAINT ck_gda_metadata_change_status CHECK ( + status IN ('pending', 'in_flight', 'processed', 'failed') + ), + CONSTRAINT ck_gda_metadata_change_attempts CHECK ( + attempt_count >= 0 AND max_attempts BETWEEN 1 AND 20 + ), + CONSTRAINT ck_gda_metadata_change_claim CHECK ( + (claimed_by IS NULL) = (claimed_until IS NULL) + ), + CONSTRAINT ck_gda_metadata_change_error CHECK ( + last_error_code IS NULL + OR last_error_code ~ '^[a-z0-9_]{1,64}$' + ), + CONSTRAINT ck_gda_metadata_change_state CHECK ( + ( + status = 'pending' + AND claimed_by IS NULL + AND completed_at IS NULL + AND activation_intent_sha256 IS NULL + ) + OR ( + status = 'in_flight' + AND claimed_by IS NOT NULL + AND completed_at IS NULL + AND activation_intent_sha256 IS NULL + ) + OR ( + status = 'processed' + AND claimed_by IS NULL + AND completed_at IS NOT NULL + AND last_error_code IS NULL + AND activation_intent_sha256 IS NOT NULL + ) + OR ( + status = 'failed' + AND claimed_by IS NULL + AND completed_at IS NOT NULL + AND last_error_code IS NOT NULL + AND activation_intent_sha256 IS NULL + ) + ) +); + +CREATE INDEX IF NOT EXISTS idx_gda_metadata_change_due + ON gda_control.metadata_change_outbox( + tenant_id, consumer_subject, available_at, occurred_at + ) WHERE status = 'pending'; +CREATE INDEX IF NOT EXISTS idx_gda_metadata_change_expired_claim + ON gda_control.metadata_change_outbox(tenant_id, claimed_until) + WHERE status = 'in_flight'; +CREATE INDEX IF NOT EXISTS idx_gda_metadata_change_resource + ON gda_control.metadata_change_outbox( + tenant_id, resource_urn, occurred_at DESC + ); + +ALTER TABLE gda_control.metadata_change_outbox ENABLE ROW LEVEL SECURITY; +ALTER TABLE gda_control.metadata_change_outbox FORCE ROW LEVEL SECURITY; +DROP POLICY IF EXISTS gda_metadata_change_tenant_isolation + ON gda_control.metadata_change_outbox; +CREATE POLICY gda_metadata_change_tenant_isolation + ON gda_control.metadata_change_outbox + USING (tenant_id = gda_control.current_tenant()) + WITH CHECK (tenant_id = gda_control.current_tenant()); + +CREATE OR REPLACE FUNCTION gda_control.claim_metadata_changes( + p_tenant_id TEXT, + p_consumer_subject TEXT, + p_worker_id TEXT, + p_limit INTEGER DEFAULT 10, + p_lease_seconds INTEGER DEFAULT 60 +) +RETURNS SETOF gda_control.metadata_change_outbox +LANGUAGE plpgsql +SECURITY DEFINER +SET search_path = pg_catalog, gda_control +SET row_security = on +AS $$ +BEGIN + IF gda_control.current_tenant() IS DISTINCT FROM p_tenant_id THEN + RAISE EXCEPTION 'tenant context mismatch' USING ERRCODE = '42501'; + END IF; + IF COALESCE(btrim(p_consumer_subject), '') = '' + OR p_consumer_subject !~ '^workload:.+' THEN + RAISE EXCEPTION 'metadata change consumer is invalid' + USING ERRCODE = '22023'; + END IF; + IF COALESCE(btrim(p_worker_id), '') = '' THEN + RAISE EXCEPTION 'worker identity is required' USING ERRCODE = '22023'; + END IF; + IF p_limit IS NULL OR p_limit < 1 OR p_limit > 100 THEN + RAISE EXCEPTION 'claim limit must be between 1 and 100' + USING ERRCODE = '22023'; + END IF; + IF p_lease_seconds IS NULL + OR p_lease_seconds < 5 OR p_lease_seconds > 3600 THEN + RAISE EXCEPTION 'lease must be between 5 and 3600 seconds' + USING ERRCODE = '22023'; + END IF; + + UPDATE gda_control.metadata_change_outbox + SET status = 'failed', + claimed_by = NULL, + claimed_until = NULL, + last_error_code = COALESCE(last_error_code, 'lease_expired'), + completed_at = clock_timestamp() + WHERE tenant_id = p_tenant_id + AND status = 'in_flight' + AND claimed_until <= clock_timestamp() + AND attempt_count >= max_attempts; + + RETURN QUERY + WITH candidates AS ( + SELECT event_id + FROM gda_control.metadata_change_outbox + WHERE tenant_id = p_tenant_id + AND consumer_subject = p_consumer_subject + AND attempt_count < max_attempts + AND ( + (status = 'pending' AND available_at <= clock_timestamp()) + OR + (status = 'in_flight' AND claimed_until <= clock_timestamp()) + ) + ORDER BY available_at, occurred_at, event_id + LIMIT p_limit + FOR UPDATE SKIP LOCKED + ) + UPDATE gda_control.metadata_change_outbox AS delivery + SET status = 'in_flight', + attempt_count = delivery.attempt_count + 1, + claimed_by = p_worker_id, + claimed_until = clock_timestamp() + + make_interval(secs => p_lease_seconds), + last_error_code = NULL, + completed_at = NULL + FROM candidates + WHERE delivery.tenant_id = p_tenant_id + AND delivery.event_id = candidates.event_id + RETURNING delivery.*; +END; +$$; + +CREATE OR REPLACE FUNCTION gda_control.complete_metadata_change( + p_tenant_id TEXT, + p_event_id UUID, + p_worker_id TEXT, + p_activation_intent_sha256 TEXT +) +RETURNS SETOF gda_control.metadata_change_outbox +LANGUAGE plpgsql +SECURITY DEFINER +SET search_path = pg_catalog, gda_control +SET row_security = on +AS $$ +BEGIN + IF gda_control.current_tenant() IS DISTINCT FROM p_tenant_id THEN + RAISE EXCEPTION 'tenant context mismatch' USING ERRCODE = '42501'; + END IF; + IF p_activation_intent_sha256 !~ '^[0-9a-f]{64}$' THEN + RAISE EXCEPTION 'activation intent fingerprint is invalid' + USING ERRCODE = '22023'; + END IF; + RETURN QUERY + UPDATE gda_control.metadata_change_outbox AS delivery + SET status = 'processed', + claimed_by = NULL, + claimed_until = NULL, + last_error_code = NULL, + activation_intent_sha256 = p_activation_intent_sha256, + completed_at = clock_timestamp() + WHERE delivery.tenant_id = p_tenant_id + AND delivery.event_id = p_event_id + AND delivery.status = 'in_flight' + AND delivery.claimed_by = p_worker_id + AND delivery.claimed_until > clock_timestamp() + RETURNING delivery.*; + IF NOT FOUND THEN + RAISE EXCEPTION 'metadata change claim is missing, expired, or owned by another worker' + USING ERRCODE = '40001'; + END IF; +END; +$$; + +CREATE OR REPLACE FUNCTION gda_control.fail_metadata_change( + p_tenant_id TEXT, + p_event_id UUID, + p_worker_id TEXT, + p_error_code TEXT, + p_retryable BOOLEAN DEFAULT true, + p_retry_delay_seconds INTEGER DEFAULT 30 +) +RETURNS SETOF gda_control.metadata_change_outbox +LANGUAGE plpgsql +SECURITY DEFINER +SET search_path = pg_catalog, gda_control +SET row_security = on +AS $$ +BEGIN + IF gda_control.current_tenant() IS DISTINCT FROM p_tenant_id THEN + RAISE EXCEPTION 'tenant context mismatch' USING ERRCODE = '42501'; + END IF; + IF COALESCE(p_error_code, '') !~ '^[a-z0-9_]{1,64}$' THEN + RAISE EXCEPTION 'failure code is invalid' USING ERRCODE = '22023'; + END IF; + IF p_retry_delay_seconds IS NULL + OR p_retry_delay_seconds < 0 OR p_retry_delay_seconds > 86400 THEN + RAISE EXCEPTION 'retry delay must be between 0 and 86400 seconds' + USING ERRCODE = '22023'; + END IF; + RETURN QUERY + UPDATE gda_control.metadata_change_outbox AS delivery + SET status = CASE + WHEN NOT p_retryable + OR delivery.attempt_count >= delivery.max_attempts + THEN 'failed' ELSE 'pending' END, + claimed_by = NULL, + claimed_until = NULL, + last_error_code = p_error_code, + activation_intent_sha256 = NULL, + available_at = CASE + WHEN NOT p_retryable + OR delivery.attempt_count >= delivery.max_attempts + THEN delivery.available_at + ELSE clock_timestamp() + + make_interval(secs => p_retry_delay_seconds) + END, + completed_at = CASE + WHEN NOT p_retryable + OR delivery.attempt_count >= delivery.max_attempts + THEN clock_timestamp() ELSE NULL END + WHERE delivery.tenant_id = p_tenant_id + AND delivery.event_id = p_event_id + AND delivery.status = 'in_flight' + AND delivery.claimed_by = p_worker_id + AND delivery.claimed_until > clock_timestamp() + RETURNING delivery.*; + IF NOT FOUND THEN + RAISE EXCEPTION 'metadata change claim is missing, expired, or owned by another worker' + USING ERRCODE = '40001'; + END IF; +END; +$$; + +REVOKE ALL ON TABLE gda_control.metadata_change_outbox FROM PUBLIC; +REVOKE ALL ON TABLE gda_control.metadata_change_outbox + FROM gda_control_gateway; +GRANT SELECT, INSERT ON gda_control.metadata_change_outbox + TO gda_control_gateway; + +REVOKE ALL ON FUNCTION gda_control.claim_metadata_changes( + text, text, text, integer, integer +) FROM PUBLIC; +REVOKE ALL ON FUNCTION gda_control.complete_metadata_change( + text, uuid, text, text +) FROM PUBLIC; +REVOKE ALL ON FUNCTION gda_control.fail_metadata_change( + text, uuid, text, text, boolean, integer +) FROM PUBLIC; +GRANT EXECUTE ON FUNCTION gda_control.claim_metadata_changes( + text, text, text, integer, integer +) TO gda_control_gateway; +GRANT EXECUTE ON FUNCTION gda_control.complete_metadata_change( + text, uuid, text, text +) TO gda_control_gateway; +GRANT EXECUTE ON FUNCTION gda_control.fail_metadata_change( + text, uuid, text, text, boolean, integer +) TO gda_control_gateway; diff --git a/data_agent/platform_gateway.py b/data_agent/platform_gateway.py index 0f110d59..814856b7 100644 --- a/data_agent/platform_gateway.py +++ b/data_agent/platform_gateway.py @@ -17,6 +17,16 @@ from sqlalchemy import text from sqlalchemy.exc import DBAPIError, SQLAlchemyError +from .active_metadata_change_contract import ( + ActiveMetadataRegistration, + MetadataActivationIntent, + MetadataChangeDelivery, + MetadataChangeDeliveryStatus, + MetadataChangeEvent, + build_metadata_activation_intent, + build_metadata_change_delivery, + metadata_change_binding_payload, +) from .db_engine import get_engine from .metadata_fabric_binding_contract import ( MetadataFabricApplyPlan, @@ -89,6 +99,11 @@ / "migrations" / "098_metadata_fabric_openlineage_delivery.sql" ) +ACTIVE_METADATA_CHANGE_MIGRATION = ( + Path(__file__).resolve().parent + / "migrations" + / "099_active_metadata_change_outbox.sql" +) USER_TENANT_MIGRATION = ( Path(__file__).resolve().parent / "migrations" @@ -350,6 +365,288 @@ def register_resource_version( with self._transaction(version.tenant_id) as connection: return self._put_resource_version(connection, version) + @staticmethod + def _metadata_change_from_row(row) -> MetadataChangeDelivery: + fields = { + "event", + "status", + "attempt_count", + "max_attempts", + "available_at", + "claimed_by", + "claimed_until", + "last_error_code", + "activation_intent_sha256", + "completed_at", + } + value = {name: row[name] for name in fields} + value["event"] = _as_json(value["event"]) + return MetadataChangeDelivery.model_validate(value) + + @classmethod + def _load_metadata_change_delivery( + cls, connection, tenant_id: str, event_id: UUID + ) -> MetadataChangeDelivery | None: + row = connection.execute( + text( + """ + SELECT event, status, attempt_count, max_attempts, + available_at, claimed_by, claimed_until, + last_error_code, activation_intent_sha256, completed_at + FROM gda_control.metadata_change_outbox + WHERE tenant_id = :tenant_id AND event_id = :event_id + """ + ), + {"tenant_id": tenant_id, "event_id": event_id}, + ).mappings().one_or_none() + if row is None: + return None + return cls._metadata_change_from_row(row) + + def _put_metadata_change_event( + self, + connection, + event: MetadataChangeEvent, + *, + max_attempts: int, + ) -> GatewayWriteResult: + delivery = build_metadata_change_delivery( + event, + max_attempts=max_attempts, + ) + inserted = connection.execute( + text( + """ + INSERT INTO gda_control.metadata_change_outbox ( + tenant_id, event_id, event_type, resource_urn, + resource_version_id, version_key, + predecessor_version_id, content_sha256, + producer_subject, consumer_subject, occurred_at, + event, event_sha256, status, attempt_count, + max_attempts, available_at, claimed_by, claimed_until, + last_error_code, activation_intent_sha256, completed_at + ) VALUES ( + :tenant_id, :event_id, :event_type, :resource_urn, + :resource_version_id, :version_key, + :predecessor_version_id, :content_sha256, + :producer_subject, :consumer_subject, :occurred_at, + CAST(:event AS jsonb), :event_sha256, :status, + :attempt_count, :max_attempts, :available_at, + :claimed_by, :claimed_until, :last_error_code, + :activation_intent_sha256, :completed_at + ) + ON CONFLICT DO NOTHING + RETURNING event_id + """ + ), + { + **event.model_dump( + mode="python", + by_alias=False, + exclude={"event_schema"}, + ), + "event": _json(event.model_dump(mode="json", by_alias=True)), + "status": delivery.status.value, + "attempt_count": delivery.attempt_count, + "max_attempts": delivery.max_attempts, + "available_at": delivery.available_at, + "claimed_by": delivery.claimed_by, + "claimed_until": delivery.claimed_until, + "last_error_code": delivery.last_error_code, + "activation_intent_sha256": delivery.activation_intent_sha256, + "completed_at": delivery.completed_at, + }, + ).first() + stored = self._load_metadata_change_delivery( + connection, + event.tenant_id, + event.event_id, + ) + if ( + stored is None + or metadata_change_binding_payload(stored) + != metadata_change_binding_payload(delivery) + ): + raise GatewayConflictError( + "MetadataChangeEvent identity already has different content" + ) + return GatewayWriteResult(stored, inserted is not None) + + def register_resource_version_with_metadata_event( + self, + registration: ActiveMetadataRegistration, + *, + max_attempts: int = 5, + ) -> GatewayWriteResult: + try: + registration = ActiveMetadataRegistration.model_validate( + registration.model_dump(mode="json", by_alias=True) + ) + build_metadata_change_delivery( + registration.event, + max_attempts=max_attempts, + ) + except ValueError as exc: + raise GatewayValidationError( + "Active Metadata registration is not content-bound" + ) from exc + with self._transaction(registration.resource_version.tenant_id) as connection: + version_result = self._put_resource_version( + connection, + registration.resource_version, + ) + event_result = self._put_metadata_change_event( + connection, + registration.event, + max_attempts=max_attempts, + ) + if version_result.created != event_result.created: + raise GatewayConflictError( + "ResourceVersion and MetadataChangeEvent creation state diverged" + ) + stored = ActiveMetadataRegistration( + resource_version=version_result.value, + event=event_result.value.event, + ) + return GatewayWriteResult(stored, version_result.created) + + def get_metadata_change_delivery( + self, + tenant_id: str, + event_id: UUID, + ) -> MetadataChangeDelivery: + tenant = _TENANT_ADAPTER.validate_python(tenant_id) + with self._transaction(tenant) as connection: + delivery = self._load_metadata_change_delivery( + connection, + tenant, + event_id, + ) + if delivery is None: + raise GatewayNotFoundError( + "MetadataChangeEvent delivery was not found" + ) + return delivery + + def claim_metadata_changes( + self, + tenant_id: str, + worker_id: str, + *, + consumer_subject: str, + limit: int = 10, + lease_seconds: int = 60, + ) -> list[MetadataChangeDelivery]: + with self._transaction(tenant_id) as connection: + rows = connection.execute( + text( + """ + SELECT * FROM gda_control.claim_metadata_changes( + :tenant_id, :consumer_subject, :worker_id, + :limit, :lease_seconds + ) + """ + ), + { + "tenant_id": tenant_id, + "consumer_subject": consumer_subject, + "worker_id": worker_id, + "limit": limit, + "lease_seconds": lease_seconds, + }, + ).mappings().all() + return [self._metadata_change_from_row(row) for row in rows] + + def complete_metadata_change( + self, + tenant_id: str, + event_id: UUID, + *, + worker_id: str, + activation_intent: MetadataActivationIntent, + ) -> MetadataChangeDelivery: + try: + activation_intent = MetadataActivationIntent.model_validate( + activation_intent.model_dump(mode="json", by_alias=True) + ) + except ValueError as exc: + raise GatewayValidationError( + "metadata activation intent is not content-bound" + ) from exc + with self._transaction(tenant_id) as connection: + claimed = self._load_metadata_change_delivery( + connection, + tenant_id, + event_id, + ) + if claimed is None: + raise GatewayNotFoundError( + "MetadataChangeEvent delivery was not found" + ) + expected = build_metadata_activation_intent( + claimed.event, + routed_by=claimed.event.consumer_subject, + ) + if ( + claimed.status != MetadataChangeDeliveryStatus.IN_FLIGHT + or activation_intent != expected + ): + raise GatewayValidationError( + "activation intent does not match the claimed metadata change" + ) + row = connection.execute( + text( + """ + SELECT * FROM gda_control.complete_metadata_change( + :tenant_id, :event_id, :worker_id, + :activation_intent_sha256 + ) + """ + ), + { + "tenant_id": tenant_id, + "event_id": event_id, + "worker_id": worker_id, + "activation_intent_sha256": activation_intent.intent_sha256, + }, + ).mappings().one() + return self._metadata_change_from_row(row) + + def fail_metadata_change( + self, + tenant_id: str, + event_id: UUID, + *, + worker_id: str, + error_code: str, + retryable: bool = True, + retry_delay_seconds: int = 30, + ) -> MetadataChangeDelivery: + if not re.fullmatch(r"[a-z0-9_]{1,64}", error_code): + raise GatewayValidationError( + "metadata change failure code is invalid" + ) + with self._transaction(tenant_id) as connection: + row = connection.execute( + text( + """ + SELECT * FROM gda_control.fail_metadata_change( + :tenant_id, :event_id, :worker_id, :error_code, + :retryable, :retry_delay_seconds + ) + """ + ), + { + "tenant_id": tenant_id, + "event_id": event_id, + "worker_id": worker_id, + "error_code": error_code, + "retryable": retryable, + "retry_delay_seconds": retry_delay_seconds, + }, + ).mappings().one() + return self._metadata_change_from_row(row) + @staticmethod def _load_definition( connection, tenant_id: str, definition_version_id: UUID @@ -1889,6 +2186,7 @@ def build_gateway_report( success_migration: Path | None = None, binding_migration: Path | None = None, lineage_migration: Path | None = None, + active_metadata_migration: Path | None = None, gateway_source: Path | None = None, routes_source: Path | None = None, command_consumer_source: Path | None = None, @@ -1910,6 +2208,9 @@ def build_gateway_report( "lineage_migration": ( lineage_migration or METADATA_FABRIC_LINEAGE_MIGRATION ).resolve(), + "active_metadata_migration": ( + active_metadata_migration or ACTIVE_METADATA_CHANGE_MIGRATION + ).resolve(), "gateway_source": (gateway_source or Path(__file__)).resolve(), "routes_source": (routes_source or GATEWAY_ROUTES_SOURCE).resolve(), "command_consumer_source": ( @@ -1989,6 +2290,17 @@ def build_gateway_report( "FORCE ROW LEVEL SECURITY", "GRANT SELECT, INSERT ON gda_control.metadata_fabric_lineage_outbox", ), + "active_metadata_migration": ( + "CREATE TABLE IF NOT EXISTS gda_control.metadata_change_outbox", + "FOREIGN KEY (", + "resource_version_id, content_sha256", + "FOR UPDATE SKIP LOCKED", + "claim_metadata_changes", + "complete_metadata_change", + "fail_metadata_change", + "ALTER TABLE gda_control.metadata_change_outbox FORCE ROW LEVEL SECURITY", + "GRANT SELECT, INSERT ON gda_control.metadata_change_outbox", + ), "gateway_source": ( 'SET LOCAL ROLE "{GATEWAY_DATABASE_ROLE}"', "SELECT set_config('app.current_tenant', :tenant, true)", @@ -2005,6 +2317,10 @@ def build_gateway_report( "def claim_metadata_fabric_lineage(", "def complete_metadata_fabric_lineage(", "def fail_metadata_fabric_lineage(", + "def register_resource_version_with_metadata_event(", + "def claim_metadata_changes(", + "def complete_metadata_change(", + "def fail_metadata_change(", ), "routes_source": ( 'base = "/api/platform/v1"', @@ -2057,6 +2373,7 @@ def build_gateway_report( or forbidden in texts.get("success_migration", "") or forbidden in texts.get("binding_migration", "") or forbidden in texts.get("lineage_migration", "") + or forbidden in texts.get("active_metadata_migration", "") ): errors.append(f"gateway role contains forbidden privilege: {forbidden}") consumer_source = texts.get("command_consumer_source", "") diff --git a/data_agent/platform_truth.py b/data_agent/platform_truth.py index b7837e30..b55c3b70 100644 --- a/data_agent/platform_truth.py +++ b/data_agent/platform_truth.py @@ -748,6 +748,26 @@ def _config( ), "Protected identity/TLS, production object storage and full Spark/Flink conformance", ), + RuntimeSpec( + "metadata_active_metadata_outbox_rehearsal", + "active_metadata_outbox_rehearsal", + "governed", + "evidence_durable", + "committed local PostgreSQL Active Metadata outbox evidence", + "metadata-platform", + "local_verification_only", + ( + "data_agent/metadata_fabric_active_metadata_outbox.py", + "scripts/metadata-fabric-active-metadata-outbox.sh", + ), + ( + ( + "data_agent/metadata_fabric_active_metadata_outbox.py", + "def run_local_rehearsal", + ), + ), + "Managed consumer submission to DolphinScheduler with protected identity", + ), RuntimeSpec( "datalake_monitor", "monitor_loop", diff --git a/data_agent/test_active_metadata_change_contract.py b/data_agent/test_active_metadata_change_contract.py new file mode 100644 index 00000000..e6f2b335 --- /dev/null +++ b/data_agent/test_active_metadata_change_contract.py @@ -0,0 +1,139 @@ +from datetime import datetime, timedelta, timezone +from uuid import UUID + +import pytest +from pydantic import ValidationError + +from data_agent.active_metadata_change_contract import ( + ActiveMetadataContractError, + MetadataActivationIntent, + MetadataChangeDelivery, + MetadataChangeEvent, + build_active_metadata_registration, + build_metadata_activation_intent, + build_metadata_change_delivery, +) +from data_agent.platform_contracts import ResourceVersion + + +TENANT = "tenant-a" +VERSION_ID = UUID("00000000-0000-4000-8000-000000000001") +NOW = datetime(2026, 7, 30, 8, 0, tzinfo=timezone.utc) +CONSUMER = "workload:metadata-router" + + +def _version(**overrides) -> ResourceVersion: + values = { + "tenant_id": TENANT, + "resource_urn": "gda://tenant-a/dataset/parcels", + "resource_version_id": VERSION_ID, + "version_key": "snapshot-1", + "content_sha256": "a" * 64, + "authority_version_ref": {"snapshot_id": 1}, + "created_by": "human:operator", + "created_at": NOW, + } + values.update(overrides) + return ResourceVersion(**values) + + +def test_registration_event_and_activation_intent_are_deterministic(): + first = build_active_metadata_registration( + _version(), + consumer_subject=CONSUMER, + ) + second = build_active_metadata_registration( + _version(), + consumer_subject=CONSUMER, + ) + intent = build_metadata_activation_intent( + first.event, + routed_by=CONSUMER, + ) + + assert first == second + assert str(first.event.event_id) == "23bce695-edf5-53ef-b266-f053628f3446" + assert first.event.event_sha256 == ( + "21fd3a5bc5446869412787bdee548bc58c876b980901f437fce99e166cc3e3d0" + ) + assert intent.route == "metadata_fabric.projection_plan" + assert intent.provider_apply_authorized is False + assert intent.provider_mutations_executed is False + assert intent.production_ingestion_verified is False + assert intent.intent_sha256 == ( + "169ac6b822d9af2ff75071eccc69a9468bffcacbe576bdad8430b95215d1a88c" + ) + + +def test_event_and_activation_intent_reject_content_tampering(): + registration = build_active_metadata_registration( + _version(), + consumer_subject=CONSUMER, + ) + event_payload = registration.event.model_dump(mode="json", by_alias=True) + event_payload["version_key"] = "snapshot-tampered" + with pytest.raises(ValidationError, match="SHA-256"): + MetadataChangeEvent.model_validate(event_payload) + + intent = build_metadata_activation_intent( + registration.event, + routed_by=CONSUMER, + ) + intent_payload = intent.model_dump(mode="json", by_alias=True) + intent_payload["resource_urn"] = "gda://tenant-a/dataset/private" + with pytest.raises(ValidationError, match="SHA-256"): + MetadataActivationIntent.model_validate(intent_payload) + + +def test_authenticated_producer_and_exact_consumer_are_required(): + with pytest.raises(ActiveMetadataContractError, match="authenticated subject"): + build_active_metadata_registration( + _version(created_by="anonymous"), + consumer_subject=CONSUMER, + ) + + registration = build_active_metadata_registration( + _version(), + consumer_subject=CONSUMER, + ) + with pytest.raises(ActiveMetadataContractError, match="event consumer"): + build_metadata_activation_intent( + registration.event, + routed_by="workload:other-router", + ) + + +def test_delivery_state_machine_rejects_incoherent_claims_and_terminal_state(): + event = build_active_metadata_registration( + _version(), + consumer_subject=CONSUMER, + ).event + pending = build_metadata_change_delivery(event, max_attempts=3) + assert pending.status.value == "pending" + assert pending.available_at == event.occurred_at + + payload = pending.model_dump(mode="json", by_alias=True) + payload.update( + { + "status": "in_flight", + "attempt_count": 1, + "claimed_by": "worker:router-1", + "claimed_until": (NOW + timedelta(minutes=1)).isoformat(), + } + ) + claimed = MetadataChangeDelivery.model_validate(payload) + assert claimed.status.value == "in_flight" + + payload["claimed_until"] = None + with pytest.raises(ValidationError, match="claim fields"): + MetadataChangeDelivery.model_validate(payload) + + terminal = pending.model_dump(mode="json", by_alias=True) + terminal.update( + { + "status": "processed", + "completed_at": (NOW + timedelta(minutes=2)).isoformat(), + } + ) + with pytest.raises(ValidationError, match="processed metadata change"): + MetadataChangeDelivery.model_validate(terminal) diff --git a/data_agent/test_metadata_fabric_active_metadata_outbox.py b/data_agent/test_metadata_fabric_active_metadata_outbox.py new file mode 100644 index 00000000..5bba7362 --- /dev/null +++ b/data_agent/test_metadata_fabric_active_metadata_outbox.py @@ -0,0 +1,45 @@ +import json +from copy import deepcopy + +from data_agent import metadata_fabric_active_metadata_outbox as outbox + + +def test_static_contract_binds_transactional_event_and_safe_activation_route(): + report = outbox.build_contract_report() + + assert report["status"] == "valid" + assert report["errors"] == [] + assert report["activation_route"] == "metadata_fabric.projection_plan" + assert report["consumer_subject"] == outbox.CONSUMER_SUBJECT + assert report["production_ready"] is False + + +def test_checked_evidence_is_current_content_bound_and_locally_scoped(): + evidence = json.loads( + outbox.DEFAULT_EVIDENCE_PATH.read_text(encoding="utf-8") + ) + + assert outbox.validate_rehearsal_evidence(evidence) == [] + assert evidence["local_postgresql_active_metadata_loop_verified"] is True + assert evidence["transactional_outbox_verified"] is True + assert evidence["legacy_backfill_blocked"] is True + assert evidence["provider_apply_authorized"] is False + assert evidence["provider_mutations_executed"] is False + assert evidence["production_ingestion_verified"] is False + assert evidence["production_scheduler_submission_verified"] is False + assert evidence["production_ready"] is False + + +def test_evidence_validation_rejects_tampering_and_production_overclaim(): + evidence = json.loads( + outbox.DEFAULT_EVIDENCE_PATH.read_text(encoding="utf-8") + ) + tampered = deepcopy(evidence) + tampered["authoritative_event_count"] = 2 + tampered["production_ready"] = True + + errors = outbox.validate_rehearsal_evidence(tampered) + + assert "Active Metadata evidence SHA-256 does not match" in errors + assert "local Active Metadata evidence may not claim production_ready" in errors + assert "Active Metadata evidence must contain exactly one event" in errors diff --git a/data_agent/test_metadata_fabric_active_metadata_outbox_postgres.py b/data_agent/test_metadata_fabric_active_metadata_outbox_postgres.py new file mode 100644 index 00000000..422a98b8 --- /dev/null +++ b/data_agent/test_metadata_fabric_active_metadata_outbox_postgres.py @@ -0,0 +1,111 @@ +import os +from uuid import uuid4 + +import pytest +from sqlalchemy import create_engine, text +from sqlalchemy.engine import make_url +from sqlalchemy.exc import DBAPIError + +from data_agent.metadata_fabric_active_metadata_outbox import ( + CONSUMER_SUBJECT, + TENANT, + build_active_metadata_bundle, + run_local_rehearsal, + validate_rehearsal_evidence, +) +from data_agent.platform_gateway import GatewayNotFoundError, PlatformGateway + + +DATABASE_URL = os.environ.get("DATABASE_URL") + + +def _temporary_database_url() -> tuple[object, str, str]: + admin_url = make_url(DATABASE_URL) + admin_engine = create_engine(admin_url, isolation_level="AUTOCOMMIT") + with admin_engine.connect() as connection: + is_superuser = connection.exec_driver_sql( + "SELECT rolsuper FROM pg_roles WHERE rolname = current_user" + ).scalar_one() + if not is_superuser: + admin_engine.dispose() + pytest.skip("Active Metadata test requires a PostgreSQL superuser") + database_name = f"gda_active_metadata_{uuid4().hex}" + connection.exec_driver_sql(f'CREATE DATABASE "{database_name}"') + database_url = admin_url.set(database=database_name).render_as_string( + hide_password=False + ) + return admin_engine, database_name, database_url + + +def _drop_temporary_database(admin_engine, database_name: str) -> None: + with admin_engine.connect() as connection: + connection.execute( + text( + """ + SELECT pg_terminate_backend(pid) + FROM pg_stat_activity + WHERE datname = :database_name + AND pid <> pg_backend_pid() + """ + ), + {"database_name": database_name}, + ) + connection.exec_driver_sql(f'DROP DATABASE "{database_name}"') + admin_engine.dispose() + + +@pytest.mark.skipif(not DATABASE_URL, reason="DATABASE_URL is not configured") +def test_postgres_active_metadata_outbox_is_atomic_scoped_and_retryable(): + admin_engine, database_name, database_url = _temporary_database_url() + engine = None + try: + evidence = run_local_rehearsal(database_url) + + assert validate_rehearsal_evidence(evidence) == [] + assert evidence["first_registration_created"] is True + assert evidence["exact_replay_created"] is False + assert evidence["final_attempt_count"] == 3 + assert evidence["legacy_backfill_blocked"] is True + assert evidence["authoritative_event_count"] == 1 + + bundle = build_active_metadata_bundle() + engine = create_engine(database_url) + gateway = PlatformGateway(engine) + stored = gateway.get_metadata_change_delivery( + TENANT, + bundle.registration.event.event_id, + ) + assert stored.status.value == "processed" + assert gateway.claim_metadata_changes( + TENANT, + "worker:post-test", + consumer_subject=CONSUMER_SUBJECT, + ) == [] + with pytest.raises(GatewayNotFoundError): + gateway.get_metadata_change_delivery( + "active-metadata-isolated", + bundle.registration.event.event_id, + ) + + with gateway._transaction(TENANT) as connection: + for statement in ( + """ + UPDATE gda_control.metadata_change_outbox + SET attempt_count = attempt_count + 1 + WHERE event_id = :event_id + """, + """ + DELETE FROM gda_control.metadata_change_outbox + WHERE event_id = :event_id + """, + ): + with pytest.raises(DBAPIError): + with connection.begin_nested(): + connection.execute( + text(statement), + {"event_id": bundle.registration.event.event_id}, + ) + finally: + if engine is not None: + engine.dispose() + _drop_temporary_database(admin_engine, database_name) diff --git a/data_agent/test_platform_gateway.py b/data_agent/test_platform_gateway.py index 9b84e8d5..bb356e7e 100644 --- a/data_agent/test_platform_gateway.py +++ b/data_agent/test_platform_gateway.py @@ -21,6 +21,7 @@ quality_result_fingerprint, ) from data_agent.platform_gateway import ( + ACTIVE_METADATA_CHANGE_MIGRATION, COMMAND_OUTBOX_MIGRATION, DefinitionRegistration, GATEWAY_ROLE_MIGRATION, @@ -549,3 +550,17 @@ def test_platform_gateway_static_contract_and_fail_closed_role(tmp_path): unsafe_report = build_gateway_report(command_migration=unsafe_command) assert unsafe_report["status"] == "invalid" assert "command_migration" in unsafe_report["missing_markers"] + + unsafe_active_metadata = tmp_path / "unsafe_active_metadata.sql" + unsafe_active_metadata.write_text( + ACTIVE_METADATA_CHANGE_MIGRATION.read_text(encoding="utf-8").replace( + "FOR UPDATE SKIP LOCKED", + "FOR UPDATE", + ), + encoding="utf-8", + ) + unsafe_report = build_gateway_report( + active_metadata_migration=unsafe_active_metadata + ) + assert unsafe_report["status"] == "invalid" + assert "active_metadata_migration" in unsafe_report["missing_markers"] diff --git a/data_agent/test_platform_truth.py b/data_agent/test_platform_truth.py index 7ff050d6..3c3f7526 100644 --- a/data_agent/test_platform_truth.py +++ b/data_agent/test_platform_truth.py @@ -242,6 +242,11 @@ def test_repository_source_access_and_runtime_baselines_match(): and item["production_role"] == "local_verification_only" for item in static_report["runtime"]["inventory"] ) + assert any( + item["runtime_id"] == "metadata_active_metadata_outbox_rehearsal" + and item["production_role"] == "local_verification_only" + for item in static_report["runtime"]["inventory"] + ) def test_runtime_report_detects_unregistered_background_mechanism(tmp_path): diff --git a/docs/architecture-decisions/adr-060-transactional-active-metadata-change-outbox.md b/docs/architecture-decisions/adr-060-transactional-active-metadata-change-outbox.md new file mode 100644 index 00000000..f4e6d56d --- /dev/null +++ b/docs/architecture-decisions/adr-060-transactional-active-metadata-change-outbox.md @@ -0,0 +1,62 @@ +# ADR-060: Transactional Active Metadata Change Outbox + +- Status: Accepted for local AR-1 verification +- Date: 2026-07-30 +- Owners: Metadata Platform / Data Platform + +## Context + +The Metadata Fabric can reconcile provider state, generate projection plans, +persist provider bindings, and deliver OpenLineage. It did not yet have an +authoritative change signal connecting a newly registered ResourceVersion to +an Active Metadata action. Polling the ledger would lose producer intent, +while creating an event after the version commit would permit missing events +and fabricated historical backfill. + +## Decision + +For the first event type, `resource_version.registered`, create the +ResourceVersion and its deterministic `MetadataChangeEvent` in one PostgreSQL +transaction. The event ID is derived from the ResourceVersion ID, and the +event fingerprint binds tenant, resource/version identity, predecessor, +content checksum, authenticated producer, workload consumer, and occurrence +time. + +Store delivery state in migration 099 under tenant-forced RLS. The gateway may +only select and insert directly; claim, complete, and fail transitions use +security-definer functions with workload scoping, worker ownership, leases, +bounded attempts, retry delay, and terminal finality. Exact replay creates +nothing. If only the ResourceVersion already exists, the transaction fails and +does not synthesize a historical event. + +The consumer output is a deterministic `metadata_fabric.projection_plan` +activation intent. It explicitly carries +`provider_apply_authorized=false`, `provider_mutations_executed=false`, and +`production_ingestion_verified=false`. This slice adds no resident worker or +scheduler. A later managed consumer must submit authorized work through +DolphinScheduler rather than execute provider mutations itself. + +## Consequences + +- Active Metadata now has a durable event spine tied to the platform version + authority instead of a catalog polling convention. +- Delivery is at least once. Event and activation fingerprints provide stable + idempotency; they do not claim network exactly once. +- Existing ResourceVersions remain historical records and are not silently + backfilled with events. +- Provider policy, production workload identity, protected scheduling, + production ingestion, and provider mutation evidence remain future gates. + +## Local Verification + +The checked PostgreSQL 16 evidence proves atomic registration, exact replay, +consumer scoping, wrong-worker rejection, retry, lease-expiry reclaim, +processed finality, tenant isolation, forced RLS, direct update/delete denial, +and rollback of a legacy-version event attempt. It records one authoritative +event, three delivery attempts, contract fingerprint +`3429cd1d7fc5015dab7dfd27b3972c4628238bc18d13c76ad28bb16697898e75`, +and evidence fingerprint +`d85a4575a6103e2f7107f8e11153c080a430e82f6cf6db547295f14ef909e96a`. + +This is local transactional evidence only and does not establish production +Active Metadata readiness. diff --git a/docs/evidence/metadata-fabric-active-metadata-outbox-2026-07-30.json b/docs/evidence/metadata-fabric-active-metadata-outbox-2026-07-30.json new file mode 100644 index 00000000..7c87aea0 --- /dev/null +++ b/docs/evidence/metadata-fabric-active-metadata-outbox-2026-07-30.json @@ -0,0 +1,33 @@ +{ + "activation_intent_sha256": "0380e458ac706d6e1f66f0a320177f1ea7956c15a93faa654d8e01323258fce9", + "activation_route": "metadata_fabric.projection_plan", + "authoritative_event_count": 1, + "contract_sha256": "3429cd1d7fc5015dab7dfd27b3972c4628238bc18d13c76ad28bb16697898e75", + "cross_tenant_read_blocked": true, + "errors": [], + "event_id": "52073ce1-db1f-526c-a6a7-51976c2ec8cd", + "event_sha256": "08087307bdd6694a3b2de2176fdb02ab7db720728b65b6da641e0b25dfa5edcf", + "evidence_sha256": "d85a4575a6103e2f7107f8e11153c080a430e82f6cf6db547295f14ef909e96a", + "exact_replay_created": false, + "final_attempt_count": 3, + "first_registration_created": true, + "force_rls_verified": true, + "gateway_select_insert_only_verified": true, + "lease_expiry_reclaim_verified": true, + "legacy_backfill_blocked": true, + "legacy_event_transaction_rolled_back": true, + "local_postgresql_active_metadata_loop_verified": true, + "processed_delivery_not_reclaimed": true, + "processed_replay_created": false, + "production_ingestion_verified": false, + "production_ready": false, + "production_scheduler_submission_verified": false, + "provider_apply_authorized": false, + "provider_mutations_executed": false, + "retry_pending_verified": true, + "schema": "gda.active_metadata_outbox_evidence.v1", + "status": "local_postgresql_active_metadata_loop_verified", + "transactional_outbox_verified": true, + "wrong_consumer_claim_blocked": true, + "wrong_worker_completion_blocked": true +} diff --git a/docs/roadmap-ar0-platform-truth-2026-07-24.md b/docs/roadmap-ar0-platform-truth-2026-07-24.md index f88babbc..7170c43b 100644 --- a/docs/roadmap-ar0-platform-truth-2026-07-24.md +++ b/docs/roadmap-ar0-platform-truth-2026-07-24.md @@ -199,7 +199,7 @@ Temporal 继续保持目标组件状态,不在这一包并行接入。OpenMeta 当前完成仅指本地合同、授权 evidence、outbox/callback 代码、数据库成功终局门、托管 worker 代码、默认关闭的部署模板及离线 activation/release preflight、candidate/registry/provenance/artifact-release/live observation evidence gate、合成 golden slice、定向测试、真实 PostgreSQL 16 事务边界和 canonical mainline 治理。`candidate_validated`、`registry_subject_bound`、本地合成 `provenance_verified`、`ready_for_activation`、`ready_for_staging_apply`、`verified_for_staging_apply` 和本地 live collection 都不等于真实镜像已 attested 或 staging 已部署;真实 IAM/OIDC 与 service token 生命周期、首次 GHCR publish/verify、真实 provenance artifact verify、registry-backed live staging revision、worker/callback 扩容运行、golden slice staging 运行链、受保护 release/live evidence provenance、独立 DolphinScheduler metadata PostgreSQL 和真实数据终局证据仍属于 4.7 后续切片。 -### 4.8 Metadata Fabric Bridge M1 + M2 + M3-13(本地 Spark uncertain-commit reconcile 已验证,生产验证待执行) +### 4.8 Metadata Fabric Bridge M1 + M2 + M3-14(本地 Active Metadata 事务 outbox 已验证,生产验证待执行) 第八块回到 AR-1 的 metadata control plane,以 [ADR-036](architecture-decisions/adr-036-read-only-metadata-fabric-bridge-contract.md) 固定 OpenMetadata + Gravitino + GDA Control Ledger 的首条 table slice: @@ -233,8 +233,9 @@ Temporal 继续保持目标组件状态,不在这一包并行接入。OpenMeta 28. [ADR-057](architecture-decisions/adr-057-production-object-store-readiness-gate.md) 已将 M3-10 evidence、S3-compatible provider/account/region/bucket、独立 failure domain、OIDC workload federation、精确八项 S3 permission、TLS/private path、KMS、versioning、cross-region replication、strong read/list consistency、tenant isolation、owner/SLO/runbook 和 26 项 protected attestation check 冻结为 fail-closed profile。当前 profile fingerprint 为 `668e194b3c688307014148391e7f389c9d6e9ca69c95d7b4cc92b4acae93181a`,report fingerprint 为 `85362dd10b7dc565f9fa567673d90b774cdec714bd1e70fb2c3c83c1af48b5ea`,合同有效但 43 项生产输入仍 blocked,全部 production claims 为 `false`。这只是 provider-neutral 决策和验收合同:没有选择、部署或验证 AWS S3、华为云 OBS 或其他生产对象存储;原生非 S3 provider 必须进入新的 conformance slice。 29. [ADR-058](architecture-decisions/adr-058-local-spark-commit-failure-recovery.md) 已在 Spark driver 的 loopback Iceberg REST proxy 中于 provider 转发前注入 HTTP 503。baseline 为 1 个 append snapshot、2 行和 1 个 referenced Parquet;失败调用经过精确 2 次 503 后,snapshot/row/file 均零漂移;对同一 `spark-recovery` 行做一次显式重试后为父子相连的 2 个 append snapshots、3 行和 2 个 referenced Parquet。直接 MinIO inventory 精确为 2 data + 3 metadata + 4 manifest = 9 objects,没有孤儿 data file;namespace、两块 PV 和 port-forward 均清理。contract fingerprint 为 `6d8944ab80246dc65891aa81118cb8b73f7ecad699be9a2af5e62d8260c41002`,evidence fingerprint 为 `39571cdac1e4043bcfc2d03a73b2b12ff925210daf8ae36bc640b8cb14d89401`。该结果只证明已知 pre-forward 失败下的本地原子性和一次显式重试,不证明 uncertain commit reconciliation、网络 exactly-once、生产对象存储或完整 engine conformance。 30. [ADR-059](architecture-decisions/adr-059-local-spark-uncertain-commit-reconciliation.md) 已将一个 armed commit 转发给 Gravitino,并在 provider 返回 200 后丢弃成功响应、向 Spark 返回 Iceberg `CommitStateUnknownException` 所需的 HTTP 504;一次传输重试被抑制。Spark 不重提逻辑写,而是 readback 得到父子相连的 2 个 append snapshots、3 行和 2 个 referenced Parquet,决策为 `committed_do_not_resubmit`、`write_resubmitted=false`。MinIO inventory 为 2 data + 3 metadata + 4 manifest = 9 objects;Job `Complete 1/1`,namespace、两块 PV 和 port-forward 均清理。contract fingerprint 为 `7a8d75a1d6b4558b982c6c3242d8d356c5046955f8aae7a45e5c297b6f4d4132`,evidence fingerprint 为 `d6462fff78d07047311b1f715d5f2c7f08c0ce8fbdd5c8b26a3d95ddc3474786`。该结果只证明一个本地 append 的确定性 readback/no-resubmit,不证明持久 reconcile controller、并发写、进程崩溃恢复、网络 exactly-once 或生产能力。 +31. [ADR-060](architecture-decisions/adr-060-transactional-active-metadata-change-outbox.md) 已新增 migration 099、内容绑定 `MetadataChangeEvent`、deterministic activation intent 与 PlatformGateway 原子注册/claim/fail/complete API。真实 PostgreSQL 16 演练中,ResourceVersion 与事件同事务创建,精确 replay 与 processed replay 均不新增事件;错误 consumer/worker 被拒绝,一次 retry 和一次强制租约过期后由第三个 worker 完成,旧 ResourceVersion 的补事件尝试整笔回滚。最终只有 1 条权威事件、3 次 attempt,FORCE RLS、跨租户拒绝与 gateway 无直接 UPDATE/DELETE 均通过。contract fingerprint 为 `3429cd1d7fc5015dab7dfd27b3972c4628238bc18d13c76ad28bb16697898e75`,evidence fingerprint 为 `d85a4575a6103e2f7107f8e11153c080a430e82f6cf6db547295f14ef909e96a`。该切片只生成 `metadata_fabric.projection_plan` 意图,不新增常驻 consumer、不提交 DolphinScheduler、不授权或执行 provider mutation,也不证明 production ingestion。 -此处 M1 只证明静态合同和只读 HTTP 边界;M2a 只证明本地 live foundation 与 PVC 重挂载连续性;M2b-1/M2b-2 分别限定在同集群新 PVC 和同集群隔离 repository;M2b-3 的 `local_cross_cluster_recovery_verified=true` 只限定在 `local_same_host_distinct_kubernetes_clusters_external_s3_repository`;M2c-1/M2c-2/M2c-3 分别限定本地 provider metrics、临时双周期 OTel 和单 job scrape recovery;M2c-4/M2d-2 只证明 production observability/NetworkPolicy profile 与 attestation 合同可校验;M2d-1 只证明本地两节点 kindnet 的隔离合成流量;M3-1 的 terminal evidence 与 M3-2 的 PolicyDecision/Approval 仍是 deterministic local fixtures。M3-2 只把 projection 写入本地 provider 并证明 retained target 的单次零写入 replay;M3-3 只把该本地 evidence 对应的 binding 写入临时 GDA Control 账本;M3-4 只向无认证 loopback receiver 发送精确 candidate 并验证 503 后幂等恢复;M3-5 只证明 OpenMetadata 在 provider 强制默认 role 之上的项目新增 grant 限定为 `table/Create`,以及本地 JWT 轮换/吊销和越权拒绝;M3-6 只证明隔离 Gravitino Basic IdP 的 bounded table-create、catalog-create 拒绝、登录轮换/吊销和完整清理;M3-7 只证明 pending production identity profile、profile-bound attestation 和派生 claim 的 fail-closed 合同可校验,没有部署或证明真实身份路径;M3-8 只证明同一 Docker Desktop 集群内 Basic 用户、JDBC metadata 与 file warehouse PVC 在受控 Pod restart 后连续;M3-9 只证明同节点共享 RWO PVC 的 Spark interoperability;M3-10 移除了该共享 PVC,并证明同一 Docker Desktop 主机/集群内 Spark 与 MinIO 的跨节点 S3-compatible 互操作,但不证明生产云对象存储、独立 failure domain、持久 identity binding、Flink 或完整 engine conformance;M3-11 只冻结 provider-neutral production object-store profile、精确 attestation binding 与 fail-closed claims,没有选择 provider、部署 bucket/KMS/policy 或提交真实 attestation;M3-12 只证明同一本地路径的 pre-forward commit failure 不改变可见 table state,随后一次显式重试产生一个新 snapshot/row,且无孤儿 data file;M3-13 只证明单次本地 append 在 provider 200 响应丢失并映射为 commit-state-unknown 后,可以由即时 table readback 判定 committed 且不重提,不覆盖持久 controller、进程崩溃、并发写或任意 mutation。M3-2 ingestion 仍使用 bootstrap admin,生产持久 binding、ResourceVersion 和 legacy authority 都未写入;生产对象存储、双 provider/生产最小权限、protected workload identity、OIDC、TLS、生产持久 catalog、tenant isolation、真实 receiver/alert/SLO、受保护 provider policy、生产故障注入、source-loss recovery、cancel/reconcile/lineage、完整 Spark/Flink conformance、生产 ingest、四项 production gate 和 `production_ready` 仍为 `false`。 +此处 M1 只证明静态合同和只读 HTTP 边界;M2a 只证明本地 live foundation 与 PVC 重挂载连续性;M2b-1/M2b-2 分别限定在同集群新 PVC 和同集群隔离 repository;M2b-3 的 `local_cross_cluster_recovery_verified=true` 只限定在 `local_same_host_distinct_kubernetes_clusters_external_s3_repository`;M2c-1/M2c-2/M2c-3 分别限定本地 provider metrics、临时双周期 OTel 和单 job scrape recovery;M2c-4/M2d-2 只证明 production observability/NetworkPolicy profile 与 attestation 合同可校验;M2d-1 只证明本地两节点 kindnet 的隔离合成流量;M3-1 的 terminal evidence 与 M3-2 的 PolicyDecision/Approval 仍是 deterministic local fixtures。M3-2 只把 projection 写入本地 provider 并证明 retained target 的单次零写入 replay;M3-3 只把该本地 evidence 对应的 binding 写入临时 GDA Control 账本;M3-4 只向无认证 loopback receiver 发送精确 candidate 并验证 503 后幂等恢复;M3-5 只证明 OpenMetadata 在 provider 强制默认 role 之上的项目新增 grant 限定为 `table/Create`,以及本地 JWT 轮换/吊销和越权拒绝;M3-6 只证明隔离 Gravitino Basic IdP 的 bounded table-create、catalog-create 拒绝、登录轮换/吊销和完整清理;M3-7 只证明 pending production identity profile、profile-bound attestation 和派生 claim 的 fail-closed 合同可校验,没有部署或证明真实身份路径;M3-8 只证明同一 Docker Desktop 集群内 Basic 用户、JDBC metadata 与 file warehouse PVC 在受控 Pod restart 后连续;M3-9 只证明同节点共享 RWO PVC 的 Spark interoperability;M3-10 移除了该共享 PVC,并证明同一 Docker Desktop 主机/集群内 Spark 与 MinIO 的跨节点 S3-compatible 互操作,但不证明生产云对象存储、独立 failure domain、持久 identity binding、Flink 或完整 engine conformance;M3-11 只冻结 provider-neutral production object-store profile、精确 attestation binding 与 fail-closed claims,没有选择 provider、部署 bucket/KMS/policy 或提交真实 attestation;M3-12 只证明同一本地路径的 pre-forward commit failure 不改变可见 table state,随后一次显式重试产生一个新 snapshot/row,且无孤儿 data file;M3-13 只证明单次本地 append 在 provider 200 响应丢失并映射为 commit-state-unknown 后,可以由即时 table readback 判定 committed 且不重提,不覆盖持久 controller、进程崩溃、并发写或任意 mutation;M3-14 只证明 ResourceVersion 注册与 Active Metadata 事件在本地 PostgreSQL 同事务创建,并验证租户/workload scoped claim/retry/complete,不包含常驻 consumer、DolphinScheduler submission、provider policy 或 provider mutation。M3-2 ingestion 仍使用 bootstrap admin,生产持久 binding、ResourceVersion 和 legacy authority 都未写入;生产对象存储、双 provider/生产最小权限、protected workload identity、OIDC、TLS、生产持久 catalog、tenant isolation、真实 receiver/alert/SLO、受保护 provider policy、生产故障注入、source-loss recovery、cancel/reconcile/lineage、完整 Spark/Flink conformance、生产 ingest、四项 production gate 和 `production_ready` 仍为 `false`。 ## 5. 重新评估条件 diff --git a/docs/system-of-record-matrix-2026-07-24.md b/docs/system-of-record-matrix-2026-07-24.md index 199c1c3c..96002392 100644 --- a/docs/system-of-record-matrix-2026-07-24.md +++ b/docs/system-of-record-matrix-2026-07-24.md @@ -2,9 +2,9 @@ 日期:2026-07-30 -阶段: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、M3-8 local Gravitino JDBC restart continuity、M3-9 local Spark/Iceberg REST interoperability、M3-10 local cross-node Spark/object-store interoperability、M3-11 production object-store readiness contract、M3-12 local Spark commit-failure recovery 与 M3-13 local uncertain-commit reconciliation 已验证,生产 provider ingestion、生产观测、生产 policy/tenant isolation、生产 identity/object-store 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、M3-8 local Gravitino JDBC restart continuity、M3-9 local Spark/Iceberg REST interoperability、M3-10 local cross-node Spark/object-store interoperability、M3-11 production object-store readiness contract、M3-12 local Spark commit-failure recovery、M3-13 local uncertain-commit reconciliation 与 M3-14 local Active Metadata transactional outbox 已验证,生产 provider ingestion、生产观测、生产 policy/tenant isolation、生产 identity/object-store attestation、生产 consumer/scheduler 和生产切换仍 `in_progress` -适用分支:`feat/ar1-metadata-fabric-spark-uncertain-commit-reconciliation` +适用分支:`feat/ar1-metadata-fabric-active-metadata-outbox` ## 判定规则 @@ -20,18 +20,18 @@ | SQL schema 历史 | PostgreSQL `schema_migrations`,以完整 migration ID + checksum 为权威 | migration CLI 的 JSON 报告 | 保持现有 ledger;任何 drift fail closed | Data Platform | AR-0,已验证 | | 部署配置策略 | Compose/K8s/进程环境;`platform_truth.CONFIG_SPECS` 定义关键类型与策略;DolphinScheduler worker 有默认零副本、外部 ConfigMap/Secret 驱动的 Kustomize 模板、静态 validator 和 staging activation preflight | `.env` 仅补默认;脱敏 snapshot、Secret key attestation、未扩容 Deployment 和 `ready_for_activation` 都是观测/模板 | 版本化 DeploymentProfile + secret reference;部署环境始终优先;模板或 preflight 通过都不等于环境已启用 | Platform/SRE/Security | AR-0,部分实现;worker 模板/preflight 本地已验证 | | 环境发布与晋级 | 本地 candidate/registry/provenance/release/live 合同已绑定 publisher、verifier、OCI 和 manifest identity;canonical `main@0182406`、archive refs、三组 active ruleset 与 `staging-provenance` protected environment 已建立,但尚无成功 publisher/verifier 或 deployment | 旧 mainline、feature branch、CI artifact、JSON、离线 report 和合成 `verified_for_staging_apply` 都不能单独成为发布权威;publisher SHA、verifier SHA 与 branch lineage 必须分别验证 | 由受保护 environment 的 DeploymentRevision 绑定 OCI、provenance artifact、release manifest 与全部 live verdict | Platform/SRE/Security/Repository Owner | AR-1 mainline 治理已恢复 -> 首次 GHCR publish/verify -> 真实 staging | -| 后台运行时清单 | `platform_truth.RUNTIME_INVENTORY` 是代码层登记;`gda_control` 已有受控 PlatformRun 写入口;DolphinScheduler managed worker 已登记但尚无生产调用方;M2b recovery runner、M2c-1 provider probe、M2c-2 `_OtelPortForward`、M2c-3 failure rehearsal、M2d-1 NetworkPolicy rehearsal、M3-8 JDBC restart runner、M3-9 Spark interoperability runner、M3-10 Spark/object-store runner 与 M3-12 Spark commit-failure runner 均登记为 `local_verification_only`,不是 scheduler、worker、持续监控、生产 policy/catalog controller 或状态权威 | AST primitive report、worker status JSON、FrameworkAttemptObservation、DolphinScheduler instance state、本地 recovery/metrics/network-policy/catalog/interoperability/failure evidence | PlatformRun ledger 唯一登记最终状态;framework/provider attempt 只能回报观测;本地演练进程与 evidence 不得变成生产控制器、监控后端、catalog authority 或 tenant-isolation 权威 | Platform Architecture | AR-1 adapter/worker 本地已验证;metadata recovery/metrics/policy/catalog/interoperability runner 仅本地验证 -> staging 控制链待接入 | +| 后台运行时清单 | `platform_truth.RUNTIME_INVENTORY` 是代码层登记;`gda_control` 已有受控 PlatformRun 写入口;DolphinScheduler managed worker 已登记但尚无生产调用方;Metadata Fabric recovery/metrics/policy/catalog/interoperability/failure/uncertain-commit 与 Active Metadata outbox rehearsal 均登记为 `local_verification_only`,不是 scheduler、worker、持续监控、生产 policy/catalog controller 或状态权威 | AST primitive report、worker status JSON、FrameworkAttemptObservation、DolphinScheduler instance state、本地 recovery/metrics/network-policy/catalog/interoperability/failure/outbox evidence | PlatformRun ledger 唯一登记最终状态;framework/provider attempt 只能回报观测;本地演练进程与 evidence 不得变成生产控制器、监控后端、catalog authority 或 tenant-isolation 权威 | Platform Architecture | AR-1 adapter/worker 本地已验证;metadata runner 仅本地验证 -> staging 控制链待接入 | | 原始文件/对象 | 当前 local uploads、S3/MinIO/OBS 均可能被直接写入,权威边界未统一 | 临时上传、下载缓存、预览文件 | Landing object 以 immutable URI + checksum + retention 为权威;本地 scratch 可删除 | Data Platform | AR-2 | | 湖仓表与 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 已验证 -> 生产切换待验收 | +| 数据资产身份与版本 | `gda_control.resource/resource_version` 已实现 identity、hash、predecessor、tenant FK 和幂等 gateway 写入;M3-14 新入口将新 ResourceVersion 与确定性 `resource_version.registered` 事件同事务创建,拒绝为旧版本事后补事件;`agent_data_assets`、`agent_asset_versions` 仍是兼容写路径 | UI catalog、search index、STAC、Active Metadata delivery state | GDA ledger 管身份与版本绑定;只有新原子注册路径产生权威变化事件,旧行不得伪装成实时事件;OpenMetadata 管治理目录,Gravitino 管技术对象映射 | Metadata Platform | AR-1 gateway/M3-14 本地事务已验证 -> 生产写入口切换待验收 | | 技术元数据 | M1 已冻结 Gravitino table ref/reconciliation;M2 已验证本地 foundation/recovery/metrics/policy 和 production readiness contracts;M3-1 固定 technical projection intent,M3-2 已在 Gravitino memory catalog 创建/read-back,M3-3 将验证后的 ref 追加到 tenant-scoped 本地 binding ledger;M3-6 又在隔离 Gravitino Basic IdP 中验证 bounded table-create、catalog-create 拒绝、密码轮换/用户吊销和完整清理;M3-7 已冻结 production identity profile/attestation gate;M3-8 已验证同一 Basic role 与 Iceberg JDBC table 在 PostgreSQL/Gravitino Pod restart 后保持;M3-9 已验证 Spark 经标准 Iceberg REST 对同一 JDBC catalog 做 read/write/schema evolution/snapshot/time travel;M3-10 已移除共享 warehouse PVC,并由跨节点 MinIO 对象检查与 Gravitino API 回读验证 Spark 结果;M3-11 已冻结 provider-neutral production object-store profile/attestation gate;M3-12 已验证 pre-forward 503 下失败提交零可见漂移、一次显式重试和无孤儿 data file | harvester 结果、合成 response、本地 sandbox/recovery/metrics/policy/ingestion/identity/JDBC restart/Spark/object-store/commit-failure observation、projection plan、provider evidence、binding ledger 与 readiness report | 源系统技术对象是原始证据;Gravitino 映射并联邦,不能覆盖业务 ResourceVersion;GDA binding ledger 只记录已验证关系,本地 memory/JDBC/file-PVC/MinIO catalog evidence 不得冒充生产持久技术权威;Basic IdP、无认证 REST/HTTP、本地主机对象存储、pending profile 和合成 attestation 都不是生产身份或生产 storage,M3-11 合同不构成 provider selection/deployment,M3-12 本地证据不构成网络 exactly-once 或生产 reconcile | Metadata Platform | AR-1 M1/M2 + M3-6 identity + M3-7 gate + M3-8 local persistence + M3-9/M3-10 interoperability + M3-11 object-store gate + M3-12 commit-failure recovery 已验证 -> 受保护身份/生产对象存储 attestation/uncertain outcome reconcile/完整 Spark-Flink conformance 待执行 | | 治理目录 | M1 已冻结 OpenMetadata table ref/reconciliation;M2 已验证本地 foundation/recovery/metrics/policy 和 production readiness contracts;M3-2 已用 bootstrap admin 创建目标并回读真实 UUID;M3-3 将该 UUID 经 evidence gate 追加到本地 GDA binding ledger;M3-4 将精确 OpenLineage candidate 经 outbox 投递到本地 HTTP receiver;M3-5 已验证临时非管理员 bot 的 scoped `table/Create` grant、policy-create 拒绝及 JWT 轮换/吊销;M3-7 将其 allow/deny 范围纳入双 provider production identity gate | 搜索/页面视图、合成 response、本地 sandbox/recovery/metrics/policy/ingestion/identity observation、projection/provider evidence、binding ledger、lineage outbox/receipt、OpenLineage event 与 readiness report | OpenMetadata 为 owner/glossary/classification/quality discoverability 权威;GDA ledger 保留审批/provider identity,outbox 只拥有投递状态,receiver 拥有接收状态;pending identity profile 和合成 attestation 均不反写 ResourceVersion 或建立生产权威 | Governance | AR-1 M1/M2 + M3-5 local bounded identity + M3-7 readiness contract 已验证 -> protected identity ingestion/生产持久 binding/受保护 production receiver 待执行 | | 血缘 | `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/生产切换待验收 | | 调度与补数 | 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;其他 WebSocket/bot/feedback 多为 best effort | command delivery status、消费者 claim、worker status JSON、WebSocket 消息 | command/event 与源事实同事务入 outbox,幂等 consumer 交付;worker status、outbox 状态、缓存或 socket 都不是 Run/业务权威 | Platform/Integrations | AR-1 command delivery/worker 代码已验证 -> staging worker/callback 待部署 | +| 事件交付 | Standards outbox 已数据库耐久;`platform_command_outbox` 支持 DolphinScheduler dispatch/reconcile;M3-14 `metadata_change_outbox` 将新 ResourceVersion 与内容绑定事件同事务写入,并提供 tenant/workload scoped claim/retry/complete、租约和终态,当前没有常驻 consumer | command/metadata delivery status、消费者 claim、activation intent、worker status JSON、WebSocket 消息 | command/event 与源事实同事务入 outbox,幂等 consumer 交付;Active Metadata consumer 只生成 projection-plan intent 并通过 DolphinScheduler 提交授权工作,不能自行执行 provider mutation | Platform/Integrations/Metadata Platform | AR-1 command worker 代码与 M3-14 本地 outbox 已验证 -> staging managed consumer/scheduler 待部署 | | 质量结果 | `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 待接入 | | 标准与语义定义 | `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;M3-8 证明同一 Gravitino Basic role 在本地 JDBC restart 后连续 | session/cache、前端菜单权限、本地 provider identity/JDBC restart 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、restart continuity 或人工批准均不可替代 | Security | AR-1 local identities/persistence + production readiness contract 已验证 -> protected 双 provider IAM/attestation 待执行 | @@ -57,6 +57,7 @@ 12. QualityResult evaluator 必须是 workload,且成功终局中的 evaluator 不能等于 Run workload;该代码级职责分离不替代生产 IAM。 13. `candidate_validated`、`registry_subject_bound`、GitHub provenance action 成功、CI artifact、离线 preflight、未独立 attested 的 live observation JSON 或人工批准都不能单独授权 production;缺少同一 source revision 的 OCI subject 独立验证、registry/live revision/identity/health/golden-slice 绑定及受保护 provenance 时,promotion 必须失败。 14. Metadata Fabric M1 只允许 OpenMetadata/Gravitino GET;M2 只执行本地 foundation/recovery/metrics/policy 演练或验证 production readiness profile;M3-1 只从 synthetic terminal evidence 生成 plan/candidate;M3-2 只允许 exact local PolicyDecision/Approval 后向本地 provider 写 projection;M3-3 只将同一 source evidence 经 PlatformGateway 写入临时 append-only binding ledger;M3-4 只经 tenant-scoped outbox 向无认证 loopback receiver 投递精确 candidate,并验证 at-least-once + receiver idempotency;M3-5 只证明临时 OpenMetadata bot 在 provider 强制 `DefaultBotRole` 之上的项目新增 grant 是 `table/Create`,并验证 policy-create 拒绝与本地 JWT 轮换/吊销;M3-6 只证明隔离 Gravitino Basic user 的 bounded table-create、catalog-create 拒绝、密码轮换/用户吊销和完整清理;M3-7 只冻结 production identity profile、精确 attestation binding 与 fail-closed 派生 claims,既不部署 identity path,也不提交真实 production attestation;M3-8 只证明 Docker Desktop 单集群中 Basic role、PostgreSQL JDBC metadata 与 file warehouse PVC 在受控 Pod restart 后连续;M3-9 只证明同节点共享 RWO PVC 的 Spark/Iceberg REST 互操作;M3-10 移除 Spark/Gravitino 共享 warehouse PVC,且只证明同一 Docker Desktop 主机/集群内跨节点 MinIO 的 read/write/schema evolution/snapshot/time travel 与对象级 metadata 一致;M3-11 只冻结 S3-compatible production profile、精确 attestation binding 与 fail-closed claims,既不选择/部署 provider,也不创建 bucket/KMS/policy 或提交真实 attestation;M3-12 只在 Spark driver loopback proxy 中于转发前注入 503,证明失败尝试零可见漂移、随后一次显式重试和直接对象清单无孤儿 data file;M3-13 只在 provider 200 响应丢失后以 HTTP 504 触发 commit-state-unknown,并对一个本地 append 做即时 readback/no-resubmit,不构成持久生产 reconcile controller、crash/concurrency proof 或网络 exactly-once。生产对象存储、identity/TLS、受保护环境故障注入、source-loss recovery、cancel/reconcile/lineage 与 Flink 仍未证明。Gravitino `1.3.0` Basic IdP 不算 OIDC,生产必须明确选择并证明 custom OIDC authenticator 或 identity-aware proxy。本地 bootstrap provisioner、Basic IdP、loopback/cluster HTTP、memory/file-backed JDBC catalog、同节点 PVC/MinIO、临时 identity/ledger/outbox、loopback receiver、pending profile、合成 attestation 和 local evidence 都不等于双 provider/生产最小权限、protected workload identity/OIDC、生产持久 catalog/binding、TLS、受保护 OpenLineage receiver、tenant isolation、alert/SLO、生产 ingestion/conformance 或生产写权威。 +15. M3-14 只证明本地 PostgreSQL 16 中 ResourceVersion 与 `resource_version.registered` 事件同事务创建,以及 tenant/workload scoped claim、retry、lease reclaim 和 complete;它没有部署常驻 consumer、没有提交 DolphinScheduler、没有取得 provider apply 授权、没有执行 provider mutation,也不构成 production Active Metadata readiness。 ## 已建立的 AR-0/AR-1 entry 证据 @@ -93,6 +94,7 @@ - Metadata Fabric M3-11 已建立 production object-store profile/attestation gate;checked-in profile fingerprint 为 `668e194b3c688307014148391e7f389c9d6e9ca69c95d7b4cc92b4acae93181a`,report fingerprint 为 `85362dd10b7dc565f9fa567673d90b774cdec714bd1e70fb2c3c83c1af48b5ea`,`profile_valid=true`,43 项 provider/identity/transport/encryption/durability/consistency/tenancy/operations 外部输入以 blockers 暴露,`ready_for_protected_verification=false`、`attestation_valid=false`、`production_object_store_gate_passed=false`、`production_ready=false`。该合同绑定 M3-10 evidence,但没有选择或部署 provider;合成完整 attestation 只验证门禁逻辑,不计入生产证据。下一项真实证据是经 owner 批准并物化的 provider profile,以及来自 `production-object-store` 受保护环境、绑定当前 source/profile 并通过全部 26 项检查的 attestation。 - Metadata Fabric M3-12 已在本地 Spark driver loopback Iceberg REST proxy 中于 provider 转发前注入 HTTP 503。baseline 为 1 个 append snapshot、2 行和 1 个 referenced Parquet;失败写经 2 次 503 后 snapshot/row/file 零漂移;对同一逻辑行一次显式重试后为父子相连的 2 个 append snapshots、3 行和 2 个 referenced Parquet。直接 MinIO inventory 精确为 2 data + 3 metadata + 4 manifest = 9 objects,没有孤儿 data file;namespace、两块 PV 和 port-forward 均清理。contract fingerprint 为 `6d8944ab80246dc65891aa81118cb8b73f7ecad699be9a2af5e62d8260c41002`,evidence fingerprint 为 `39571cdac1e4043bcfc2d03a73b2b12ff925210daf8ae36bc640b8cb14d89401`。这只证明已知 pre-forward failure 的本地原子性与一次显式重试,不证明 provider uncertain outcome reconcile、网络 exactly-once、生产对象存储、cancel/lineage、Flink 或完整 Spark conformance。 - Metadata Fabric M3-13 已在本地 Spark driver loopback proxy 将 armed commit 转发给 Gravitino,并在 provider 200 后丢弃成功响应、返回 HTTP 504;Iceberg 将其映射为 commit-state-unknown,一次传输重试被抑制。Spark 只读 readback 后输出 `committed_do_not_resubmit` 和 `write_resubmitted=false`,最终为父子相连的 2 个 append snapshots、3 行、2 个 referenced Parquet;MinIO 为 2 data + 3 metadata + 4 manifest = 9 objects。Job `Complete 1/1`,namespace、两块 PV 和 port-forward 均清理。contract fingerprint 为 `7a8d75a1d6b4558b982c6c3242d8d356c5046955f8aae7a45e5c297b6f4d4132`,evidence fingerprint 为 `d6462fff78d07047311b1f715d5f2c7f08c0ce8fbdd5c8b26a3d95ddc3474786`。这只证明一个本地 append 的确定性 readback/no-resubmit;持久 controller、crash/concurrency、网络 exactly-once、生产对象存储、cancel/lineage、Flink 和完整 Spark conformance 仍未证明。 +- Metadata Fabric M3-14 已新增内容绑定 `MetadataChangeEvent`、migration 099 transactional outbox 与 PlatformGateway 原子注册/claim/fail/complete。真实 PostgreSQL 16 演练验证首次注册 `created=true`、pending/processed 精确 replay 均 `created=false`、错误 consumer/worker 拒绝、retry、lease-expiry reclaim、第三次 attempt 完成、processed 不再认领、旧 ResourceVersion 补事件整笔回滚、跨租户不可见、FORCE RLS 和 gateway 无直接 UPDATE/DELETE;最终只有 1 条权威事件。contract fingerprint 为 `3429cd1d7fc5015dab7dfd27b3972c4628238bc18d13c76ad28bb16697898e75`,evidence fingerprint 为 `d85a4575a6103e2f7107f8e11153c080a430e82f6cf6db547295f14ef909e96a`。激活意图只路由到 `metadata_fabric.projection_plan`,`provider_apply_authorized=false`、`provider_mutations_executed=false`、`production_ingestion_verified=false`、`production_scheduler_submission_verified=false`、`production_ready=false`。 ## 下一验收证据 @@ -100,6 +102,7 @@ - 真实 provenance artifact verify、受保护 overlay 的 `verified_for_staging_apply` release report,以及 staging/production 的 schema、config/runtime snapshot、registry/live DeploymentRevision 绑定、release/live artifact attestation 和环境 compare 报告; - staging 的 migration role、应用 login membership、连接池 role/tenant 复位、双租户 API 和 success finalization 运行产物; - DolphinScheduler adapter 的真实 IAM/OIDC、service token provisioning/轮换、provider 最小权限、binding artifact staging 接入、managed outbox worker/provider callback 实际扩容部署、唯一 worker ID、status/lease 故障恢复和无双写证据; +- Active Metadata managed consumer 的受保护 workload identity、DolphinScheduler submission、幂等 projection execution、policy/approval gate、provider read-back、重试/死信告警与生产 SLO 证据; - 首条真实图斑链对 golden slice 的 output hash、独立质量结果/evidence、血缘、发布 revision 和 rollback 演练; - OpenMetadata/Gravitino 的 source host/cluster 外生产 backup account/bucket、已批准 provider profile 与受保护对象存储 attestation、生产对象存储、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、受保护环境 commit failure injection、provider uncertain outcome reconcile、cancel/reconcile/lineage 和完整 Spark/Flink 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-8 本地 JDBC restart continuity、M3-9 本地同节点 Spark interoperability、M3-10 本地同主机跨节点 MinIO interoperability、M3-11 pending object-store profile/合成 attestation、M3-12 local pre-forward commit-failure recovery 与 M3-13 local uncertain-commit readback/no-resubmit 均不计入生产退出门; - DolphinScheduler/Temporal sandbox 的独立数据库、备份恢复、身份、版本和升级责任证明;DolphinScheduler standalone/H2 不计入此退出门。 diff --git a/scripts/metadata-fabric-active-metadata-outbox.sh b/scripts/metadata-fabric-active-metadata-outbox.sh new file mode 100755 index 00000000..1bf8621c --- /dev/null +++ b/scripts/metadata-fabric-active-metadata-outbox.sh @@ -0,0 +1,4 @@ +#!/usr/bin/env bash +set -euo pipefail + +python -m data_agent.metadata_fabric_active_metadata_outbox "$@" From 2506bc5dfeb192317d643ad7893496520e6de39b Mon Sep 17 00:00:00 2001 From: Ning Zhou Date: Thu, 30 Jul 2026 10:36:25 +0800 Subject: [PATCH 2/3] fix(platform): stabilize active metadata evidence fingerprint --- data_agent/metadata_fabric_active_metadata_outbox.py | 7 +++++-- data_agent/test_metadata_fabric_active_metadata_outbox.py | 4 ++++ .../adr-060-transactional-active-metadata-change-outbox.md | 4 ++-- .../metadata-fabric-active-metadata-outbox-2026-07-30.json | 4 ++-- docs/roadmap-ar0-platform-truth-2026-07-24.md | 2 +- docs/system-of-record-matrix-2026-07-24.md | 2 +- 6 files changed, 15 insertions(+), 8 deletions(-) diff --git a/data_agent/metadata_fabric_active_metadata_outbox.py b/data_agent/metadata_fabric_active_metadata_outbox.py index e0a8fef2..31023834 100644 --- a/data_agent/metadata_fabric_active_metadata_outbox.py +++ b/data_agent/metadata_fabric_active_metadata_outbox.py @@ -168,12 +168,15 @@ def build_contract_report() -> dict[str, Any]: for name, path in paths.items(): if not path.is_file(): errors.append(f"{name} is missing") - files[name] = {"path": path.resolve().as_posix(), "sha256": None} + files[name] = { + "path": path.resolve().relative_to(REPO_ROOT).as_posix(), + "sha256": None, + } continue raw = path.read_bytes() source = raw.decode("utf-8") files[name] = { - "path": path.resolve().as_posix(), + "path": path.resolve().relative_to(REPO_ROOT).as_posix(), "sha256": hashlib.sha256(raw).hexdigest(), } missing = [marker for marker in required[name] if marker not in source] diff --git a/data_agent/test_metadata_fabric_active_metadata_outbox.py b/data_agent/test_metadata_fabric_active_metadata_outbox.py index 5bba7362..f4e50fbf 100644 --- a/data_agent/test_metadata_fabric_active_metadata_outbox.py +++ b/data_agent/test_metadata_fabric_active_metadata_outbox.py @@ -12,6 +12,10 @@ def test_static_contract_binds_transactional_event_and_safe_activation_route(): assert report["activation_route"] == "metadata_fabric.projection_plan" assert report["consumer_subject"] == outbox.CONSUMER_SUBJECT assert report["production_ready"] is False + assert all( + not item["path"].startswith("/") + for item in report["files"].values() + ) def test_checked_evidence_is_current_content_bound_and_locally_scoped(): diff --git a/docs/architecture-decisions/adr-060-transactional-active-metadata-change-outbox.md b/docs/architecture-decisions/adr-060-transactional-active-metadata-change-outbox.md index f4e6d56d..92bc40d1 100644 --- a/docs/architecture-decisions/adr-060-transactional-active-metadata-change-outbox.md +++ b/docs/architecture-decisions/adr-060-transactional-active-metadata-change-outbox.md @@ -54,9 +54,9 @@ consumer scoping, wrong-worker rejection, retry, lease-expiry reclaim, processed finality, tenant isolation, forced RLS, direct update/delete denial, and rollback of a legacy-version event attempt. It records one authoritative event, three delivery attempts, contract fingerprint -`3429cd1d7fc5015dab7dfd27b3972c4628238bc18d13c76ad28bb16697898e75`, +`b4b9d77d716f0e5d62389378aec60bdc366b2d847ace07a6b476310eb3ff6732`, and evidence fingerprint -`d85a4575a6103e2f7107f8e11153c080a430e82f6cf6db547295f14ef909e96a`. +`fcd44cf9043e144016644b0d04e72f348a777668d853081d9d21ded341bb2b59`. This is local transactional evidence only and does not establish production Active Metadata readiness. diff --git a/docs/evidence/metadata-fabric-active-metadata-outbox-2026-07-30.json b/docs/evidence/metadata-fabric-active-metadata-outbox-2026-07-30.json index 7c87aea0..544faaf8 100644 --- a/docs/evidence/metadata-fabric-active-metadata-outbox-2026-07-30.json +++ b/docs/evidence/metadata-fabric-active-metadata-outbox-2026-07-30.json @@ -2,12 +2,12 @@ "activation_intent_sha256": "0380e458ac706d6e1f66f0a320177f1ea7956c15a93faa654d8e01323258fce9", "activation_route": "metadata_fabric.projection_plan", "authoritative_event_count": 1, - "contract_sha256": "3429cd1d7fc5015dab7dfd27b3972c4628238bc18d13c76ad28bb16697898e75", + "contract_sha256": "b4b9d77d716f0e5d62389378aec60bdc366b2d847ace07a6b476310eb3ff6732", "cross_tenant_read_blocked": true, "errors": [], "event_id": "52073ce1-db1f-526c-a6a7-51976c2ec8cd", "event_sha256": "08087307bdd6694a3b2de2176fdb02ab7db720728b65b6da641e0b25dfa5edcf", - "evidence_sha256": "d85a4575a6103e2f7107f8e11153c080a430e82f6cf6db547295f14ef909e96a", + "evidence_sha256": "fcd44cf9043e144016644b0d04e72f348a777668d853081d9d21ded341bb2b59", "exact_replay_created": false, "final_attempt_count": 3, "first_registration_created": true, diff --git a/docs/roadmap-ar0-platform-truth-2026-07-24.md b/docs/roadmap-ar0-platform-truth-2026-07-24.md index 7170c43b..51aebd6a 100644 --- a/docs/roadmap-ar0-platform-truth-2026-07-24.md +++ b/docs/roadmap-ar0-platform-truth-2026-07-24.md @@ -233,7 +233,7 @@ Temporal 继续保持目标组件状态,不在这一包并行接入。OpenMeta 28. [ADR-057](architecture-decisions/adr-057-production-object-store-readiness-gate.md) 已将 M3-10 evidence、S3-compatible provider/account/region/bucket、独立 failure domain、OIDC workload federation、精确八项 S3 permission、TLS/private path、KMS、versioning、cross-region replication、strong read/list consistency、tenant isolation、owner/SLO/runbook 和 26 项 protected attestation check 冻结为 fail-closed profile。当前 profile fingerprint 为 `668e194b3c688307014148391e7f389c9d6e9ca69c95d7b4cc92b4acae93181a`,report fingerprint 为 `85362dd10b7dc565f9fa567673d90b774cdec714bd1e70fb2c3c83c1af48b5ea`,合同有效但 43 项生产输入仍 blocked,全部 production claims 为 `false`。这只是 provider-neutral 决策和验收合同:没有选择、部署或验证 AWS S3、华为云 OBS 或其他生产对象存储;原生非 S3 provider 必须进入新的 conformance slice。 29. [ADR-058](architecture-decisions/adr-058-local-spark-commit-failure-recovery.md) 已在 Spark driver 的 loopback Iceberg REST proxy 中于 provider 转发前注入 HTTP 503。baseline 为 1 个 append snapshot、2 行和 1 个 referenced Parquet;失败调用经过精确 2 次 503 后,snapshot/row/file 均零漂移;对同一 `spark-recovery` 行做一次显式重试后为父子相连的 2 个 append snapshots、3 行和 2 个 referenced Parquet。直接 MinIO inventory 精确为 2 data + 3 metadata + 4 manifest = 9 objects,没有孤儿 data file;namespace、两块 PV 和 port-forward 均清理。contract fingerprint 为 `6d8944ab80246dc65891aa81118cb8b73f7ecad699be9a2af5e62d8260c41002`,evidence fingerprint 为 `39571cdac1e4043bcfc2d03a73b2b12ff925210daf8ae36bc640b8cb14d89401`。该结果只证明已知 pre-forward 失败下的本地原子性和一次显式重试,不证明 uncertain commit reconciliation、网络 exactly-once、生产对象存储或完整 engine conformance。 30. [ADR-059](architecture-decisions/adr-059-local-spark-uncertain-commit-reconciliation.md) 已将一个 armed commit 转发给 Gravitino,并在 provider 返回 200 后丢弃成功响应、向 Spark 返回 Iceberg `CommitStateUnknownException` 所需的 HTTP 504;一次传输重试被抑制。Spark 不重提逻辑写,而是 readback 得到父子相连的 2 个 append snapshots、3 行和 2 个 referenced Parquet,决策为 `committed_do_not_resubmit`、`write_resubmitted=false`。MinIO inventory 为 2 data + 3 metadata + 4 manifest = 9 objects;Job `Complete 1/1`,namespace、两块 PV 和 port-forward 均清理。contract fingerprint 为 `7a8d75a1d6b4558b982c6c3242d8d356c5046955f8aae7a45e5c297b6f4d4132`,evidence fingerprint 为 `d6462fff78d07047311b1f715d5f2c7f08c0ce8fbdd5c8b26a3d95ddc3474786`。该结果只证明一个本地 append 的确定性 readback/no-resubmit,不证明持久 reconcile controller、并发写、进程崩溃恢复、网络 exactly-once 或生产能力。 -31. [ADR-060](architecture-decisions/adr-060-transactional-active-metadata-change-outbox.md) 已新增 migration 099、内容绑定 `MetadataChangeEvent`、deterministic activation intent 与 PlatformGateway 原子注册/claim/fail/complete API。真实 PostgreSQL 16 演练中,ResourceVersion 与事件同事务创建,精确 replay 与 processed replay 均不新增事件;错误 consumer/worker 被拒绝,一次 retry 和一次强制租约过期后由第三个 worker 完成,旧 ResourceVersion 的补事件尝试整笔回滚。最终只有 1 条权威事件、3 次 attempt,FORCE RLS、跨租户拒绝与 gateway 无直接 UPDATE/DELETE 均通过。contract fingerprint 为 `3429cd1d7fc5015dab7dfd27b3972c4628238bc18d13c76ad28bb16697898e75`,evidence fingerprint 为 `d85a4575a6103e2f7107f8e11153c080a430e82f6cf6db547295f14ef909e96a`。该切片只生成 `metadata_fabric.projection_plan` 意图,不新增常驻 consumer、不提交 DolphinScheduler、不授权或执行 provider mutation,也不证明 production ingestion。 +31. [ADR-060](architecture-decisions/adr-060-transactional-active-metadata-change-outbox.md) 已新增 migration 099、内容绑定 `MetadataChangeEvent`、deterministic activation intent 与 PlatformGateway 原子注册/claim/fail/complete API。真实 PostgreSQL 16 演练中,ResourceVersion 与事件同事务创建,精确 replay 与 processed replay 均不新增事件;错误 consumer/worker 被拒绝,一次 retry 和一次强制租约过期后由第三个 worker 完成,旧 ResourceVersion 的补事件尝试整笔回滚。最终只有 1 条权威事件、3 次 attempt,FORCE RLS、跨租户拒绝与 gateway 无直接 UPDATE/DELETE 均通过。contract fingerprint 为 `b4b9d77d716f0e5d62389378aec60bdc366b2d847ace07a6b476310eb3ff6732`,evidence fingerprint 为 `fcd44cf9043e144016644b0d04e72f348a777668d853081d9d21ded341bb2b59`。该切片只生成 `metadata_fabric.projection_plan` 意图,不新增常驻 consumer、不提交 DolphinScheduler、不授权或执行 provider mutation,也不证明 production ingestion。 此处 M1 只证明静态合同和只读 HTTP 边界;M2a 只证明本地 live foundation 与 PVC 重挂载连续性;M2b-1/M2b-2 分别限定在同集群新 PVC 和同集群隔离 repository;M2b-3 的 `local_cross_cluster_recovery_verified=true` 只限定在 `local_same_host_distinct_kubernetes_clusters_external_s3_repository`;M2c-1/M2c-2/M2c-3 分别限定本地 provider metrics、临时双周期 OTel 和单 job scrape recovery;M2c-4/M2d-2 只证明 production observability/NetworkPolicy profile 与 attestation 合同可校验;M2d-1 只证明本地两节点 kindnet 的隔离合成流量;M3-1 的 terminal evidence 与 M3-2 的 PolicyDecision/Approval 仍是 deterministic local fixtures。M3-2 只把 projection 写入本地 provider 并证明 retained target 的单次零写入 replay;M3-3 只把该本地 evidence 对应的 binding 写入临时 GDA Control 账本;M3-4 只向无认证 loopback receiver 发送精确 candidate 并验证 503 后幂等恢复;M3-5 只证明 OpenMetadata 在 provider 强制默认 role 之上的项目新增 grant 限定为 `table/Create`,以及本地 JWT 轮换/吊销和越权拒绝;M3-6 只证明隔离 Gravitino Basic IdP 的 bounded table-create、catalog-create 拒绝、登录轮换/吊销和完整清理;M3-7 只证明 pending production identity profile、profile-bound attestation 和派生 claim 的 fail-closed 合同可校验,没有部署或证明真实身份路径;M3-8 只证明同一 Docker Desktop 集群内 Basic 用户、JDBC metadata 与 file warehouse PVC 在受控 Pod restart 后连续;M3-9 只证明同节点共享 RWO PVC 的 Spark interoperability;M3-10 移除了该共享 PVC,并证明同一 Docker Desktop 主机/集群内 Spark 与 MinIO 的跨节点 S3-compatible 互操作,但不证明生产云对象存储、独立 failure domain、持久 identity binding、Flink 或完整 engine conformance;M3-11 只冻结 provider-neutral production object-store profile、精确 attestation binding 与 fail-closed claims,没有选择 provider、部署 bucket/KMS/policy 或提交真实 attestation;M3-12 只证明同一本地路径的 pre-forward commit failure 不改变可见 table state,随后一次显式重试产生一个新 snapshot/row,且无孤儿 data file;M3-13 只证明单次本地 append 在 provider 200 响应丢失并映射为 commit-state-unknown 后,可以由即时 table readback 判定 committed 且不重提,不覆盖持久 controller、进程崩溃、并发写或任意 mutation;M3-14 只证明 ResourceVersion 注册与 Active Metadata 事件在本地 PostgreSQL 同事务创建,并验证租户/workload scoped claim/retry/complete,不包含常驻 consumer、DolphinScheduler submission、provider policy 或 provider mutation。M3-2 ingestion 仍使用 bootstrap admin,生产持久 binding、ResourceVersion 和 legacy authority 都未写入;生产对象存储、双 provider/生产最小权限、protected workload identity、OIDC、TLS、生产持久 catalog、tenant isolation、真实 receiver/alert/SLO、受保护 provider policy、生产故障注入、source-loss recovery、cancel/reconcile/lineage、完整 Spark/Flink conformance、生产 ingest、四项 production gate 和 `production_ready` 仍为 `false`。 diff --git a/docs/system-of-record-matrix-2026-07-24.md b/docs/system-of-record-matrix-2026-07-24.md index 96002392..42b0c4aa 100644 --- a/docs/system-of-record-matrix-2026-07-24.md +++ b/docs/system-of-record-matrix-2026-07-24.md @@ -94,7 +94,7 @@ - Metadata Fabric M3-11 已建立 production object-store profile/attestation gate;checked-in profile fingerprint 为 `668e194b3c688307014148391e7f389c9d6e9ca69c95d7b4cc92b4acae93181a`,report fingerprint 为 `85362dd10b7dc565f9fa567673d90b774cdec714bd1e70fb2c3c83c1af48b5ea`,`profile_valid=true`,43 项 provider/identity/transport/encryption/durability/consistency/tenancy/operations 外部输入以 blockers 暴露,`ready_for_protected_verification=false`、`attestation_valid=false`、`production_object_store_gate_passed=false`、`production_ready=false`。该合同绑定 M3-10 evidence,但没有选择或部署 provider;合成完整 attestation 只验证门禁逻辑,不计入生产证据。下一项真实证据是经 owner 批准并物化的 provider profile,以及来自 `production-object-store` 受保护环境、绑定当前 source/profile 并通过全部 26 项检查的 attestation。 - Metadata Fabric M3-12 已在本地 Spark driver loopback Iceberg REST proxy 中于 provider 转发前注入 HTTP 503。baseline 为 1 个 append snapshot、2 行和 1 个 referenced Parquet;失败写经 2 次 503 后 snapshot/row/file 零漂移;对同一逻辑行一次显式重试后为父子相连的 2 个 append snapshots、3 行和 2 个 referenced Parquet。直接 MinIO inventory 精确为 2 data + 3 metadata + 4 manifest = 9 objects,没有孤儿 data file;namespace、两块 PV 和 port-forward 均清理。contract fingerprint 为 `6d8944ab80246dc65891aa81118cb8b73f7ecad699be9a2af5e62d8260c41002`,evidence fingerprint 为 `39571cdac1e4043bcfc2d03a73b2b12ff925210daf8ae36bc640b8cb14d89401`。这只证明已知 pre-forward failure 的本地原子性与一次显式重试,不证明 provider uncertain outcome reconcile、网络 exactly-once、生产对象存储、cancel/lineage、Flink 或完整 Spark conformance。 - Metadata Fabric M3-13 已在本地 Spark driver loopback proxy 将 armed commit 转发给 Gravitino,并在 provider 200 后丢弃成功响应、返回 HTTP 504;Iceberg 将其映射为 commit-state-unknown,一次传输重试被抑制。Spark 只读 readback 后输出 `committed_do_not_resubmit` 和 `write_resubmitted=false`,最终为父子相连的 2 个 append snapshots、3 行、2 个 referenced Parquet;MinIO 为 2 data + 3 metadata + 4 manifest = 9 objects。Job `Complete 1/1`,namespace、两块 PV 和 port-forward 均清理。contract fingerprint 为 `7a8d75a1d6b4558b982c6c3242d8d356c5046955f8aae7a45e5c297b6f4d4132`,evidence fingerprint 为 `d6462fff78d07047311b1f715d5f2c7f08c0ce8fbdd5c8b26a3d95ddc3474786`。这只证明一个本地 append 的确定性 readback/no-resubmit;持久 controller、crash/concurrency、网络 exactly-once、生产对象存储、cancel/lineage、Flink 和完整 Spark conformance 仍未证明。 -- Metadata Fabric M3-14 已新增内容绑定 `MetadataChangeEvent`、migration 099 transactional outbox 与 PlatformGateway 原子注册/claim/fail/complete。真实 PostgreSQL 16 演练验证首次注册 `created=true`、pending/processed 精确 replay 均 `created=false`、错误 consumer/worker 拒绝、retry、lease-expiry reclaim、第三次 attempt 完成、processed 不再认领、旧 ResourceVersion 补事件整笔回滚、跨租户不可见、FORCE RLS 和 gateway 无直接 UPDATE/DELETE;最终只有 1 条权威事件。contract fingerprint 为 `3429cd1d7fc5015dab7dfd27b3972c4628238bc18d13c76ad28bb16697898e75`,evidence fingerprint 为 `d85a4575a6103e2f7107f8e11153c080a430e82f6cf6db547295f14ef909e96a`。激活意图只路由到 `metadata_fabric.projection_plan`,`provider_apply_authorized=false`、`provider_mutations_executed=false`、`production_ingestion_verified=false`、`production_scheduler_submission_verified=false`、`production_ready=false`。 +- Metadata Fabric M3-14 已新增内容绑定 `MetadataChangeEvent`、migration 099 transactional outbox 与 PlatformGateway 原子注册/claim/fail/complete。真实 PostgreSQL 16 演练验证首次注册 `created=true`、pending/processed 精确 replay 均 `created=false`、错误 consumer/worker 拒绝、retry、lease-expiry reclaim、第三次 attempt 完成、processed 不再认领、旧 ResourceVersion 补事件整笔回滚、跨租户不可见、FORCE RLS 和 gateway 无直接 UPDATE/DELETE;最终只有 1 条权威事件。contract fingerprint 为 `b4b9d77d716f0e5d62389378aec60bdc366b2d847ace07a6b476310eb3ff6732`,evidence fingerprint 为 `fcd44cf9043e144016644b0d04e72f348a777668d853081d9d21ded341bb2b59`。激活意图只路由到 `metadata_fabric.projection_plan`,`provider_apply_authorized=false`、`provider_mutations_executed=false`、`production_ingestion_verified=false`、`production_scheduler_submission_verified=false`、`production_ready=false`。 ## 下一验收证据 From 41eb17003a18a95cf2014bb2a8d36091a36f3560 Mon Sep 17 00:00:00 2001 From: Ning Zhou Date: Thu, 30 Jul 2026 11:48:20 +0800 Subject: [PATCH 3/3] test(platform): advance migration catalog expectation --- data_agent/test_platform_contracts.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/data_agent/test_platform_contracts.py b/data_agent/test_platform_contracts.py index 495fbd9e..c60f564f 100644 --- a/data_agent/test_platform_contracts.py +++ b/data_agent/test_platform_contracts.py @@ -472,7 +472,7 @@ def test_control_ledger_contract_and_migration_catalog_are_valid(): assert report["contract_count"] == 16 assert report["migration"]["sha256"] == migration["checksum"] assert migrations[-1]["migration_id"] == ( - "098_metadata_fabric_openlineage_delivery" + "099_active_metadata_change_outbox" )