diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 27309d68..0dd3cfd2 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -125,6 +125,9 @@ jobs: - name: Validate metadata fabric binding ledger evidence run: python -m data_agent.metadata_fabric_binding_ledger validate + - name: Validate metadata fabric OpenLineage delivery evidence + run: python -m data_agent.metadata_fabric_lineage_delivery validate + - name: Validate DolphinScheduler adapter boundary run: python -m data_agent.dolphinscheduler_adapter validate @@ -146,6 +149,11 @@ jobs: DATABASE_URL: postgresql://postgres:postgres@localhost:5432/gis_agent_test run: python -m pytest data_agent/test_metadata_fabric_binding_ledger_postgres.py -q + - name: Verify metadata fabric OpenLineage delivery on PostgreSQL + env: + 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: Run required platform tests env: DATABASE_URL: postgresql://postgres:postgres@localhost:5432/gis_agent_test @@ -174,6 +182,7 @@ 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_metadata_fabric_lineage_delivery.py \ data_agent/test_metadata_fabric_otel_failure_rehearsal.py \ data_agent/test_metadata_fabric_otel_metrics.py \ data_agent/test_metadata_fabric_provider_metrics.py \ diff --git a/data_agent/metadata_fabric_lineage_delivery.py b/data_agent/metadata_fabric_lineage_delivery.py new file mode 100644 index 00000000..0d00fe18 --- /dev/null +++ b/data_agent/metadata_fabric_lineage_delivery.py @@ -0,0 +1,613 @@ +"""Rehearse M3-4 OpenLineage delivery over a real local HTTP boundary.""" + +from __future__ import annotations + +import argparse +import hashlib +import json +import threading +from datetime import timedelta +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer +from pathlib import Path +from typing import Any + +from pydantic import BaseModel, ConfigDict +from sqlalchemy import create_engine, text + +from . import metadata_fabric_ingestion as ingestion +from . import metadata_fabric_binding_ledger as binding_ledger +from .metadata_fabric_binding_contract import ( + parse_metadata_fabric_execution_plan_artifact, +) +from .metadata_fabric_lineage_delivery_contract import ( + DEFAULT_TARGET_NAME, + MetadataFabricLineageDelivery, + build_metadata_fabric_lineage_delivery, +) +from .metadata_fabric_lineage_emitter import ( + LineageEmitterProfile, + MetadataFabricLineageConsumer, + OpenLineageHttpEmitter, +) +from .platform_contracts import canonical_json_bytes, canonical_json_fingerprint +from .platform_gateway import ( + DefinitionRegistration, + GatewayNotFoundError, + PlatformGateway, +) + +CONTRACT_SCHEMA = "gda.metadata_fabric_openlineage_delivery_contract.v1" +EVIDENCE_SCHEMA = "gda.metadata_fabric_openlineage_delivery_evidence.v1" +EMITTER_ACTOR = "workload:gda-lineage-emitter-local" +WORKER_ID = "worker:gda-lineage-emitter-local-1" + +REPO_ROOT = Path(__file__).resolve().parent.parent +DEFAULT_SOURCE_EVIDENCE = ( + REPO_ROOT / "docs/evidence/metadata-fabric-binding-ledger-2026-07-28.json" +) +DEFAULT_EVIDENCE_PATH = ( + REPO_ROOT + / "docs/evidence/metadata-fabric-openlineage-delivery-2026-07-28.json" +) +DEFAULT_WRAPPER_PATH = ( + REPO_ROOT / "scripts/metadata-fabric-openlineage-delivery.sh" +) +MIGRATIONS = tuple( + Path(__file__).resolve().parent / "migrations" / filename + for filename in ( + "092_platform_control_ledger.sql", + "093_app_user_tenant_context.sql", + "094_platform_control_gateway.sql", + "095_platform_command_outbox.sql", + "096_platform_success_verdict.sql", + "097_metadata_fabric_binding_ledger.sql", + "098_metadata_fabric_openlineage_delivery.sql", + ) +) + + +class MetadataFabricLineageDeliveryError(RuntimeError): + """The bounded local OpenLineage delivery rehearsal failed closed.""" + + +class _FrozenModel(BaseModel): + model_config = ConfigDict(extra="forbid", frozen=True, arbitrary_types_allowed=True) + + +class LineageDeliveryBundle(_FrozenModel): + binding_bundle: binding_ledger.BindingLedgerBundle + source_plan: ingestion.MetadataFabricIngestionPlan + delivery: MetadataFabricLineageDelivery + source_evidence_sha256: str + + +def _load_json_object(path: Path) -> dict[str, Any]: + payload = json.loads(path.read_text(encoding="utf-8")) + if not isinstance(payload, dict): + raise MetadataFabricLineageDeliveryError( + f"{path.name} must contain an object" + ) + return payload + + +def _build_source_plan() -> ingestion.MetadataFabricIngestionPlan: + values = ingestion._load_contract_inputs( + ingestion.DEFAULT_PLATFORM_FIXTURE, + ingestion.DEFAULT_METADATA_FIXTURE, + ) + return ingestion.build_ingestion_plan( + metadata_resource=values[2], + target=values[3], + binding=values[4], + definition=values[5], + run=values[6], + source=values[7], + artifact=values[8], + quality=values[9], + lineage=values[10], + success=values[11], + openmetadata=values[12], + gravitino=values[13], + ) + + +def build_lineage_delivery_bundle( + source_evidence_path: Path = DEFAULT_SOURCE_EVIDENCE, +) -> LineageDeliveryBundle: + source_evidence = _load_json_object(source_evidence_path.resolve()) + errors = binding_ledger.validate_rehearsal_evidence(source_evidence) + if errors: + raise MetadataFabricLineageDeliveryError( + "M3-3 binding evidence is invalid: " + ", ".join(errors) + ) + binding_bundle = binding_ledger.build_binding_ledger_bundle() + if ( + source_evidence.get("binding_id") + != str(binding_bundle.record.binding_id) + or source_evidence.get("record_sha256") + != binding_bundle.record.record_sha256 + ): + raise MetadataFabricLineageDeliveryError( + "M3-3 evidence does not match the deterministic binding" + ) + source_plan = _build_source_plan() + apply_plan = parse_metadata_fabric_execution_plan_artifact( + binding_bundle.artifacts[0] + ) + delivery = build_metadata_fabric_lineage_delivery( + binding=binding_bundle.record, + source_plan=source_plan, + apply_plan=apply_plan, + actor_subject=EMITTER_ACTOR, + created_at=binding_bundle.record.recorded_at + timedelta(seconds=1), + target_name=DEFAULT_TARGET_NAME, + ) + return LineageDeliveryBundle( + binding_bundle=binding_bundle, + source_plan=source_plan, + delivery=delivery, + source_evidence_sha256=source_evidence["evidence_sha256"], + ) + + +def build_contract_report( + *, + source_evidence_path: Path = DEFAULT_SOURCE_EVIDENCE, + wrapper_path: Path = DEFAULT_WRAPPER_PATH, +) -> dict[str, Any]: + errors: list[str] = [] + bundle: LineageDeliveryBundle | None = None + try: + bundle = build_lineage_delivery_bundle(source_evidence_path) + except (KeyError, OSError, TypeError, ValueError) as exc: + errors.append(f"lineage delivery contract is invalid: {type(exc).__name__}") + try: + wrapper = wrapper_path.read_text(encoding="utf-8") + for marker in ( + "set -euo pipefail", + "metadata_fabric_lineage_delivery", + "docker run", + ): + if marker not in wrapper: + errors.append(f"lineage wrapper is missing marker: {marker}") + except OSError as exc: + errors.append(f"lineage wrapper is invalid: {type(exc).__name__}") + files: dict[str, dict[str, str | None]] = {} + for path in ( + Path(__file__).resolve(), + source_evidence_path.resolve(), + wrapper_path.resolve(), + *MIGRATIONS, + ): + files[path.name] = { + "path": ( + path.resolve().relative_to(REPO_ROOT).as_posix() + if path.resolve().is_relative_to(REPO_ROOT) + else path.resolve().as_posix() + ), + "sha256": ( + hashlib.sha256(path.read_bytes()).hexdigest() + if path.is_file() + else None + ), + } + stable = { + "schema": CONTRACT_SCHEMA, + "source_evidence_sha256": ( + bundle.source_evidence_sha256 if bundle is not None else None + ), + "delivery_id": ( + str(bundle.delivery.delivery_id) if bundle is not None else None + ), + "event_sha256": ( + bundle.delivery.event_sha256 if bundle is not None else None + ), + "idempotency_key": ( + bundle.delivery.idempotency_key if bundle is not None else None + ), + "files": files, + "errors": errors, + } + return { + **stable, + "status": "valid" if not errors else "invalid", + "contract_sha256": canonical_json_fingerprint(stable), + "local_wire_openlineage_delivery_verified": False, + "production_ready": False, + } + + +class _SinkState: + def __init__(self, expected: MetadataFabricLineageDelivery) -> None: + self.expected = expected + self.requests: list[dict[str, Any]] = [] + self.accepted: dict[str, str] = {} + self.errors: list[str] = [] + + +def _sink_handler(state: _SinkState) -> type[BaseHTTPRequestHandler]: + class Handler(BaseHTTPRequestHandler): + def log_message(self, *args: object) -> None: + return + + def do_POST(self) -> None: + length = int(self.headers.get("Content-Length", "0")) + body = self.rfile.read(length) + key = self.headers.get("Idempotency-Key", "") + delivery_id = self.headers.get("X-GDA-Delivery-ID", "") + event_sha = self.headers.get("X-GDA-Event-SHA256", "") + body_sha = hashlib.sha256(body).hexdigest() + expected_body = canonical_json_bytes( + state.expected.event.model_dump(mode="json", by_alias=True) + ) + request_errors: list[str] = [] + if self.path != "/api/v1/lineage": + request_errors.append("path_mismatch") + if self.headers.get("Content-Type") != "application/json": + request_errors.append("content_type_mismatch") + if key != state.expected.idempotency_key: + request_errors.append("idempotency_key_mismatch") + if delivery_id != str(state.expected.delivery_id): + request_errors.append("delivery_id_mismatch") + if event_sha != state.expected.event_sha256: + request_errors.append("event_sha_mismatch") + if body != expected_body: + request_errors.append("body_mismatch") + state.errors.extend(request_errors) + duplicate = key in state.accepted + if duplicate and state.accepted[key] != body_sha: + request_errors.append("idempotency_content_conflict") + state.errors.append("idempotency_content_conflict") + if not request_errors and not duplicate: + state.accepted[key] = body_sha + state.requests.append( + { + "body_sha256": body_sha, + "idempotency_key": key, + "delivery_id": delivery_id, + "event_sha256": event_sha, + "duplicate": duplicate, + } + ) + if request_errors: + status = 409 + response = {"accepted": False} + elif duplicate: + status = 200 + response = {"accepted": True, "duplicate": True} + else: + # Simulate receiver commit followed by a lost/failed acknowledgement. + status = 503 + response = {"accepted": True, "acknowledged": False} + payload = canonical_json_bytes(response) + self.send_response(status) + self.send_header("Content-Type", "application/json") + self.send_header("Content-Length", str(len(payload))) + self.end_headers() + self.wfile.write(payload) + + return Handler + + +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 MetadataFabricLineageDeliveryError( + "local lineage 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 _register_binding(gateway: PlatformGateway, bundle: LineageDeliveryBundle) -> None: + binding = bundle.binding_bundle + by_urn = {item.resource_urn: item for item in binding.resources} + definition_version = next( + item + for item in binding.resource_versions + if item.resource_version_id == binding.definition.definition_version_id + ) + gateway.register_definition( + DefinitionRegistration( + resource=by_urn[definition_version.resource_urn], + resource_version=definition_version, + definition=binding.definition, + ) + ) + for resource in binding.resources: + if resource.resource_urn != definition_version.resource_urn: + gateway.register_resource(resource) + for version in binding.resource_versions: + if version.resource_version_id != definition_version.resource_version_id: + gateway.register_resource_version(version) + for artifact in binding.artifacts: + gateway.record_artifact(artifact) + gateway.commit_metadata_fabric_binding(binding.record) + + +def run_local_rehearsal( + database_url: str, + *, + source_evidence_path: Path = DEFAULT_SOURCE_EVIDENCE, +) -> dict[str, Any]: + if not database_url.startswith( + ("postgresql://", "postgresql+psycopg://", "postgresql+psycopg2://") + ): + raise MetadataFabricLineageDeliveryError( + "lineage rehearsal requires a PostgreSQL database URL" + ) + bundle = build_lineage_delivery_bundle(source_evidence_path) + engine = create_engine(database_url) + server: ThreadingHTTPServer | None = None + thread: threading.Thread | None = None + try: + _apply_migrations(engine) + gateway = PlatformGateway(engine) + _register_binding(gateway, bundle) + first_enqueue = gateway.enqueue_metadata_fabric_lineage( + bundle.delivery, + source_plan=bundle.source_plan, + ) + replay_enqueue = gateway.enqueue_metadata_fabric_lineage( + bundle.delivery, + source_plan=bundle.source_plan, + ) + + state = _SinkState(bundle.delivery) + server = ThreadingHTTPServer(("127.0.0.1", 0), _sink_handler(state)) + thread = threading.Thread(target=server.serve_forever, daemon=True) + thread.start() + endpoint = f"http://127.0.0.1:{server.server_port}/api/v1/lineage" + profile = LineageEmitterProfile( + target_name=bundle.delivery.target_name, + endpoint_url=endpoint, + actor_subject=bundle.delivery.actor_subject, + ) + with OpenLineageHttpEmitter(profile) as emitter: + consumer = MetadataFabricLineageConsumer( + emitter, + gateway=gateway, + retry_delay_seconds=0, + ) + first_batch = consumer.run_once( + bundle.delivery.tenant_id, + worker_id=WORKER_ID, + ) + after_first = gateway.get_metadata_fabric_lineage_delivery( + bundle.delivery.tenant_id, + bundle.delivery.delivery_id, + ) + second_batch = consumer.run_once( + bundle.delivery.tenant_id, + worker_id=WORKER_ID, + ) + third_batch = consumer.run_once( + bundle.delivery.tenant_id, + worker_id=WORKER_ID, + ) + stored = gateway.get_metadata_fabric_lineage_delivery( + bundle.delivery.tenant_id, + bundle.delivery.delivery_id, + ) + cross_tenant_visible = True + try: + gateway.get_metadata_fabric_lineage_delivery( + "ar0-golden-isolated", + bundle.delivery.delivery_id, + ) + except GatewayNotFoundError: + cross_tenant_visible = False + with engine.connect() as connection: + restricted_mutation = connection.exec_driver_sql( + """ + SELECT NOT has_table_privilege( + 'gda_control_gateway', + 'gda_control.metadata_fabric_lineage_outbox', 'UPDATE' + ) + AND NOT has_table_privilege( + 'gda_control_gateway', + 'gda_control.metadata_fabric_lineage_outbox', 'DELETE' + ) + """ + ).scalar_one() + force_rls = connection.exec_driver_sql( + """ + SELECT relforcerowsecurity + FROM pg_class + WHERE oid = 'gda_control.metadata_fabric_lineage_outbox'::regclass + """ + ).scalar_one() + connection.rollback() + finally: + if server is not None: + server.shutdown() + server.server_close() + if thread is not None: + thread.join(timeout=5) + engine.dispose() + + expected_body_sha = hashlib.sha256( + canonical_json_bytes( + bundle.delivery.event.model_dump(mode="json", by_alias=True) + ) + ).hexdigest() + request_count = len(state.requests) + duplicate_count = sum(item["duplicate"] for item in state.requests) + verified = ( + first_enqueue.created + and not replay_enqueue.created + and first_batch.claimed == 1 + and first_batch.retry_pending == 1 + and after_first.status.value == "pending" + and after_first.attempt_count == 1 + and after_first.last_error_code == "http_5xx" + and after_first.response_status == 503 + and second_batch.delivered == 1 + and third_batch.claimed == 0 + and stored.status.value == "delivered" + and stored.attempt_count == 2 + and stored.receipt_sha256 is not None + and request_count == 2 + and len(state.accepted) == 1 + and duplicate_count == 1 + and all(item["body_sha256"] == expected_body_sha for item in state.requests) + and not state.errors + and not cross_tenant_visible + and restricted_mutation + and force_rls + ) + stable = { + "schema": EVIDENCE_SCHEMA, + "status": ( + "local_wire_openlineage_delivery_verified" if verified else "blocked" + ), + "source_evidence_sha256": bundle.source_evidence_sha256, + "binding_id": str(bundle.delivery.binding_id), + "delivery_id": str(bundle.delivery.delivery_id), + "source_plan_sha256": bundle.delivery.source_plan_sha256, + "event_sha256": bundle.delivery.event_sha256, + "idempotency_key": bundle.delivery.idempotency_key, + "target_name": bundle.delivery.target_name, + "first_enqueue_created": first_enqueue.created, + "replay_enqueue_created": replay_enqueue.created, + "first_attempt_response_status": after_first.response_status, + "first_attempt_retry_pending": first_batch.retry_pending == 1, + "final_response_status": stored.response_status, + "final_attempt_count": stored.attempt_count, + "receipt_sha256": stored.receipt_sha256, + "wire_request_count": request_count, + "receiver_unique_accept_count": len(state.accepted), + "receiver_duplicate_count": duplicate_count, + "wire_body_sha256": expected_body_sha, + "wire_body_matched": all( + item["body_sha256"] == expected_body_sha for item in state.requests + ), + "stable_idempotency_header_verified": all( + item["idempotency_key"] == bundle.delivery.idempotency_key + for item in state.requests + ), + "receiver_commit_then_failed_ack_recovered": ( + first_batch.retry_pending == 1 + and second_batch.delivered == 1 + and duplicate_count == 1 + ), + "completed_delivery_not_reclaimed": third_batch.claimed == 0, + "cross_tenant_read_blocked": not cross_tenant_visible, + "gateway_direct_update_delete_blocked": bool(restricted_mutation), + "force_rls_verified": bool(force_rls), + "transport_semantics": "at_least_once_with_receiver_idempotency", + "local_loopback_receiver": True, + "receiver_credentials_used": False, + "local_wire_openlineage_delivery_verified": verified, + "live_openlineage_emission_verified": verified, + "provider_mutations_executed": False, + "writes_to_legacy": False, + "provider_minimum_privilege_verified": False, + "oidc_verified": False, + "durable_catalog_verified": False, + "production_receiver_verified": False, + "production_ingestion_verified": False, + "production_ready": False, + "errors": [] if verified else ["local OpenLineage delivery 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("lineage delivery evidence schema does not match") + if evidence.get("evidence_sha256") != canonical_json_fingerprint(stable): + errors.append("lineage delivery evidence SHA-256 does not match") + for claim in ( + "provider_minimum_privilege_verified", + "oidc_verified", + "durable_catalog_verified", + "production_receiver_verified", + "production_ingestion_verified", + "production_ready", + ): + if evidence.get(claim) is not False: + errors.append(f"local lineage evidence may not claim {claim}") + for claim in ( + "local_wire_openlineage_delivery_verified", + "live_openlineage_emission_verified", + "wire_body_matched", + "stable_idempotency_header_verified", + "receiver_commit_then_failed_ack_recovered", + "completed_delivery_not_reclaimed", + "cross_tenant_read_blocked", + "gateway_direct_update_delete_blocked", + "force_rls_verified", + ): + if evidence.get(claim) is not True: + errors.append(f"local lineage evidence did not verify {claim}") + if evidence.get("transport_semantics") != ( + "at_least_once_with_receiver_idempotency" + ): + errors.append("lineage transport semantics are invalid") + if evidence.get("receiver_credentials_used") is not False: + errors.append("local lineage receiver must not use credentials") + if evidence.get("provider_mutations_executed") is not False: + errors.append("M3-4 must not mutate metadata providers") + if evidence.get("writes_to_legacy") is not False: + errors.append("M3-4 must not write legacy tables") + 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( + "--source-evidence", type=Path, default=DEFAULT_SOURCE_EVIDENCE + ) + 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( + "--source-evidence", type=Path, default=DEFAULT_SOURCE_EVIDENCE + ) + rehearse.add_argument("--evidence-out", type=Path, required=True) + args = parser.parse_args(argv) + + if args.command == "validate": + report = build_contract_report(source_evidence_path=args.source_evidence) + try: + evidence = _load_json_object(args.evidence) + report["errors"].extend(validate_rehearsal_evidence(evidence)) + except (OSError, ValueError) as exc: + report["errors"].append( + f"lineage delivery evidence is invalid: {type(exc).__name__}" + ) + report["status"] = "valid" if not report["errors"] else "invalid" + report["local_wire_openlineage_delivery_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, + source_evidence_path=args.source_evidence, + ) + 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/metadata_fabric_lineage_delivery_contract.py b/data_agent/metadata_fabric_lineage_delivery_contract.py new file mode 100644 index 00000000..c1b300c1 --- /dev/null +++ b/data_agent/metadata_fabric_lineage_delivery_contract.py @@ -0,0 +1,322 @@ +"""Content-bound outbox contract for Metadata Fabric OpenLineage delivery.""" + +from __future__ import annotations + +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 .metadata_fabric_binding_contract import ( + MetadataFabricApplyPlan, + MetadataFabricBindingRecord, +) +from .metadata_fabric_ingestion import ( + MetadataFabricIngestionPlan, + OpenLineageRunEvent, +) +from .platform_contracts import ( + Sha256, + TenantId, + canonical_json_fingerprint, +) + +DELIVERY_SCHEMA = "gda.metadata_fabric_openlineage_delivery.v1" +DEFAULT_TARGET_NAME = "local-openlineage-http-sink" + +NonEmptyText = Annotated[ + str, + StringConstraints(strip_whitespace=True, min_length=1, max_length=512), +] +FailureCode = Annotated[ + str, + StringConstraints( + strip_whitespace=True, + pattern=r"^[a-z0-9_]{1,64}$", + ), +] + + +class LineageDeliveryStatus(str, Enum): + PENDING = "pending" + IN_FLIGHT = "in_flight" + DELIVERED = "delivered" + FAILED = "failed" + + +class _FrozenModel(BaseModel): + model_config = ConfigDict(extra="forbid", frozen=True) + + +def _aware_utc(value: datetime) -> datetime: + if value.tzinfo is None or value.utcoffset() is None: + raise ValueError("delivery timestamps must include a timezone") + return value.astimezone(UTC) + + +def openlineage_event_sha256(event: OpenLineageRunEvent) -> str: + return canonical_json_fingerprint( + event.model_dump(mode="json", by_alias=True) + ) + + +def openlineage_delivery_id( + binding_id: UUID, + *, + target_name: str, + event_sha256: str, +) -> UUID: + return uuid5( + binding_id, + f"openlineage:{target_name}:{event_sha256}", + ) + + +def openlineage_idempotency_key( + *, + tenant_id: str, + binding_id: UUID, + target_name: str, + event_sha256: str, +) -> str: + return canonical_json_fingerprint( + { + "tenant_id": tenant_id, + "binding_id": str(binding_id), + "target_name": target_name, + "event_sha256": event_sha256, + } + ) + + +def openlineage_receipt_sha256( + delivery: "MetadataFabricLineageDelivery", + *, + response_status: int, + response_body_sha256: str, +) -> str: + return canonical_json_fingerprint( + { + "tenant_id": delivery.tenant_id, + "delivery_id": str(delivery.delivery_id), + "binding_id": str(delivery.binding_id), + "target_name": delivery.target_name, + "event_sha256": delivery.event_sha256, + "idempotency_key": delivery.idempotency_key, + "response_status": response_status, + "response_body_sha256": response_body_sha256, + } + ) + + +class MetadataFabricLineageDelivery(_FrozenModel): + delivery_schema: Literal[ + "gda.metadata_fabric_openlineage_delivery.v1" + ] = Field(default=DELIVERY_SCHEMA, alias="schema") + tenant_id: TenantId + delivery_id: UUID + binding_id: UUID + resource_version_id: UUID + run_id: UUID + source_plan_sha256: Sha256 + target_name: NonEmptyText + event: OpenLineageRunEvent + event_sha256: Sha256 + idempotency_key: Sha256 + actor_subject: NonEmptyText + status: LineageDeliveryStatus = LineageDeliveryStatus.PENDING + attempt_count: Annotated[int, Field(ge=0)] = 0 + max_attempts: Annotated[int, Field(ge=1, le=20)] = 3 + available_at: datetime + claimed_by: NonEmptyText | None = None + claimed_until: datetime | None = None + last_error_code: FailureCode | None = None + response_status: Annotated[int, Field(ge=100, le=599)] | None = None + response_body_sha256: Sha256 | None = None + receipt_sha256: Sha256 | None = None + created_at: datetime + completed_at: datetime | None = None + + @field_validator( + "available_at", "claimed_until", "created_at", "completed_at" + ) + @classmethod + def _utc_timestamps(cls, value: datetime | None) -> datetime | None: + return _aware_utc(value) if value is not None else None + + @model_validator(mode="after") + def _content_bound_delivery(self) -> Self: + if not self.actor_subject.startswith("workload:"): + raise ValueError("lineage delivery actor must use workload identity") + if self.event.run.run_id != self.run_id: + raise ValueError("lineage delivery run does not match event") + expected_event_sha = openlineage_event_sha256(self.event) + if self.event_sha256 != expected_event_sha: + raise ValueError("lineage delivery event SHA-256 does not match") + expected_id = openlineage_delivery_id( + self.binding_id, + target_name=self.target_name, + event_sha256=self.event_sha256, + ) + if self.delivery_id != expected_id: + raise ValueError("lineage delivery UUID does not match content") + expected_key = openlineage_idempotency_key( + tenant_id=self.tenant_id, + binding_id=self.binding_id, + target_name=self.target_name, + event_sha256=self.event_sha256, + ) + if self.idempotency_key != expected_key: + raise ValueError("lineage delivery idempotency key does not match") + + 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("delivery claim owner and expiry must be set together") + if self.status == LineageDeliveryStatus.PENDING: + if claimed or self.completed_at is not None: + raise ValueError("pending delivery cannot be claimed or completed") + elif self.status == LineageDeliveryStatus.IN_FLIGHT: + if not claimed or self.completed_at is not None: + raise ValueError("in-flight delivery requires an active claim") + elif claimed or self.completed_at is None: + raise ValueError("terminal delivery must release its claim") + + if self.status == LineageDeliveryStatus.DELIVERED: + if ( + self.response_status is None + or not 200 <= self.response_status < 300 + or self.response_body_sha256 is None + or self.receipt_sha256 is None + or self.last_error_code is not None + ): + raise ValueError("delivered lineage state is incomplete") + expected_receipt = openlineage_receipt_sha256( + self, + response_status=self.response_status, + response_body_sha256=self.response_body_sha256, + ) + if self.receipt_sha256 != expected_receipt: + raise ValueError("lineage delivery receipt SHA-256 does not match") + elif self.receipt_sha256 is not None: + raise ValueError("only delivered lineage may have a receipt") + + if ( + self.response_body_sha256 is not None + and self.response_status is None + ): + raise ValueError("response body fingerprint requires HTTP status") + if self.status == LineageDeliveryStatus.FAILED: + if self.last_error_code is None: + raise ValueError("failed lineage delivery requires an error code") + return self + + +def validate_delivery_source( + *, + binding: MetadataFabricBindingRecord, + source_plan: MetadataFabricIngestionPlan, + apply_plan: MetadataFabricApplyPlan, +) -> None: + expected_identity = ( + binding.tenant_id, + binding.binding.resource_urn, + binding.binding.resource_version_id, + binding.binding.content_sha256, + ) + observed_identity = ( + source_plan.tenant_id, + source_plan.resource_urn, + source_plan.resource_version_id, + source_plan.content_sha256, + ) + if observed_identity != expected_identity: + raise ValueError("OpenLineage source plan does not match binding identity") + if apply_plan.source_plan_sha256 != source_plan.plan_sha256: + raise ValueError("authorized apply plan does not bind OpenLineage source plan") + if ( + apply_plan.tenant_id != binding.tenant_id + or apply_plan.resource_version_id + != binding.binding.resource_version_id + or apply_plan.run_id != source_plan.run_id + ): + raise ValueError("authorized apply plan does not match binding lineage") + + +def build_metadata_fabric_lineage_delivery( + *, + binding: MetadataFabricBindingRecord, + source_plan: MetadataFabricIngestionPlan, + apply_plan: MetadataFabricApplyPlan, + actor_subject: str, + created_at: datetime, + target_name: str = DEFAULT_TARGET_NAME, + max_attempts: int = 3, +) -> MetadataFabricLineageDelivery: + validate_delivery_source( + binding=binding, + source_plan=source_plan, + apply_plan=apply_plan, + ) + created_at = _aware_utc(created_at) + if created_at < binding.recorded_at: + raise ValueError("lineage delivery cannot predate the binding") + event = source_plan.openlineage_event + event_sha = openlineage_event_sha256(event) + delivery_id = openlineage_delivery_id( + binding.binding_id, + target_name=target_name, + event_sha256=event_sha, + ) + return MetadataFabricLineageDelivery( + tenant_id=binding.tenant_id, + delivery_id=delivery_id, + binding_id=binding.binding_id, + resource_version_id=binding.binding.resource_version_id, + run_id=source_plan.run_id, + source_plan_sha256=source_plan.plan_sha256, + target_name=target_name, + event=event, + event_sha256=event_sha, + idempotency_key=openlineage_idempotency_key( + tenant_id=binding.tenant_id, + binding_id=binding.binding_id, + target_name=target_name, + event_sha256=event_sha, + ), + actor_subject=actor_subject, + max_attempts=max_attempts, + available_at=created_at, + created_at=created_at, + ) + + +def delivery_binding_payload( + delivery: MetadataFabricLineageDelivery, +) -> dict[str, Any]: + """Return immutable identity fields, excluding mutable delivery state.""" + return delivery.model_dump( + mode="json", + by_alias=True, + exclude={ + "status", + "attempt_count", + "available_at", + "claimed_by", + "claimed_until", + "last_error_code", + "response_status", + "response_body_sha256", + "receipt_sha256", + "created_at", + "completed_at", + }, + ) diff --git a/data_agent/metadata_fabric_lineage_emitter.py b/data_agent/metadata_fabric_lineage_emitter.py new file mode 100644 index 00000000..6ae8f11d --- /dev/null +++ b/data_agent/metadata_fabric_lineage_emitter.py @@ -0,0 +1,242 @@ +"""Strict local HTTP emitter for claimed Metadata Fabric OpenLineage events.""" + +from __future__ import annotations + +import hashlib +from dataclasses import dataclass +from typing import Annotated, Self +from urllib.parse import urlsplit +from uuid import UUID + +import httpx +from pydantic import ( + BaseModel, + ConfigDict, + Field, + StringConstraints, + model_validator, +) + +from .metadata_fabric_lineage_delivery_contract import ( + LineageDeliveryStatus, + MetadataFabricLineageDelivery, +) +from .platform_contracts import canonical_json_bytes +from .platform_gateway import PlatformGateway + +NonEmptyText = Annotated[ + str, + StringConstraints(strip_whitespace=True, min_length=1, max_length=512), +] + + +class LineageEmitterProfile(BaseModel): + """Local-only endpoint profile; credentials and redirects are unsupported.""" + + model_config = ConfigDict(extra="forbid", frozen=True) + + target_name: NonEmptyText + endpoint_url: NonEmptyText + actor_subject: NonEmptyText + timeout_seconds: Annotated[float, Field(gt=0, le=30)] = 5.0 + + @model_validator(mode="after") + def _bounded_local_endpoint(self) -> Self: + parsed = urlsplit(self.endpoint_url) + if ( + parsed.scheme != "http" + or parsed.hostname not in {"127.0.0.1", "localhost"} + or parsed.path != "/api/v1/lineage" + or parsed.username is not None + or parsed.password is not None + or parsed.query + or parsed.fragment + ): + raise ValueError( + "local OpenLineage endpoint must be loopback /api/v1/lineage" + ) + if not self.actor_subject.startswith("workload:"): + raise ValueError("lineage emitter must use workload identity") + return self + + +class LineageHttpDeliveryError(RuntimeError): + def __init__( + self, + code: str, + *, + retryable: bool, + response_status: int | None = None, + ) -> None: + super().__init__(code) + self.code = code + self.retryable = retryable + self.response_status = response_status + + +@dataclass(frozen=True) +class LineageHttpReceipt: + response_status: int + response_body_sha256: str + + +@dataclass(frozen=True) +class LineageDeliveryBatchResult: + claimed: int + delivered: int + retry_pending: int + failed: int + delivery_ids: tuple[UUID, ...] + + +class OpenLineageHttpEmitter: + def __init__( + self, + profile: LineageEmitterProfile, + *, + transport: httpx.BaseTransport | None = None, + ) -> None: + self.profile = profile + self._client = httpx.Client( + timeout=profile.timeout_seconds, + follow_redirects=False, + transport=transport, + ) + + def close(self) -> None: + self._client.close() + + def __enter__(self) -> "OpenLineageHttpEmitter": + return self + + def __exit__(self, *args: object) -> None: + self.close() + + def emit( + self, delivery: MetadataFabricLineageDelivery + ) -> LineageHttpReceipt: + if ( + delivery.status != LineageDeliveryStatus.IN_FLIGHT + or delivery.claimed_by is None + ): + raise LineageHttpDeliveryError( + "delivery_not_claimed", + retryable=False, + ) + if ( + delivery.target_name != self.profile.target_name + or delivery.actor_subject != self.profile.actor_subject + ): + raise LineageHttpDeliveryError( + "emitter_profile_mismatch", + retryable=False, + ) + body = canonical_json_bytes( + delivery.event.model_dump(mode="json", by_alias=True) + ) + try: + response = self._client.post( + self.profile.endpoint_url, + content=body, + headers={ + "Accept": "application/json", + "Content-Type": "application/json", + "Idempotency-Key": delivery.idempotency_key, + "X-GDA-Delivery-ID": str(delivery.delivery_id), + "X-GDA-Event-SHA256": delivery.event_sha256, + }, + ) + except httpx.TransportError: + raise LineageHttpDeliveryError( + "transport_error", + retryable=True, + ) from None + status = response.status_code + if 200 <= status < 300: + return LineageHttpReceipt( + response_status=status, + response_body_sha256=hashlib.sha256(response.content).hexdigest(), + ) + if status == 429: + code = "http_429" + retryable = True + elif 500 <= status < 600: + code = "http_5xx" + retryable = True + else: + code = "http_4xx" + retryable = False + raise LineageHttpDeliveryError( + code, + retryable=retryable, + response_status=status, + ) + + +class MetadataFabricLineageConsumer: + """Deliver claimed events at least once; receiver idempotency handles replay.""" + + def __init__( + self, + emitter: OpenLineageHttpEmitter, + *, + gateway: PlatformGateway, + retry_delay_seconds: int = 5, + ) -> None: + if not 0 <= retry_delay_seconds <= 86400: + raise ValueError("lineage retry delay must be between 0 and 86400") + self.emitter = emitter + self.gateway = gateway + self.retry_delay_seconds = retry_delay_seconds + + def run_once( + self, + tenant_id: str, + *, + worker_id: str, + limit: int = 10, + lease_seconds: int = 60, + ) -> LineageDeliveryBatchResult: + deliveries = self.gateway.claim_metadata_fabric_lineage( + tenant_id, + worker_id, + actor_subject=self.emitter.profile.actor_subject, + limit=limit, + lease_seconds=lease_seconds, + ) + delivered = 0 + retry_pending = 0 + failed = 0 + for delivery in deliveries: + try: + receipt = self.emitter.emit(delivery) + except LineageHttpDeliveryError as exc: + outcome = self.gateway.fail_metadata_fabric_lineage( + delivery.tenant_id, + delivery.delivery_id, + worker_id=worker_id, + error_code=exc.code, + response_status=exc.response_status, + retryable=exc.retryable, + retry_delay_seconds=self.retry_delay_seconds, + ) + if outcome.status == LineageDeliveryStatus.FAILED: + failed += 1 + else: + retry_pending += 1 + continue + self.gateway.complete_metadata_fabric_lineage( + delivery.tenant_id, + delivery.delivery_id, + worker_id=worker_id, + response_status=receipt.response_status, + response_body_sha256=receipt.response_body_sha256, + ) + delivered += 1 + return LineageDeliveryBatchResult( + claimed=len(deliveries), + delivered=delivered, + retry_pending=retry_pending, + failed=failed, + delivery_ids=tuple(item.delivery_id for item in deliveries), + ) diff --git a/data_agent/migrations/098_metadata_fabric_openlineage_delivery.sql b/data_agent/migrations/098_metadata_fabric_openlineage_delivery.sql new file mode 100644 index 00000000..31cee7eb --- /dev/null +++ b/data_agent/migrations/098_metadata_fabric_openlineage_delivery.sql @@ -0,0 +1,354 @@ +-- 098: Tenant-scoped OpenLineage HTTP delivery outbox. +-- +-- This table owns at-least-once delivery state only. The immutable binding, +-- PlatformRun and receiver remain the authorities for their own domains. + +CREATE TABLE IF NOT EXISTS gda_control.metadata_fabric_lineage_outbox ( + tenant_id TEXT NOT NULL, + delivery_id UUID PRIMARY KEY, + binding_id UUID NOT NULL, + resource_version_id UUID NOT NULL, + run_id UUID NOT NULL, + source_plan_sha256 CHAR(64) NOT NULL, + target_name TEXT NOT NULL, + event JSONB NOT NULL, + event_sha256 CHAR(64) NOT NULL, + idempotency_key CHAR(64) NOT NULL, + actor_subject TEXT NOT NULL, + status TEXT NOT NULL DEFAULT 'pending', + attempt_count INTEGER NOT NULL DEFAULT 0, + max_attempts INTEGER NOT NULL DEFAULT 3, + available_at TIMESTAMPTZ NOT NULL, + claimed_by TEXT, + claimed_until TIMESTAMPTZ, + last_error_code TEXT, + response_status INTEGER, + response_body_sha256 CHAR(64), + receipt_sha256 CHAR(64), + created_at TIMESTAMPTZ NOT NULL, + completed_at TIMESTAMPTZ, + CONSTRAINT uq_gda_lineage_delivery_tenant_id + UNIQUE (tenant_id, delivery_id), + CONSTRAINT uq_gda_lineage_delivery_event + UNIQUE (tenant_id, binding_id, target_name, event_sha256), + CONSTRAINT uq_gda_lineage_delivery_idempotency + UNIQUE (tenant_id, idempotency_key), + CONSTRAINT fk_gda_lineage_delivery_binding + FOREIGN KEY (tenant_id, binding_id) + REFERENCES gda_control.metadata_fabric_binding(tenant_id, binding_id), + CONSTRAINT ck_gda_lineage_delivery_document CHECK ( + jsonb_typeof(event) = 'object' + AND event->>'schemaURL' + = 'https://openlineage.io/spec/2-0-2/OpenLineage.json#/definitions/RunEvent' + AND event->>'eventType' = 'COMPLETE' + AND event->'run'->>'runId' = run_id::text + ), + CONSTRAINT ck_gda_lineage_delivery_sha256 CHECK ( + source_plan_sha256 ~ '^[0-9a-f]{64}$' + AND event_sha256 ~ '^[0-9a-f]{64}$' + AND idempotency_key ~ '^[0-9a-f]{64}$' + AND ( + response_body_sha256 IS NULL + OR response_body_sha256 ~ '^[0-9a-f]{64}$' + ) + AND ( + receipt_sha256 IS NULL + OR receipt_sha256 ~ '^[0-9a-f]{64}$' + ) + ), + CONSTRAINT ck_gda_lineage_delivery_actor CHECK ( + actor_subject ~ '^workload:.+' + ), + CONSTRAINT ck_gda_lineage_delivery_status CHECK ( + status IN ('pending', 'in_flight', 'delivered', 'failed') + ), + CONSTRAINT ck_gda_lineage_delivery_attempts CHECK ( + attempt_count >= 0 AND max_attempts BETWEEN 1 AND 20 + ), + CONSTRAINT ck_gda_lineage_delivery_claim CHECK ( + (claimed_by IS NULL) = (claimed_until IS NULL) + ), + CONSTRAINT ck_gda_lineage_delivery_error CHECK ( + last_error_code IS NULL + OR last_error_code ~ '^[a-z0-9_]{1,64}$' + ), + CONSTRAINT ck_gda_lineage_delivery_response CHECK ( + (response_status IS NULL OR response_status BETWEEN 100 AND 599) + AND (response_body_sha256 IS NULL OR response_status IS NOT NULL) + ), + CONSTRAINT ck_gda_lineage_delivery_state CHECK ( + ( + status = 'pending' + AND claimed_by IS NULL + AND completed_at IS NULL + AND receipt_sha256 IS NULL + ) + OR ( + status = 'in_flight' + AND claimed_by IS NOT NULL + AND completed_at IS NULL + AND receipt_sha256 IS NULL + ) + OR ( + status = 'delivered' + AND claimed_by IS NULL + AND completed_at IS NOT NULL + AND last_error_code IS NULL + AND response_status BETWEEN 200 AND 299 + AND response_body_sha256 IS NOT NULL + AND receipt_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 receipt_sha256 IS NULL + ) + ) +); + +CREATE INDEX IF NOT EXISTS idx_gda_lineage_delivery_due + ON gda_control.metadata_fabric_lineage_outbox( + tenant_id, actor_subject, available_at, created_at + ) WHERE status = 'pending'; +CREATE INDEX IF NOT EXISTS idx_gda_lineage_delivery_expired_claim + ON gda_control.metadata_fabric_lineage_outbox(tenant_id, claimed_until) + WHERE status = 'in_flight'; +CREATE INDEX IF NOT EXISTS idx_gda_lineage_delivery_binding + ON gda_control.metadata_fabric_lineage_outbox( + tenant_id, binding_id, created_at + ); + +ALTER TABLE gda_control.metadata_fabric_lineage_outbox + ENABLE ROW LEVEL SECURITY; +ALTER TABLE gda_control.metadata_fabric_lineage_outbox + FORCE ROW LEVEL SECURITY; +DROP POLICY IF EXISTS gda_lineage_delivery_tenant_isolation + ON gda_control.metadata_fabric_lineage_outbox; +CREATE POLICY gda_lineage_delivery_tenant_isolation + ON gda_control.metadata_fabric_lineage_outbox + USING (tenant_id = gda_control.current_tenant()) + WITH CHECK (tenant_id = gda_control.current_tenant()); + +CREATE OR REPLACE FUNCTION gda_control.claim_metadata_fabric_lineage( + p_tenant_id TEXT, + p_actor_subject TEXT, + p_worker_id TEXT, + p_limit INTEGER DEFAULT 10, + p_lease_seconds INTEGER DEFAULT 60 +) +RETURNS SETOF gda_control.metadata_fabric_lineage_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_actor_subject), '') = '' + OR p_actor_subject !~ '^workload:.+' THEN + RAISE EXCEPTION 'lineage actor 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_fabric_lineage_outbox + SET status = 'failed', + claimed_by = NULL, + claimed_until = NULL, + last_error_code = COALESCE(last_error_code, 'lease_expired'), + response_status = NULL, + response_body_sha256 = NULL, + 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 delivery_id + FROM gda_control.metadata_fabric_lineage_outbox + WHERE tenant_id = p_tenant_id + AND actor_subject = p_actor_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, created_at, delivery_id + LIMIT p_limit + FOR UPDATE SKIP LOCKED + ) + UPDATE gda_control.metadata_fabric_lineage_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, + response_status = NULL, + response_body_sha256 = NULL, + completed_at = NULL + FROM candidates + WHERE delivery.tenant_id = p_tenant_id + AND delivery.delivery_id = candidates.delivery_id + RETURNING delivery.*; +END; +$$; + +CREATE OR REPLACE FUNCTION gda_control.complete_metadata_fabric_lineage( + p_tenant_id TEXT, + p_delivery_id UUID, + p_worker_id TEXT, + p_response_status INTEGER, + p_response_body_sha256 TEXT, + p_receipt_sha256 TEXT +) +RETURNS SETOF gda_control.metadata_fabric_lineage_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_response_status < 200 OR p_response_status > 299 THEN + RAISE EXCEPTION 'successful response must be 2xx' + USING ERRCODE = '22023'; + END IF; + IF p_response_body_sha256 !~ '^[0-9a-f]{64}$' + OR p_receipt_sha256 !~ '^[0-9a-f]{64}$' THEN + RAISE EXCEPTION 'delivery receipt fingerprint is invalid' + USING ERRCODE = '22023'; + END IF; + RETURN QUERY + UPDATE gda_control.metadata_fabric_lineage_outbox AS delivery + SET status = 'delivered', + claimed_by = NULL, + claimed_until = NULL, + last_error_code = NULL, + response_status = p_response_status, + response_body_sha256 = p_response_body_sha256, + receipt_sha256 = p_receipt_sha256, + completed_at = clock_timestamp() + WHERE delivery.tenant_id = p_tenant_id + AND delivery.delivery_id = p_delivery_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 'lineage claim is missing, expired, or owned by another worker' + USING ERRCODE = '40001'; + END IF; +END; +$$; + +CREATE OR REPLACE FUNCTION gda_control.fail_metadata_fabric_lineage( + p_tenant_id TEXT, + p_delivery_id UUID, + p_worker_id TEXT, + p_error_code TEXT, + p_response_status INTEGER DEFAULT NULL, + p_retryable BOOLEAN DEFAULT true, + p_retry_delay_seconds INTEGER DEFAULT 30 +) +RETURNS SETOF gda_control.metadata_fabric_lineage_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_response_status IS NOT NULL + AND (p_response_status < 100 OR p_response_status > 599) THEN + RAISE EXCEPTION 'failure response status 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_fabric_lineage_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, + response_status = p_response_status, + response_body_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.delivery_id = p_delivery_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 'lineage claim is missing, expired, or owned by another worker' + USING ERRCODE = '40001'; + END IF; +END; +$$; + +REVOKE ALL ON TABLE gda_control.metadata_fabric_lineage_outbox FROM PUBLIC; +REVOKE ALL ON TABLE gda_control.metadata_fabric_lineage_outbox + FROM gda_control_gateway; +GRANT SELECT, INSERT ON gda_control.metadata_fabric_lineage_outbox + TO gda_control_gateway; + +REVOKE ALL ON FUNCTION gda_control.claim_metadata_fabric_lineage( + text, text, text, integer, integer +) FROM PUBLIC; +REVOKE ALL ON FUNCTION gda_control.complete_metadata_fabric_lineage( + text, uuid, text, integer, text, text +) FROM PUBLIC; +REVOKE ALL ON FUNCTION gda_control.fail_metadata_fabric_lineage( + text, uuid, text, text, integer, boolean, integer +) FROM PUBLIC; +GRANT EXECUTE ON FUNCTION gda_control.claim_metadata_fabric_lineage( + text, text, text, integer, integer +) TO gda_control_gateway; +GRANT EXECUTE ON FUNCTION gda_control.complete_metadata_fabric_lineage( + text, uuid, text, integer, text, text +) TO gda_control_gateway; +GRANT EXECUTE ON FUNCTION gda_control.fail_metadata_fabric_lineage( + text, uuid, text, text, integer, boolean, integer +) TO gda_control_gateway; diff --git a/data_agent/platform_gateway.py b/data_agent/platform_gateway.py index 885443ec..0f110d59 100644 --- a/data_agent/platform_gateway.py +++ b/data_agent/platform_gateway.py @@ -29,6 +29,14 @@ MetadataFabricConfigurationError, build_metadata_fabric_binding, ) +from .metadata_fabric_ingestion import MetadataFabricIngestionPlan +from .metadata_fabric_lineage_delivery_contract import ( + LineageDeliveryStatus, + MetadataFabricLineageDelivery, + delivery_binding_payload, + openlineage_receipt_sha256, + validate_delivery_source, +) from .platform_authorization import ( AuthorizationEvidenceError, parse_approval_artifact, @@ -76,6 +84,11 @@ / "migrations" / "097_metadata_fabric_binding_ledger.sql" ) +METADATA_FABRIC_LINEAGE_MIGRATION = ( + Path(__file__).resolve().parent + / "migrations" + / "098_metadata_fabric_openlineage_delivery.sql" +) USER_TENANT_MIGRATION = ( Path(__file__).resolve().parent / "migrations" @@ -1287,6 +1300,29 @@ def _load_metadata_fabric_binding( return None return cls._metadata_fabric_binding_from_row(row) + @classmethod + def _load_metadata_fabric_binding_by_id( + cls, connection, tenant_id: str, binding_id: UUID + ) -> MetadataFabricBindingRecord | None: + row = connection.execute( + text( + """ + SELECT tenant_id, binding_id, binding_document, + execution_plan_artifact_id, + policy_decision_artifact_id, approval_artifact_id, + provider_evidence_artifact_id, recorded_by, recorded_at, + record_sha256 + FROM gda_control.metadata_fabric_binding + WHERE tenant_id = :tenant_id + AND binding_id = :binding_id + """ + ), + {"tenant_id": tenant_id, "binding_id": binding_id}, + ).mappings().one_or_none() + if row is None: + return None + return cls._metadata_fabric_binding_from_row(row) + def get_metadata_fabric_binding( self, tenant_id: str, resource_version_id: UUID ) -> MetadataFabricBindingRecord: @@ -1507,6 +1543,284 @@ def commit_metadata_fabric_binding( ) return GatewayWriteResult(stored, inserted is not None) + @staticmethod + def _metadata_fabric_lineage_from_row( + row, + ) -> MetadataFabricLineageDelivery: + value = dict(row) + value["event"] = _as_json(value["event"]) + return MetadataFabricLineageDelivery.model_validate(value) + + @classmethod + def _load_metadata_fabric_lineage_delivery( + cls, connection, tenant_id: str, delivery_id: UUID + ) -> MetadataFabricLineageDelivery | None: + row = connection.execute( + text( + """ + SELECT tenant_id, delivery_id, binding_id, + resource_version_id, run_id, source_plan_sha256, + target_name, event, event_sha256, idempotency_key, + actor_subject, status, attempt_count, max_attempts, + available_at, claimed_by, claimed_until, + last_error_code, response_status, + response_body_sha256, receipt_sha256, + created_at, completed_at + FROM gda_control.metadata_fabric_lineage_outbox + WHERE tenant_id = :tenant_id + AND delivery_id = :delivery_id + """ + ), + {"tenant_id": tenant_id, "delivery_id": delivery_id}, + ).mappings().one_or_none() + if row is None: + return None + return cls._metadata_fabric_lineage_from_row(row) + + def enqueue_metadata_fabric_lineage( + self, + delivery: MetadataFabricLineageDelivery, + *, + source_plan: MetadataFabricIngestionPlan, + ) -> GatewayWriteResult: + try: + delivery = MetadataFabricLineageDelivery.model_validate( + delivery.model_dump(mode="json", by_alias=True) + ) + source_plan = MetadataFabricIngestionPlan.model_validate( + source_plan.model_dump(mode="json", by_alias=True) + ) + except ValueError as exc: + raise GatewayValidationError( + "Metadata Fabric lineage input is not content-bound" + ) from exc + if delivery.status != LineageDeliveryStatus.PENDING: + raise GatewayValidationError("new lineage delivery must be pending") + with self._transaction(delivery.tenant_id) as connection: + binding = self._load_metadata_fabric_binding_by_id( + connection, delivery.tenant_id, delivery.binding_id + ) + if binding is None: + raise GatewayValidationError( + "Metadata Fabric lineage binding was not found" + ) + artifact = self._load_artifact( + connection, + delivery.tenant_id, + binding.execution_plan_artifact_id, + ) + if artifact is None: + raise GatewayValidationError( + "Metadata Fabric lineage execution plan was not found" + ) + apply_plan = self._parse_metadata_fabric_execution_plan(artifact) + try: + validate_delivery_source( + binding=binding, + source_plan=source_plan, + apply_plan=apply_plan, + ) + except ValueError as exc: + raise GatewayValidationError(str(exc)) from exc + expected = ( + binding.tenant_id, + binding.binding_id, + binding.binding.resource_version_id, + source_plan.run_id, + source_plan.plan_sha256, + source_plan.openlineage_event, + source_plan.openlineage_event_sha256, + ) + observed = ( + delivery.tenant_id, + delivery.binding_id, + delivery.resource_version_id, + delivery.run_id, + delivery.source_plan_sha256, + delivery.event, + delivery.event_sha256, + ) + if observed != expected or delivery.created_at < binding.recorded_at: + raise GatewayValidationError( + "Metadata Fabric lineage delivery does not match the binding" + ) + if delivery.actor_subject == binding.recorded_by: + raise GatewayValidationError( + "Metadata Fabric lineage emitter must be independent" + ) + inserted = connection.execute( + text( + """ + INSERT INTO gda_control.metadata_fabric_lineage_outbox ( + tenant_id, delivery_id, binding_id, + resource_version_id, run_id, source_plan_sha256, + target_name, event, event_sha256, idempotency_key, + actor_subject, status, attempt_count, max_attempts, + available_at, claimed_by, claimed_until, + last_error_code, response_status, + response_body_sha256, receipt_sha256, + created_at, completed_at + ) VALUES ( + :tenant_id, :delivery_id, :binding_id, + :resource_version_id, :run_id, :source_plan_sha256, + :target_name, CAST(:event AS jsonb), :event_sha256, + :idempotency_key, :actor_subject, :status, + :attempt_count, :max_attempts, :available_at, + :claimed_by, :claimed_until, :last_error_code, + :response_status, :response_body_sha256, + :receipt_sha256, :created_at, :completed_at + ) + ON CONFLICT DO NOTHING + RETURNING delivery_id + """ + ), + { + **delivery.model_dump( + mode="python", + exclude={"delivery_schema", "event"}, + ), + "event": _json( + delivery.event.model_dump(mode="json", by_alias=True) + ), + "status": delivery.status.value, + }, + ).first() + stored = self._load_metadata_fabric_lineage_delivery( + connection, delivery.tenant_id, delivery.delivery_id + ) + if ( + stored is None + or delivery_binding_payload(stored) + != delivery_binding_payload(delivery) + ): + raise GatewayConflictError( + "OpenLineage delivery identity has different content" + ) + return GatewayWriteResult(stored, inserted is not None) + + def get_metadata_fabric_lineage_delivery( + self, tenant_id: str, delivery_id: UUID + ) -> MetadataFabricLineageDelivery: + tenant = _TENANT_ADAPTER.validate_python(tenant_id) + with self._transaction(tenant) as connection: + delivery = self._load_metadata_fabric_lineage_delivery( + connection, tenant, delivery_id + ) + if delivery is None: + raise GatewayNotFoundError( + "Metadata Fabric lineage delivery was not found" + ) + return delivery + + def claim_metadata_fabric_lineage( + self, + tenant_id: str, + worker_id: str, + *, + actor_subject: str, + limit: int = 10, + lease_seconds: int = 60, + ) -> list[MetadataFabricLineageDelivery]: + with self._transaction(tenant_id) as connection: + rows = connection.execute( + text( + """ + SELECT * FROM gda_control.claim_metadata_fabric_lineage( + :tenant_id, :actor_subject, :worker_id, + :limit, :lease_seconds + ) + """ + ), + { + "tenant_id": tenant_id, + "actor_subject": actor_subject, + "worker_id": worker_id, + "limit": limit, + "lease_seconds": lease_seconds, + }, + ).mappings().all() + return [ + self._metadata_fabric_lineage_from_row(row) for row in rows + ] + + def complete_metadata_fabric_lineage( + self, + tenant_id: str, + delivery_id: UUID, + *, + worker_id: str, + response_status: int, + response_body_sha256: str, + ) -> MetadataFabricLineageDelivery: + with self._transaction(tenant_id) as connection: + claimed = self._load_metadata_fabric_lineage_delivery( + connection, tenant_id, delivery_id + ) + if claimed is None: + raise GatewayNotFoundError( + "Metadata Fabric lineage delivery was not found" + ) + receipt = openlineage_receipt_sha256( + claimed, + response_status=response_status, + response_body_sha256=response_body_sha256, + ) + row = connection.execute( + text( + """ + SELECT * FROM gda_control.complete_metadata_fabric_lineage( + :tenant_id, :delivery_id, :worker_id, + :response_status, :response_body_sha256, + :receipt_sha256 + ) + """ + ), + { + "tenant_id": tenant_id, + "delivery_id": delivery_id, + "worker_id": worker_id, + "response_status": response_status, + "response_body_sha256": response_body_sha256, + "receipt_sha256": receipt, + }, + ).mappings().one() + return self._metadata_fabric_lineage_from_row(row) + + def fail_metadata_fabric_lineage( + self, + tenant_id: str, + delivery_id: UUID, + *, + worker_id: str, + error_code: str, + response_status: int | None = None, + retryable: bool = True, + retry_delay_seconds: int = 30, + ) -> MetadataFabricLineageDelivery: + if not re.fullmatch(r"[a-z0-9_]{1,64}", error_code): + raise GatewayValidationError("lineage failure code is invalid") + with self._transaction(tenant_id) as connection: + row = connection.execute( + text( + """ + SELECT * FROM gda_control.fail_metadata_fabric_lineage( + :tenant_id, :delivery_id, :worker_id, :error_code, + :response_status, :retryable, :retry_delay_seconds + ) + """ + ), + { + "tenant_id": tenant_id, + "delivery_id": delivery_id, + "worker_id": worker_id, + "error_code": error_code, + "response_status": response_status, + "retryable": retryable, + "retry_delay_seconds": retry_delay_seconds, + }, + ).mappings().one() + return self._metadata_fabric_lineage_from_row(row) + @staticmethod def _load_lineage( connection, tenant_id: str, lineage_event_id: UUID @@ -1574,6 +1888,7 @@ def build_gateway_report( command_migration: Path | None = None, success_migration: Path | None = None, binding_migration: Path | None = None, + lineage_migration: Path | None = None, gateway_source: Path | None = None, routes_source: Path | None = None, command_consumer_source: Path | None = None, @@ -1592,6 +1907,9 @@ def build_gateway_report( "binding_migration": ( binding_migration or METADATA_FABRIC_BINDING_MIGRATION ).resolve(), + "lineage_migration": ( + lineage_migration or METADATA_FABRIC_LINEAGE_MIGRATION + ).resolve(), "gateway_source": (gateway_source or Path(__file__)).resolve(), "routes_source": (routes_source or GATEWAY_ROUTES_SOURCE).resolve(), "command_consumer_source": ( @@ -1660,6 +1978,17 @@ def build_gateway_report( "ALTER TABLE gda_control.metadata_fabric_binding FORCE ROW LEVEL SECURITY", "GRANT SELECT, INSERT ON gda_control.metadata_fabric_binding", ), + "lineage_migration": ( + "CREATE TABLE IF NOT EXISTS gda_control.metadata_fabric_lineage_outbox", + "FOREIGN KEY (tenant_id, binding_id)", + "FOR UPDATE SKIP LOCKED", + "claim_metadata_fabric_lineage", + "complete_metadata_fabric_lineage", + "fail_metadata_fabric_lineage", + "ALTER TABLE gda_control.metadata_fabric_lineage_outbox", + "FORCE ROW LEVEL SECURITY", + "GRANT SELECT, INSERT ON gda_control.metadata_fabric_lineage_outbox", + ), "gateway_source": ( 'SET LOCAL ROLE "{GATEWAY_DATABASE_ROLE}"', "SELECT set_config('app.current_tenant', :tenant, true)", @@ -1672,6 +2001,10 @@ def build_gateway_report( "def finalize_run_success(", "def commit_metadata_fabric_binding(", "def get_metadata_fabric_binding(", + "def enqueue_metadata_fabric_lineage(", + "def claim_metadata_fabric_lineage(", + "def complete_metadata_fabric_lineage(", + "def fail_metadata_fabric_lineage(", ), "routes_source": ( 'base = "/api/platform/v1"', @@ -1723,6 +2056,7 @@ def build_gateway_report( or forbidden in texts.get("command_migration", "") or forbidden in texts.get("success_migration", "") or forbidden in texts.get("binding_migration", "") + or forbidden in texts.get("lineage_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 7e0b61c6..54ac20fb 100644 --- a/data_agent/platform_truth.py +++ b/data_agent/platform_truth.py @@ -588,6 +588,26 @@ def _config( ), "Protected-environment NetworkPolicy enforcement and tenant isolation gate", ), + RuntimeSpec( + "metadata_openlineage_delivery_rehearsal", + "lineage_delivery_rehearsal", + "governed", + "evidence_durable", + "temporary PostgreSQL outbox + committed local wire evidence", + "metadata-platform", + "local_verification_only", + ( + "data_agent/metadata_fabric_lineage_delivery.py", + "scripts/metadata-fabric-openlineage-delivery.sh", + ), + ( + ( + "data_agent/metadata_fabric_lineage_delivery.py", + "threading.Thread", + ), + ), + "Managed outbox worker and protected authenticated OpenLineage receiver", + ), RuntimeSpec( "datalake_monitor", "monitor_loop", @@ -652,7 +672,7 @@ def _config( "41949811ca1d12a9d8bdbd5e7ecb1ba528be7049af96742306d6f317ab0791b8" ) RUNTIME_PRIMITIVE_BASELINE_FINGERPRINT = ( - "03da85670462b4f77cf37714a9af3dc675c24b578d795efe0e2627075eb7d265" + "d6402d91e40ddb61591a7d258925d79e5eee964c3a9c0ace7de34acd10facbfd" ) _IGNORED_SOURCE_PARTS = frozenset( diff --git a/data_agent/test_metadata_fabric_lineage_delivery.py b/data_agent/test_metadata_fabric_lineage_delivery.py new file mode 100644 index 00000000..75e9031d --- /dev/null +++ b/data_agent/test_metadata_fabric_lineage_delivery.py @@ -0,0 +1,156 @@ +import hashlib +import json +from datetime import timedelta + +import httpx +import pytest +from pydantic import ValidationError + +from data_agent.metadata_fabric_lineage_delivery import ( + DEFAULT_EVIDENCE_PATH, + build_lineage_delivery_bundle, + build_contract_report, + validate_rehearsal_evidence, +) +from data_agent.metadata_fabric_lineage_delivery_contract import ( + MetadataFabricLineageDelivery, + build_metadata_fabric_lineage_delivery, +) +from data_agent.metadata_fabric_lineage_emitter import ( + LineageEmitterProfile, + LineageHttpDeliveryError, + OpenLineageHttpEmitter, +) +from data_agent.platform_contracts import canonical_json_bytes + + +def _claimed_delivery(): + delivery = build_lineage_delivery_bundle().delivery + payload = delivery.model_dump(mode="json", by_alias=True) + payload.update( + { + "status": "in_flight", + "attempt_count": 1, + "claimed_by": "worker:test-lineage-1", + "claimed_until": ( + delivery.created_at + timedelta(minutes=1) + ).isoformat(), + } + ) + return MetadataFabricLineageDelivery.model_validate(payload) + + +def test_delivery_is_deterministic_and_bound_to_authorized_source_plan(): + first = build_lineage_delivery_bundle() + second = build_lineage_delivery_bundle() + + assert first == second + assert str(first.delivery.delivery_id) == ( + "49a54408-b3a8-5843-a27d-6395c080af99" + ) + assert first.delivery.event_sha256 == ( + "4929e51c4126e09415a9fc1578c9401077c5d7c374294e70deeebd29c8216dd2" + ) + assert first.delivery.idempotency_key == ( + "e1a2862b7e246b3717ee2e65cf1a765a40865fdce13eed1b129319d9772c0073" + ) + assert first.delivery.actor_subject != ( + first.binding_bundle.record.recorded_by + ) + + +def test_delivery_rejects_event_and_binding_tampering(): + bundle = build_lineage_delivery_bundle() + payload = bundle.delivery.model_dump(mode="json", by_alias=True) + payload["event"]["job"]["name"] = "tampered" + with pytest.raises(ValidationError, match="event SHA-256"): + MetadataFabricLineageDelivery.model_validate(payload) + + apply_plan_artifact = bundle.binding_bundle.artifacts[0] + from data_agent.metadata_fabric_binding_contract import ( + parse_metadata_fabric_execution_plan_artifact, + ) + + apply_plan = parse_metadata_fabric_execution_plan_artifact( + apply_plan_artifact + ) + mismatched_plan = bundle.source_plan.model_copy( + update={"resource_version_id": bundle.source_plan.source_resource_version_id} + ) + with pytest.raises(ValueError, match="binding identity"): + build_metadata_fabric_lineage_delivery( + binding=bundle.binding_bundle.record, + source_plan=mismatched_plan, + apply_plan=apply_plan, + actor_subject=bundle.delivery.actor_subject, + created_at=bundle.delivery.created_at, + ) + + +def test_emitter_profile_rejects_remote_credentials_and_redirect_targets(): + common = { + "target_name": "sink", + "actor_subject": "workload:lineage", + } + for endpoint in ( + "https://example.com/api/v1/lineage", + "http://user:pass@127.0.0.1/api/v1/lineage", + "http://127.0.0.1/redirect", + "http://127.0.0.1/api/v1/lineage?token=x", + ): + with pytest.raises(ValidationError, match="loopback"): + LineageEmitterProfile(endpoint_url=endpoint, **common) + + +def test_http_emitter_uses_canonical_body_and_stable_idempotency_headers(): + delivery = _claimed_delivery() + success_body = b'{"duplicate":true}' + responses = [ + httpx.Response(503, json={"accepted": True}), + httpx.Response(200, content=success_body), + ] + requests: list[httpx.Request] = [] + + def handler(request: httpx.Request) -> httpx.Response: + requests.append(request) + return responses.pop(0) + + profile = LineageEmitterProfile( + target_name=delivery.target_name, + endpoint_url="http://127.0.0.1:9999/api/v1/lineage", + actor_subject=delivery.actor_subject, + ) + with OpenLineageHttpEmitter( + profile, transport=httpx.MockTransport(handler) + ) as emitter: + with pytest.raises(LineageHttpDeliveryError) as first: + emitter.emit(delivery) + receipt = emitter.emit(delivery) + + assert first.value.code == "http_5xx" + assert first.value.retryable is True + expected = canonical_json_bytes( + delivery.event.model_dump(mode="json", by_alias=True) + ) + assert [request.content for request in requests] == [expected, expected] + assert all( + request.headers["Idempotency-Key"] == delivery.idempotency_key + for request in requests + ) + assert receipt.response_status == 200 + assert receipt.response_body_sha256 == hashlib.sha256( + success_body + ).hexdigest() + + +def test_lineage_delivery_contract_and_committed_evidence_validate(): + report = build_contract_report() + evidence = json.loads(DEFAULT_EVIDENCE_PATH.read_text(encoding="utf-8")) + + assert report["status"] == "valid" + assert report["errors"] == [] + assert validate_rehearsal_evidence(evidence) == [] + tampered = {**evidence, "receiver_unique_accept_count": 2} + assert "lineage delivery evidence SHA-256 does not match" in ( + validate_rehearsal_evidence(tampered) + ) diff --git a/data_agent/test_metadata_fabric_lineage_delivery_postgres.py b/data_agent/test_metadata_fabric_lineage_delivery_postgres.py new file mode 100644 index 00000000..316835c2 --- /dev/null +++ b/data_agent/test_metadata_fabric_lineage_delivery_postgres.py @@ -0,0 +1,129 @@ +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_lineage_delivery import ( + WORKER_ID, + build_lineage_delivery_bundle, + run_local_rehearsal, + validate_rehearsal_evidence, +) +from data_agent.platform_gateway import ( + GatewayNotFoundError, + GatewayValidationError, + 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("lineage delivery test requires a PostgreSQL superuser") + database_name = f"gda_lineage_{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_lineage_delivery_is_tenant_scoped_retryable_and_idempotent(): + 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_attempt_response_status"] == 503 + assert evidence["final_response_status"] == 200 + assert evidence["receiver_unique_accept_count"] == 1 + assert evidence["receiver_duplicate_count"] == 1 + + bundle = build_lineage_delivery_bundle() + engine = create_engine(database_url) + gateway = PlatformGateway(engine) + stored = gateway.get_metadata_fabric_lineage_delivery( + bundle.delivery.tenant_id, + bundle.delivery.delivery_id, + ) + assert stored.status.value == "delivered" + assert stored.attempt_count == 2 + assert gateway.claim_metadata_fabric_lineage( + bundle.delivery.tenant_id, + WORKER_ID, + actor_subject=bundle.delivery.actor_subject, + ) == [] + replay = gateway.enqueue_metadata_fabric_lineage( + bundle.delivery, + source_plan=bundle.source_plan, + ) + assert replay.created is False + assert replay.value == stored + + with pytest.raises(GatewayNotFoundError): + gateway.get_metadata_fabric_lineage_delivery( + "isolated-tenant", + bundle.delivery.delivery_id, + ) + wrong_source = bundle.source_plan.model_copy( + update={ + "resource_version_id": bundle.source_plan.source_resource_version_id + } + ) + with pytest.raises(GatewayValidationError, match="content-bound"): + gateway.enqueue_metadata_fabric_lineage( + bundle.delivery, + source_plan=wrong_source, + ) + + with gateway._transaction(bundle.delivery.tenant_id) as connection: + for statement in ( + """ + UPDATE gda_control.metadata_fabric_lineage_outbox + SET attempt_count = attempt_count + 1 + WHERE delivery_id = :delivery_id + """, + """ + DELETE FROM gda_control.metadata_fabric_lineage_outbox + WHERE delivery_id = :delivery_id + """, + ): + with pytest.raises(DBAPIError): + with connection.begin_nested(): + connection.execute( + text(statement), + {"delivery_id": bundle.delivery.delivery_id}, + ) + finally: + if engine is not None: + engine.dispose() + _drop_temporary_database(admin_engine, database_name) diff --git a/data_agent/test_platform_contracts.py b/data_agent/test_platform_contracts.py index ba7afc88..495fbd9e 100644 --- a/data_agent/test_platform_contracts.py +++ b/data_agent/test_platform_contracts.py @@ -471,7 +471,9 @@ def test_control_ledger_contract_and_migration_catalog_are_valid(): assert report["status"] == "valid" assert report["contract_count"] == 16 assert report["migration"]["sha256"] == migration["checksum"] - assert migrations[-1]["migration_id"] == "097_metadata_fabric_binding_ledger" + assert migrations[-1]["migration_id"] == ( + "098_metadata_fabric_openlineage_delivery" + ) def test_sql_contract_has_tenant_fks_rls_append_only_and_no_legacy_backfill(): diff --git a/docs/architecture-decisions/adr-050-idempotent-openlineage-http-delivery.md b/docs/architecture-decisions/adr-050-idempotent-openlineage-http-delivery.md new file mode 100644 index 00000000..f802adad --- /dev/null +++ b/docs/architecture-decisions/adr-050-idempotent-openlineage-http-delivery.md @@ -0,0 +1,96 @@ +# ADR-050: Idempotent OpenLineage HTTP Delivery + +**Status**: Accepted + +**Date**: 2026-07-28 + +**Decision owners**: Data Platform, Metadata Platform, Data Governance, Security, Platform Architecture + +**Related decisions**: [ADR-006](adr-006-openmetadata-governance-and-active-metadata-platform.md) · [ADR-020](adr-020-platform-resource-run-and-evidence-contracts.md) · [ADR-025](adr-025-platform-command-outbox-and-callback.md) · [ADR-047](adr-047-deterministic-metadata-fabric-ingestion-projection.md) · [ADR-048](adr-048-local-authorized-metadata-fabric-ingestion-replay.md) · [ADR-049](adr-049-tenant-scoped-metadata-fabric-binding-ledger.md) + +## Context + +M3-1 produced a content-bound OpenLineage `COMPLETE` RunEvent candidate. M3-2 applied the associated provider projections after policy and approval validation. M3-3 then persisted the verified provider binding through `PlatformGateway`. The event was still never sent over a wire. + +An HTTP timeout or failed acknowledgement can occur after a receiver commits an event. Treating that outcome as success can lose events; blindly retrying without a stable key can create duplicate lineage. Network-level exactly-once delivery is not available, so the platform needs explicit at-least-once delivery state and receiver idempotency. + +## Options Considered + +| Option | Benefit | Cost/risk | Decision | +|---|---|---|---| +| Send directly from provider apply code | Short path | Couples provider mutation to receiver availability and loses durable retry state | Rejected | +| Mark HTTP 2xx directly on the immutable binding | No new table | Mutates the wrong authority and cannot represent claims, leases or retry | Rejected | +| Introduce Kafka or another event platform now | Mature transport | Adds an unproven second durability and operations boundary | Rejected | +| Add a narrow PostgreSQL outbox with stable HTTP idempotency | Reuses the control database, RLS and lease patterns; preserves authority boundaries | At-least-once delivery still requires receiver dedupe | Adopted | + +## Decision + +### 1. Delivery is derived from the authorized binding chain + +Migration 098 adds `gda_control.metadata_fabric_lineage_outbox`. A deterministic delivery UUID and idempotency key bind the tenant, M3-3 binding UUID, target name and canonical OpenLineage event SHA-256. Only one matching event can be enqueued for that binding and target. + +`PlatformGateway` reloads the binding and its execution-plan Artifact, parses the authorized apply plan and verifies the complete M3-1 source plan before insert. Tenant, ResourceVersion, run, source plan, event and event fingerprint must all match. The lineage emitter workload must be independent from the provider/binding recorder. Exact enqueue replay returns the existing row; different content conflicts or fails validation. + +### 2. PostgreSQL owns delivery state, not lineage truth + +The outbox has `pending`, `in_flight`, `delivered` and `failed` states, bounded attempts, worker identity and lease expiry. `FOR UPDATE SKIP LOCKED` permits multiple workers without double claim. Expired claims are reclaimable until the attempt limit. + +The gateway role receives only table `SELECT/INSERT`. Three `SECURITY DEFINER` functions are the only mutation path for claim, completion and failure. They require transaction-local tenant context, the current claim owner and an unexpired lease. `FORCE ROW LEVEL SECURITY` prevents cross-tenant visibility. + +This state is transport state only. It does not replace the immutable binding, PlatformRun, LineageEvent or receiver state. + +### 3. HTTP is canonical, bounded and idempotent + +The emitter sends canonical OpenLineage JSON to `/api/v1/lineage` with `Content-Type: application/json`, the deterministic `Idempotency-Key`, delivery UUID and event SHA headers. Only a 2xx response completes the outbox. Transport errors, 429 and 5xx are retryable; other 4xx responses fail closed. Response bodies are not stored, only a SHA-256 and content-bound receipt. + +The M3-4 local profile accepts only an unauthenticated loopback HTTP endpoint and rejects credentials, redirects, remote hosts, query strings and fragments. This is a deliberate rehearsal boundary, not a production endpoint profile. + +### 4. The rehearsal injects a commit-then-failed-ack outcome + +The local receiver commits the first idempotency key and exact event, then returns 503. The outbox retains the event for retry. The second request carries the same key and body; the receiver returns a duplicate 200 without accepting a second lineage event. A completed row cannot be reclaimed. + +This proves at-least-once delivery with receiver idempotency across the real Python HTTP stack. It does not prove network exactly-once semantics. + +## Verification + +The local PostgreSQL 16 and loopback HTTP rehearsal records: + +- source M3-3 evidence SHA `518bfed363aba34e539ada19ea1dc708bacc9eba6578ccab165d11bccfc05223`; +- deterministic delivery UUID `49a54408-b3a8-5843-a27d-6395c080af99`; +- OpenLineage event SHA `4929e51c4126e09415a9fc1578c9401077c5d7c374294e70deeebd29c8216dd2`; +- idempotency key `e1a2862b7e246b3717ee2e65cf1a765a40865fdce13eed1b129319d9772c0073`; +- first acknowledgement 503, final acknowledgement 200 and two attempts; +- two wire requests but exactly one receiver acceptance and one duplicate response; +- receipt SHA `13852b5ebee6d0a9546914cd2e2080678910d7fba8dc3dc21e5a049a4220257e`; +- exact enqueue replay `created=false`, completed delivery not reclaimed; +- FORCE RLS, cross-tenant rejection and no direct gateway UPDATE/DELETE; +- evidence SHA `8fa87a34a39b900df0673f11d0301c9f5155ce64ff9502125478ec59a3f0fdb6`. + +Unit tests cover deterministic identity, content tampering, endpoint restrictions and canonical wire headers/body. PostgreSQL integration covers retry completion, replay, source-plan mismatch, claim finality, RLS and direct mutation rejection. + +## Claim Boundary + +Allowed now: + +- the exact M3-1 OpenLineage candidate can be gated by the M3-3 binding and delivered over a real local HTTP connection; +- a receiver commit followed by failed acknowledgement is recovered with a stable idempotency key and no second receiver acceptance; +- tenant-scoped PostgreSQL outbox state supports bounded claim, lease, retry, completion and failure; +- M3-4 performs no OpenMetadata, Gravitino or legacy mutation. + +Fixed false now: + +- protected workload OIDC, TLS, receiver authentication and credential rotation; +- a production OpenLineage/OpenMetadata receiver, receiver HA or durable receiver storage; +- production outbox deployment, worker scaling, alerting, backup/recovery and SLO; +- OpenMetadata minimum privilege, Gravitino authentication and durable catalog conformance; +- production ingestion and `production_ready`. + +## Consequences + +**Positive**: lineage delivery no longer depends on a single synchronous HTTP outcome, and the exact event is traceable back through the authorized provider binding. + +**Negative**: receiver idempotency remains mandatory. A non-idempotent production receiver cannot safely consume this at-least-once transport. + +**Mitigation**: production activation requires a protected receiver profile, workload identity, TLS, receiver-specific idempotency/conformance evidence, managed worker deployment and operational gates. + +**Revisit trigger**: replace PostgreSQL polling only when measured delivery volume or SLO requires another transport and a migration preserves tenant, idempotency, receipt and replay semantics. diff --git a/docs/evidence/metadata-fabric-openlineage-delivery-2026-07-28.json b/docs/evidence/metadata-fabric-openlineage-delivery-2026-07-28.json new file mode 100644 index 00000000..947a9dc1 --- /dev/null +++ b/docs/evidence/metadata-fabric-openlineage-delivery-2026-07-28.json @@ -0,0 +1,44 @@ +{ + "binding_id": "9580cd65-9fd9-5216-90a5-1fd6837e6cfb", + "completed_delivery_not_reclaimed": true, + "cross_tenant_read_blocked": true, + "delivery_id": "49a54408-b3a8-5843-a27d-6395c080af99", + "durable_catalog_verified": false, + "errors": [], + "event_sha256": "4929e51c4126e09415a9fc1578c9401077c5d7c374294e70deeebd29c8216dd2", + "evidence_sha256": "8fa87a34a39b900df0673f11d0301c9f5155ce64ff9502125478ec59a3f0fdb6", + "final_attempt_count": 2, + "final_response_status": 200, + "first_attempt_response_status": 503, + "first_attempt_retry_pending": true, + "first_enqueue_created": true, + "force_rls_verified": true, + "gateway_direct_update_delete_blocked": true, + "idempotency_key": "e1a2862b7e246b3717ee2e65cf1a765a40865fdce13eed1b129319d9772c0073", + "live_openlineage_emission_verified": true, + "local_loopback_receiver": true, + "local_wire_openlineage_delivery_verified": true, + "oidc_verified": false, + "production_ingestion_verified": false, + "production_ready": false, + "production_receiver_verified": false, + "provider_minimum_privilege_verified": false, + "provider_mutations_executed": false, + "receipt_sha256": "13852b5ebee6d0a9546914cd2e2080678910d7fba8dc3dc21e5a049a4220257e", + "receiver_commit_then_failed_ack_recovered": true, + "receiver_credentials_used": false, + "receiver_duplicate_count": 1, + "receiver_unique_accept_count": 1, + "replay_enqueue_created": false, + "schema": "gda.metadata_fabric_openlineage_delivery_evidence.v1", + "source_evidence_sha256": "518bfed363aba34e539ada19ea1dc708bacc9eba6578ccab165d11bccfc05223", + "source_plan_sha256": "a5c8ef636c03a38d0c6edaacff7d1edeba9c4b8a7f1491c493e9308257c5a94d", + "stable_idempotency_header_verified": true, + "status": "local_wire_openlineage_delivery_verified", + "target_name": "local-openlineage-http-sink", + "transport_semantics": "at_least_once_with_receiver_idempotency", + "wire_body_matched": true, + "wire_body_sha256": "4929e51c4126e09415a9fc1578c9401077c5d7c374294e70deeebd29c8216dd2", + "wire_request_count": 2, + "writes_to_legacy": false +} diff --git a/docs/roadmap-ar0-platform-truth-2026-07-24.md b/docs/roadmap-ar0-platform-truth-2026-07-24.md index 00b909f5..507daf64 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-3(本地 binding ledger 已验证) +### 4.8 Metadata Fabric Bridge M1 + M2 + M3-4(本地 OpenLineage wire delivery 已验证) 第八块回到 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: @@ -223,8 +223,9 @@ Temporal 继续保持目标组件状态,不在这一包并行接入。OpenMeta 18. [ADR-047](architecture-decisions/adr-047-deterministic-metadata-fabric-ingestion-projection.md) 已将同一地类图斑 target ResourceVersion 的 output Artifact、passed QualityResult、独立 evaluator、LineageEvent、RunSuccessEvidence 和 M1 binding 合成为确定性双 provider projection plan 与 OpenLineage COMPLETE candidate;plan fingerprint 为 `a5c8ef636c03a38d0c6edaacff7d1edeba9c4b8a7f1491c493e9308257c5a94d`,相同 observation replay 为 `no_op`,任何 owner/domain/tag、Gravitino revision 或 target inventory 漂移均 blocked。该合同不含 provider mutation client,`provider_apply_authorized=false`、`writes_to_gda_control=false`、`writes_to_legacy=false`。 19. [ADR-048](architecture-decisions/adr-048-local-authorized-metadata-fabric-ingestion-replay.md) 已将同一 plan 绑定到精确 PolicyDecision、独立 ApprovalRecord 与 execution-plan Artifact,在本地 OpenMetadata/Gravitino 以 natural key 创建并 read-back;首次 apply 创建 8 个目标层级对象,第二次 replay 为 `no_op/0 mutations`,OpenMetadata provider UUID 与 binding candidate 均来自真实回读。apply plan fingerprint 为 `241cb2018c093f76378d265ab8fb617d161c1be7bd4effa6fad361e9db7522c4`,authorization fingerprint 为 `7bc8f577cbdea8d9979b2606278a52176cc2d723a6159c4e1f35ada0f5bb6db0`,evidence fingerprint 为 `3d5fb07267680520d2f03bf27f354787b7253210eb93ab85aae83d5f5a714dbe`。partial inventory 在 mutation 前 blocked,第二 provider 失败会反向补偿当前 attempt 创建的对象;binding candidate 未写 GDA Control。 20. [ADR-049](architecture-decisions/adr-049-tenant-scoped-metadata-fabric-binding-ledger.md) 已新增 migration 097 与 `PlatformGateway` binding commit/read:真实 provider refs、target/source/definition version、execution-plan、精确 PolicyDecision、独立 Approval 和 provider evidence 必须在同一 tenant 下完整匹配后才可追加。空临时 PostgreSQL 首次提交为 `created=true`、第二次精确 replay 为 `created=false`,FORCE RLS、跨租户不可见和 gateway 无 UPDATE/DELETE 均通过;binding UUID 为 `9580cd65-9fd9-5216-90a5-1fd6837e6cfb`,record SHA 为 `19bdbddedc27d2ed8a35119e8f065a47a02345f9bbd3a51075856cb9587f4176`,evidence SHA 为 `518bfed363aba34e539ada19ea1dc708bacc9eba6578ccab165d11bccfc05223`。M3-3 不调用 provider、不写 legacy,也不覆盖含 synthetic UUID 的既有 Resource。 +21. [ADR-050](architecture-decisions/adr-050-idempotent-openlineage-http-delivery.md) 已新增 migration 098 与 tenant-scoped lineage outbox;Gateway 只有在 M3-3 binding、execution-plan 和完整 M3-1 source plan 精确匹配后才可 enqueue。真实 loopback HTTP 演练让接收端先提交事件再返回 503,第二次以同一 `Idempotency-Key` 重发并返回 duplicate 200;共 2 个 wire requests、1 次唯一接受、最终 2 attempts/delivered,完成项不再 claim。delivery UUID 为 `49a54408-b3a8-5843-a27d-6395c080af99`,event SHA 为 `4929e51c4126e09415a9fc1578c9401077c5d7c374294e70deeebd29c8216dd2`,evidence SHA 为 `8fa87a34a39b900df0673f11d0301c9f5155ce64ff9502125478ec59a3f0fdb6`。该结论是本地 `at_least_once_with_receiver_idempotency`,不是网络 exactly-once 或生产 receiver 证明。 -此处 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 账本,不再次调用 provider。OpenMetadata 使用 bootstrap admin,Gravitino 未认证且 catalog backend 为 memory,restart persistence 未验证。生产持久 binding、OpenLineage candidate、ResourceVersion 和 legacy authority 都未写入;生产最小权限、OIDC、持久 catalog、TLS、tenant isolation、真实 alert/SLO、受保护 provider policy、live OpenLineage、生产 ingest/conformance、两项 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 后幂等恢复,不再次调用 metadata provider。OpenMetadata 使用 bootstrap admin,Gravitino 未认证且 catalog backend 为 memory,restart persistence 未验证。生产持久 binding、ResourceVersion 和 legacy authority 都未写入;生产最小权限、OIDC、持久 catalog、TLS、tenant isolation、真实 receiver/alert/SLO、受保护 provider policy、生产 OpenLineage、生产 ingest/conformance、两项 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 6f8ba4d8..2da1f2a1 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-28 -阶段: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 已验证,生产 provider ingestion、生产观测、生产 policy/tenant isolation 和生产切换仍 `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 已验证,生产 provider ingestion、生产观测、生产 policy/tenant isolation 和生产切换仍 `in_progress` -适用分支:`feat/ar1-metadata-fabric-binding-ledger` +适用分支:`feat/ar1-metadata-fabric-openlineage-delivery` ## 判定规则 @@ -26,7 +26,7 @@ | 在线空间数据 | 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 已验证 -> 生产切换待验收 | | 技术元数据 | 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 | harvester 结果、合成 response、本地 sandbox/recovery/metrics/policy/ingestion observation、projection plan、provider evidence、binding ledger 与 readiness report | 源系统技术对象是原始证据;Gravitino 映射并联邦,不能覆盖业务 ResourceVersion;GDA binding ledger 只记录已验证关系,本地 memory catalog/evidence 不得冒充生产持久技术权威 | Metadata Platform | AR-1 M1/M2 + M3-3 local ledger 已验证 -> 认证持久 catalog/生产 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,OpenLineage 未发送 | 搜索/页面视图、合成 response、本地 sandbox/recovery/metrics/policy/ingestion observation、projection/provider evidence、binding ledger 与 OpenLineage candidate | OpenMetadata 为 owner/glossary/classification/quality discoverability 权威;GDA ledger 保留审批与 provider identity 关系,不反写 ResourceVersion;本地 admin apply/临时账本不等于生产权威 | Governance | AR-1 M1/M2 + M3-3 local ledger 已验证 -> 最小权限 ingestion/生产持久 binding/live OpenLineage 待执行 | +| 治理目录 | 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 | 搜索/页面视图、合成 response、本地 sandbox/recovery/metrics/policy/ingestion observation、projection/provider evidence、binding ledger、lineage outbox/receipt 与 OpenLineage event | OpenMetadata 为 owner/glossary/classification/quality discoverability 权威;GDA ledger 保留审批/provider identity,outbox 只拥有投递状态,receiver 拥有接收状态;均不反写 ResourceVersion | Governance | AR-1 M1/M2 + M3-4 local wire delivery 已验证 -> 最小权限 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/生产切换待验收 | @@ -56,7 +56,7 @@ 11. `platform_command_outbox` 只拥有投递状态;callback 只触发 reconcile,不能把 provider payload 直接写成 PlatformRun 状态或平台终局。 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,且不调用 provider、不写 legacy。本地 bootstrap admin、未认证 Gravitino、memory catalog、临时 ledger、pending profile、合成 attestation 和 local evidence 都不等于生产最小权限/OIDC、持久 catalog/binding、live OpenLineage、tenant isolation、alert/SLO、生产 ingestion/conformance 或生产写权威。 +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。本地 bootstrap admin、未认证 Gravitino、memory catalog、临时 ledger/outbox、loopback receiver、pending profile、合成 attestation 和 local evidence 都不等于生产最小权限/OIDC、持久 catalog/binding、受保护 OpenLineage receiver、tenant isolation、alert/SLO、生产 ingestion/conformance 或生产写权威。 ## 已建立的 AR-0/AR-1 entry 证据 @@ -83,6 +83,7 @@ - Metadata Fabric M3-1 已将地类图斑 output Artifact、passed QualityResult、独立 evaluator、LineageEvent、RunSuccessEvidence 与 M1 binding 合成为 2 个 provider projection 和 OpenLineage COMPLETE candidate;plan fingerprint 为 `a5c8ef636c03a38d0c6edaacff7d1edeba9c4b8a7f1491c493e9308257c5a94d`,相同 observation replay 为 `no_op`,replay fingerprint 为 `c33857b2ae75f1106ed7d59e8e53296a3f76f4b90ef386238159b328d47c57ca`。`provider_apply_authorized=false`、`provider_mutations_executed=false`、`live_provider_ingestion_verified=false`、`production_ready=false`。 - Metadata Fabric M3-2 已将同一 plan 绑定 execution-plan Artifact、精确 PolicyDecision 与独立 ApprovalRecord,在本地 provider 以 natural key 创建后 read-back;首次 apply 为 `created/8 mutations`,第二次为 `no_op/0 mutations`,真实 OpenMetadata UUID 进入未持久化 binding candidate。partial inventory 在 mutation 前阻断,失败 attempt 反向补偿,端口转发已退出且 evidence 不含凭据;evidence fingerprint 为 `3d5fb07267680520d2f03bf27f354787b7253210eb93ab85aae83d5f5a714dbe`。`provider_minimum_privilege_verified=false`、`oidc_verified=false`、`gravitino_authentication_verified=false`、`binding_persisted_to_gda_control=false`、`live_openlineage_emission_verified=false`、`production_ingestion_verified=false`、`production_ready=false`。 - Metadata Fabric M3-3 已新增 tenant-scoped append-only binding ledger,将 M3-2 的真实 OpenMetadata UUID、Gravitino ref、execution-plan、PolicyDecision、Approval 与 provider evidence 经 PlatformGateway 精确校验后落账;空临时 PostgreSQL 首次提交 `created=true`、replay `created=false`,跨租户读取、UPDATE/DELETE 均拒绝。binding SHA 为 `125d7197f05ff9c37999a94d090d123dcf905480b776da0738d9625ab5045598`,record SHA 为 `19bdbddedc27d2ed8a35119e8f065a47a02345f9bbd3a51075856cb9587f4176`,evidence SHA 为 `518bfed363aba34e539ada19ea1dc708bacc9eba6578ccab165d11bccfc05223`。该演练不调用 provider,`provider_minimum_privilege_verified=false`、`oidc_verified=false`、`durable_catalog_verified=false`、`live_openlineage_emission_verified=false`、`production_ingestion_verified=false`、`production_ready=false`。 +- Metadata Fabric M3-4 已新增 tenant-scoped lineage outbox、Gateway claim/complete/fail 和严格 loopback HTTP emitter;接收端第一次提交 event 后返回 503,第二次以同一 idempotency key/body 重发并只返回 duplicate 200,形成 2 wire requests、1 unique accept、2 attempts 和不可再 claim 的 delivered receipt。delivery UUID 为 `49a54408-b3a8-5843-a27d-6395c080af99`,event SHA 为 `4929e51c4126e09415a9fc1578c9401077c5d7c374294e70deeebd29c8216dd2`,evidence SHA 为 `8fa87a34a39b900df0673f11d0301c9f5155ce64ff9502125478ec59a3f0fdb6`。`local_wire_openlineage_delivery_verified=true` 仅表示真实本机 HTTP wire;`provider_minimum_privilege_verified=false`、`oidc_verified=false`、`durable_catalog_verified=false`、`production_receiver_verified=false`、`production_ingestion_verified=false`、`production_ready=false`。 ## 下一验收证据 @@ -91,5 +92,5 @@ - 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 故障恢复和无双写证据; - 首条真实图斑链对 golden slice 的 output hash、独立质量结果/evidence、血缘、发布 revision 和 rollback 演练; -- OpenMetadata/Gravitino 的 source host/cluster 外生产 backup account/bucket、KMS/TLS/workload identity、PITR/source-loss recovery、RPO/RTO、OIDC、受保护环境 provider NetworkPolicy/tenant isolation、upgrade/rollback、registry provenance、持续 metrics backend/retention/query、真实 alert delivery/SLO owner/runbook,以及受保护 PolicyDecision/Approval、最小权限 provider ingestion、生产持久 binding、live OpenLineage、无双写 read-back 和 conformance;M1 fixture、M2 本地 evidence/readiness contracts、M3-1 projection candidate、M3-2 local replay 与 M3-3 临时 binding ledger 均不计入生产退出门; +- OpenMetadata/Gravitino 的 source host/cluster 外生产 backup account/bucket、KMS/TLS/workload identity、PITR/source-loss recovery、RPO/RTO、OIDC、受保护环境 provider NetworkPolicy/tenant isolation、upgrade/rollback、registry provenance、持续 metrics backend/retention/query、真实 alert delivery/SLO owner/runbook,以及受保护 PolicyDecision/Approval、最小权限 provider ingestion、生产持久 binding、受保护 production OpenLineage receiver、无双写 read-back 和 conformance;M1 fixture、M2 本地 evidence/readiness contracts、M3-1 projection candidate、M3-2 local replay、M3-3 临时 binding ledger 与 M3-4 loopback delivery 均不计入生产退出门; - DolphinScheduler/Temporal sandbox 的独立数据库、备份恢复、身份、版本和升级责任证明;DolphinScheduler standalone/H2 不计入此退出门。 diff --git a/scripts/metadata-fabric-openlineage-delivery.sh b/scripts/metadata-fabric-openlineage-delivery.sh new file mode 100755 index 00000000..8b8a41bb --- /dev/null +++ b/scripts/metadata-fabric-openlineage-delivery.sh @@ -0,0 +1,52 @@ +#!/usr/bin/env bash +set -euo pipefail + +REPO_ROOT="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)" +COMMON_GIT_DIR="$(git -C "${REPO_ROOT}" rev-parse --path-format=absolute --git-common-dir 2>/dev/null || true)" +SHARED_ROOT="" +if [ -n "${COMMON_GIT_DIR}" ]; then + SHARED_ROOT="$(cd "${COMMON_GIT_DIR}/.." && pwd)" +fi +if [ -n "${PYTHON:-}" ]; then + PYTHON_BIN="${PYTHON}" +elif [ -x "${REPO_ROOT}/.venv/bin/python" ]; then + PYTHON_BIN="${REPO_ROOT}/.venv/bin/python" +elif [ -n "${SHARED_ROOT}" ] && [ -x "${SHARED_ROOT}/.venv/bin/python" ]; then + PYTHON_BIN="${SHARED_ROOT}/.venv/bin/python" +else + PYTHON_BIN="python" +fi +EVIDENCE_OUT="${GDA_LINEAGE_DELIVERY_EVIDENCE_OUT:-${REPO_ROOT}/docs/evidence/metadata-fabric-openlineage-delivery-2026-07-28.json}" +CONTAINER_NAME="gda-metadata-lineage-${$}-${RANDOM}" +DATABASE_NAME="gda_lineage" +DATABASE_PASSWORD="$(openssl rand -hex 24)" + +cleanup() { + docker rm -f "${CONTAINER_NAME}" >/dev/null 2>&1 || true +} +trap cleanup EXIT INT TERM + +docker run --detach --rm \ + --name "${CONTAINER_NAME}" \ + --env "POSTGRES_PASSWORD=${DATABASE_PASSWORD}" \ + --env "POSTGRES_DB=${DATABASE_NAME}" \ + --publish 127.0.0.1::5432 \ + postgres:16-alpine >/dev/null + +for _attempt in $(seq 1 60); do + if docker exec "${CONTAINER_NAME}" pg_isready \ + --username postgres --dbname "${DATABASE_NAME}" >/dev/null 2>&1; then + break + fi + sleep 1 +done +docker exec "${CONTAINER_NAME}" pg_isready \ + --username postgres --dbname "${DATABASE_NAME}" >/dev/null + +HOST_PORT="$(docker port "${CONTAINER_NAME}" 5432/tcp | sed -E 's/.*:([0-9]+)$/\1/')" +DATABASE_URL="postgresql+psycopg2://postgres:${DATABASE_PASSWORD}@127.0.0.1:${HOST_PORT}/${DATABASE_NAME}" + +cd "${REPO_ROOT}" +"${PYTHON_BIN}" -m data_agent.metadata_fabric_lineage_delivery rehearse \ + --database-url "${DATABASE_URL}" \ + --evidence-out "${EVIDENCE_OUT}"