diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 894d2971..d083ca5c 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -173,6 +173,9 @@ jobs: - name: Validate metadata fabric Active Metadata consumer evidence run: python -m data_agent.metadata_fabric_active_metadata_consumer validate + - name: Validate metadata fabric Active Metadata authorization evidence + run: python -m data_agent.metadata_fabric_active_metadata_authorization validate + - name: Validate Active Metadata consumer deployment boundary run: python -m data_agent.active_metadata_consumer_deployment validate @@ -212,6 +215,11 @@ jobs: DATABASE_URL: postgresql://postgres:postgres@localhost:5432/gis_agent_test run: python -m pytest data_agent/test_active_metadata_consumer_postgres.py -q + - name: Verify Active Metadata authorization on PostgreSQL + env: + DATABASE_URL: postgresql://postgres:postgres@localhost:5432/gis_agent_test + run: python -m pytest data_agent/test_active_metadata_authorization_postgres.py -q + - name: Run required platform tests env: DATABASE_URL: postgresql://postgres:postgres@localhost:5432/gis_agent_test @@ -241,11 +249,13 @@ jobs: data_agent/test_metadata_fabric_ingestion_replay.py \ data_agent/test_metadata_fabric_binding_contract.py \ data_agent/test_active_metadata_change_contract.py \ + data_agent/test_active_metadata_authorization.py \ data_agent/test_active_metadata_consumer.py \ data_agent/test_active_metadata_consumer_deployment.py \ data_agent/test_active_metadata_consumer_worker.py \ data_agent/test_metadata_fabric_active_metadata_outbox.py \ data_agent/test_metadata_fabric_active_metadata_consumer.py \ + data_agent/test_metadata_fabric_active_metadata_authorization.py \ data_agent/test_metadata_fabric_lineage_delivery.py \ data_agent/test_metadata_fabric_provider_identity.py \ data_agent/test_metadata_fabric_gravitino_identity.py \ @@ -268,6 +278,7 @@ jobs: data_agent/test_platform_truth.py \ data_agent/test_route_registration.py \ data_agent/test_s3_obs.py \ + data_agent/test_spatial_dataset_bundle.py \ data_agent/test_staging_candidate_evidence.py \ data_agent/test_staging_deployment_bundle.py \ data_agent/test_staging_live_evidence.py \ diff --git a/data_agent/active_metadata_authorization.py b/data_agent/active_metadata_authorization.py new file mode 100644 index 00000000..8cf624e0 --- /dev/null +++ b/data_agent/active_metadata_authorization.py @@ -0,0 +1,266 @@ +"""Content-bound authorization for promoting an Active Metadata request.""" + +from __future__ import annotations + +from datetime import UTC, datetime +from typing import Any, Literal, Self +from uuid import UUID, uuid5 + +from pydantic import BaseModel, ConfigDict, Field, field_validator, model_validator + +from .active_metadata_change_contract import ( + METADATA_PROJECTION_ROUTE, + MetadataActivationRequest, + WorkloadSubject, +) +from .platform_authorization import ( + AuthorizationEvidenceError, + parse_policy_decision_artifact, + validate_run_authorization_evidence, +) +from .platform_contracts import ( + Artifact, + PlatformDefinitionVersion, + PlatformRun, + ResourceVersion, + RunStatus, + Sha256, + TenantId, + canonical_json_fingerprint, +) + +AUTHORIZATION_SCHEMA = "gda.metadata_activation_authorization.v1" +DISPATCH_ACTION = "dolphinscheduler.dispatch" + + +class MetadataActivationAuthorizationError(RuntimeError): + """The activation request is not bound to valid dispatch evidence.""" + + +class MetadataActivationAuthorization(BaseModel): + """Immutable proof that one inert request may enqueue one dispatch.""" + + model_config = ConfigDict(extra="forbid", frozen=True) + + authorization_schema: Literal[ + "gda.metadata_activation_authorization.v1" + ] = Field(default=AUTHORIZATION_SCHEMA, alias="schema") + authorization_id: UUID + tenant_id: TenantId + request_id: UUID + request_sha256: Sha256 + resource_urn: str + resource_version_id: UUID + content_sha256: Sha256 + definition_version_id: UUID + definition_sha256: Sha256 + run_id: UUID + execution_plan_artifact_id: UUID + execution_plan_sha256: Sha256 + policy_decision_artifact_id: UUID + policy_decision_sha256: Sha256 + approval_artifact_id: UUID + approval_sha256: Sha256 + command_id: UUID + route: Literal["metadata_fabric.projection_plan"] = METADATA_PROJECTION_ROUTE + status: Literal["authorized_for_dispatch"] = "authorized_for_dispatch" + authorized_by: WorkloadSubject + authorized_at: datetime + scheduler_command_enqueued: Literal[True] = True + provider_apply_authorized: Literal[False] = False + provider_mutations_executed: Literal[False] = False + production_scheduler_submission_verified: Literal[False] = False + production_ingestion_verified: Literal[False] = False + production_ready: Literal[False] = False + authorization_sha256: Sha256 + + @field_validator("authorized_at") + @classmethod + def _utc_authorized_at(cls, value: datetime) -> datetime: + if value.tzinfo is None or value.utcoffset() is None: + raise ValueError("authorization time must include a timezone") + return value.astimezone(UTC) + + @model_validator(mode="after") + def _content_bound(self) -> Self: + identity = _authorization_identity(self.model_dump(mode="json", by_alias=True)) + expected_id = uuid5( + self.request_id, + f"metadata-activation-authorization:{canonical_json_fingerprint(identity)}", + ) + if self.authorization_id != expected_id: + raise ValueError("authorization ID does not match its evidence binding") + stable = self.model_dump( + mode="json", + by_alias=True, + exclude={"authorization_sha256"}, + ) + if self.authorization_sha256 != canonical_json_fingerprint(stable): + raise ValueError("authorization SHA-256 does not match") + return self + + +def dispatch_dedupe_key(run_id: UUID, execution_plan_artifact_id: UUID) -> str: + return f"dolphinscheduler.dispatch:{run_id}:{execution_plan_artifact_id}" + + +def dispatch_command_id(run_id: UUID, execution_plan_artifact_id: UUID) -> UUID: + dedupe_key = dispatch_dedupe_key(run_id, execution_plan_artifact_id) + return uuid5(run_id, dedupe_key) + + +def _authorization_identity(values: dict[str, Any]) -> dict[str, Any]: + excluded = { + "authorization_id", + "authorization_sha256", + "authorized_at", + "status", + "scheduler_command_enqueued", + "provider_apply_authorized", + "provider_mutations_executed", + "production_scheduler_submission_verified", + "production_ingestion_verified", + "production_ready", + } + return {key: value for key, value in values.items() if key not in excluded} + + +def build_metadata_activation_authorization( + request: MetadataActivationRequest, + resource_version: ResourceVersion, + definition: PlatformDefinitionVersion, + run: PlatformRun, + execution_plan_artifact: Artifact, + policy_decision_artifact: Artifact, + approval_artifact: Artifact, + *, + authorized_by: str, + authorized_at: datetime, +) -> MetadataActivationAuthorization: + """Validate the complete chain and return its deterministic authorization.""" + if authorized_at.tzinfo is None or authorized_at.utcoffset() is None: + raise MetadataActivationAuthorizationError( + "authorization time must include a timezone" + ) + authorized_at = authorized_at.astimezone(UTC) + if not authorized_by.startswith("workload:"): + raise MetadataActivationAuthorizationError( + "activation authorizer must use workload identity" + ) + intent = request.intent + if ( + resource_version.tenant_id != intent.tenant_id + or resource_version.resource_urn != intent.resource_urn + or resource_version.resource_version_id != intent.resource_version_id + or resource_version.content_sha256 != intent.content_sha256 + ): + raise MetadataActivationAuthorizationError( + "activation request does not match the ResourceVersion" + ) + if ( + definition.tenant_id != intent.tenant_id + or definition.definition_version_id != run.definition_version_id + or definition.orchestration_class != run.orchestration_class + or definition.orchestration_class.value != "dataops" + or definition.capability_id != intent.route + ): + raise MetadataActivationAuthorizationError( + "activation request does not match the DataOps DefinitionVersion" + ) + run_actor = ( + f"{run.subject_context.subject_type.value}:{run.subject_context.subject_id}" + ) + if ( + run.tenant_id != intent.tenant_id + or run.status != RunStatus.ACCEPTED + or run.subject_context.subject_type.value != "workload" + or authorized_at < run.submitted_at + or intent.resource_version_id + not in {binding.resource_version_id for binding in run.input_bindings} + ): + raise MetadataActivationAuthorizationError( + "activation request does not match an accepted workload Run input" + ) + if run.policy_refs is None: + raise MetadataActivationAuthorizationError( + "activation Run requires immutable policy references" + ) + if ( + run.policy_refs.policy_decision_artifact_id + != policy_decision_artifact.artifact_id + or run.policy_refs.approval_artifact_id != approval_artifact.artifact_id + ): + raise MetadataActivationAuthorizationError( + "activation Run policy references do not match supplied evidence" + ) + try: + decision, approval = validate_run_authorization_evidence( + run, + policy_decision_artifact, + approval_artifact, + execution_plan_artifact, + at=authorized_at, + expected_action=DISPATCH_ACTION, + ) + except AuthorizationEvidenceError as exc: + raise MetadataActivationAuthorizationError(str(exc)) from exc + decision = parse_policy_decision_artifact(policy_decision_artifact) + independent_subjects = {run_actor, decision.evaluator_subject} + if approval is None: + raise MetadataActivationAuthorizationError( + "Active Metadata dispatch requires approval evidence" + ) + independent_subjects.add(approval.approver_subject) + if authorized_by in independent_subjects: + raise MetadataActivationAuthorizationError( + "activation authorizer must be independent from execution and review" + ) + if execution_plan_artifact.created_at > authorized_at: + raise MetadataActivationAuthorizationError( + "execution plan artifact postdates authorization" + ) + + command_id = dispatch_command_id(run.run_id, execution_plan_artifact.artifact_id) + values: dict[str, Any] = { + "tenant_id": intent.tenant_id, + "request_id": request.request_id, + "request_sha256": request.request_sha256, + "resource_urn": intent.resource_urn, + "resource_version_id": intent.resource_version_id, + "content_sha256": intent.content_sha256, + "definition_version_id": definition.definition_version_id, + "definition_sha256": definition.definition_sha256, + "run_id": run.run_id, + "execution_plan_artifact_id": execution_plan_artifact.artifact_id, + "execution_plan_sha256": execution_plan_artifact.content_sha256, + "policy_decision_artifact_id": policy_decision_artifact.artifact_id, + "policy_decision_sha256": policy_decision_artifact.content_sha256, + "approval_artifact_id": approval_artifact.artifact_id, + "approval_sha256": approval_artifact.content_sha256, + "command_id": command_id, + "authorized_by": authorized_by, + "authorized_at": authorized_at, + } + json_values = MetadataActivationAuthorization.model_construct( + authorization_id=UUID(int=0), + authorization_sha256="0" * 64, + **values, + ).model_dump(mode="json", by_alias=True) + authorization_id = uuid5( + request.request_id, + "metadata-activation-authorization:" + + canonical_json_fingerprint(_authorization_identity(json_values)), + ) + stable_model = MetadataActivationAuthorization.model_construct( + authorization_id=authorization_id, + authorization_sha256="0" * 64, + **values, + ) + stable = stable_model.model_dump( + mode="json", by_alias=True, exclude={"authorization_sha256"} + ) + return MetadataActivationAuthorization( + authorization_id=authorization_id, + authorization_sha256=canonical_json_fingerprint(stable), + **values, + ) diff --git a/data_agent/metadata_fabric_active_metadata_authorization.py b/data_agent/metadata_fabric_active_metadata_authorization.py new file mode 100644 index 00000000..d1da1a2e --- /dev/null +++ b/data_agent/metadata_fabric_active_metadata_authorization.py @@ -0,0 +1,759 @@ +"""Validate and rehearse evidence-bound Active Metadata dispatch promotion.""" + +from __future__ import annotations + +import argparse +import hashlib +import json +from dataclasses import dataclass +from datetime import UTC, datetime, timedelta +from pathlib import Path +from typing import Any +from uuid import UUID + +from sqlalchemy import create_engine, text +from sqlalchemy.exc import DBAPIError + +from .active_metadata_authorization import ( + MetadataActivationAuthorization, + build_metadata_activation_authorization, +) +from .active_metadata_change_contract import ( + ActiveMetadataRegistration, + MetadataActivationRequest, + build_active_metadata_registration, + build_metadata_activation_intent, + build_metadata_activation_request, +) +from .platform_authorization import ( + build_approval_artifact, + build_policy_decision_artifact, +) +from .platform_contracts import ( + ApprovalRecord, + Artifact, + PlatformDefinitionVersion, + PlatformRun, + PolicyDecision, + Resource, + ResourceVersion, + RunPolicyReferences, + SubjectContext, + canonical_json_fingerprint, + platform_definition_fingerprint, +) +from .platform_gateway import ( + DefinitionRegistration, + GatewayNotFoundError, + GatewayValidationError, + PlatformGateway, +) +from .spatial_dataset_bundle import ( + build_shapefile_bundle_inventory, + validate_shapefile_bundle_inventory, +) + +CONTRACT_SCHEMA = "gda.active_metadata_authorization_contract.v1" +EVIDENCE_SCHEMA = "gda.active_metadata_authorization_evidence.v1" +REPO_ROOT = Path(__file__).resolve().parent.parent +DEFAULT_EVIDENCE_PATH = ( + REPO_ROOT + / "docs/evidence/metadata-fabric-active-metadata-authorization-2026-07-30.json" +) +DEFAULT_WRAPPER_PATH = ( + REPO_ROOT / "scripts/metadata-fabric-active-metadata-authorization.sh" +) +TENANT = "metadata-authorization-local" +CONSUMER_SUBJECT = "workload:active-metadata-consumer" +WORKER = "worker:active-metadata-consumer-1" +AUTHORIZER = "workload:metadata-activation-authorizer" +SOURCE_ID = UUID("a6000000-0000-4000-8000-000000000001") +DEFINITION_ID = UUID("a6000000-0000-4000-8000-000000000002") +RUN_ID = UUID("a6000000-0000-4000-8000-000000000003") +PLAN_ID = UUID("a6000000-0000-4000-8000-000000000004") +REHEARSAL_TIME = datetime(2026, 7, 30, 8, 0, tzinfo=UTC) +MIGRATIONS = tuple( + Path(__file__).resolve().parent / "migrations" / filename + for filename in ( + "092_platform_control_ledger.sql", + "093_app_user_tenant_context.sql", + "094_platform_control_gateway.sql", + "095_platform_command_outbox.sql", + "096_platform_success_verdict.sql", + "099_active_metadata_change_outbox.sql", + "100_active_metadata_activation_request.sql", + "101_active_metadata_authorization.sql", + ) +) + + +class ActiveMetadataAuthorizationEvidenceError(RuntimeError): + """The authorization contract or local rehearsal failed closed.""" + + +@dataclass(frozen=True) +class AuthorizationBundle: + source_resource: Resource + registration: ActiveMetadataRegistration + request: MetadataActivationRequest + definition_registration: DefinitionRegistration + execution_plan: Artifact + policy_decision: Artifact + approval: Artifact + run: PlatformRun + authorization: MetadataActivationAuthorization + + +def build_authorization_bundle(content_sha256: str) -> AuthorizationBundle: + source_urn = f"gda://{TENANT}/dataset/chongqing-cultural-districts" + source_resource = Resource( + tenant_id=TENANT, + resource_urn=source_urn, + resource_kind="dataset", + authority_system="local_acceptance_bundle", + authority_locator="chongqing-cultural-districts", + owner_ref="team:metadata-platform", + governance_ref={"claim_level": "acceptance_input_only"}, + ) + source_version = ResourceVersion( + tenant_id=TENANT, + resource_urn=source_urn, + resource_version_id=SOURCE_ID, + version_key="cultural-district-bundle-v1", + content_sha256=content_sha256, + authority_version_ref={ + "source_label": "chongqing-central-cultural-districts", + "path_committed": False, + }, + created_by="workload:metadata-registrar", + created_at=REHEARSAL_TIME - timedelta(hours=2), + ) + registration = build_active_metadata_registration( + source_version, + consumer_subject=CONSUMER_SUBJECT, + ) + request = build_metadata_activation_request( + build_metadata_activation_intent( + registration.event, + routed_by=CONSUMER_SUBJECT, + ) + ) + + definition_urn = f"gda://{TENANT}/definition/metadata-projection" + definition_document = {"tasks": ["project-governance-metadata"]} + input_contract = {"metadata_change": "dataset"} + output_contract = {"projection_plan": "artifact"} + definition_sha256 = platform_definition_fingerprint( + orchestration_class="dataops", + capability_id="metadata_fabric.projection_plan", + portability_class="portable", + definition_document=definition_document, + input_contract=input_contract, + output_contract=output_contract, + ) + definition_resource = Resource( + tenant_id=TENANT, + resource_urn=definition_urn, + resource_kind="definition", + authority_system="gda", + authority_locator="definition/metadata-projection", + owner_ref="team:metadata-platform", + ) + definition_version = ResourceVersion( + tenant_id=TENANT, + resource_urn=definition_urn, + resource_version_id=DEFINITION_ID, + version_key="v1", + content_sha256=definition_sha256, + authority_version_ref={"definition_revision": 1}, + created_by="workload:metadata-definition-registrar", + created_at=REHEARSAL_TIME - timedelta(hours=2), + ) + definition = PlatformDefinitionVersion( + tenant_id=TENANT, + definition_urn=definition_urn, + definition_version_id=DEFINITION_ID, + orchestration_class="dataops", + capability_id="metadata_fabric.projection_plan", + portability_class="portable", + definition_document=definition_document, + input_contract=input_contract, + output_contract=output_contract, + definition_sha256=definition_sha256, + ) + definition_registration = DefinitionRegistration( + resource=definition_resource, + resource_version=definition_version, + definition=definition, + ) + plan_manifest = { + "schema": "gda.metadata_projection_execution_plan.v1", + "route": "metadata_fabric.projection_plan", + } + execution_plan = Artifact( + tenant_id=TENANT, + artifact_id=PLAN_ID, + artifact_key="metadata-projection-plan", + artifact_role="execution_plan", + storage_uri=f"postgresql://gda-control/execution-plans/{TENANT}/{PLAN_ID}", + media_type="application/vnd.gda.metadata-projection-plan+json", + content_sha256=canonical_json_fingerprint(plan_manifest), + size_bytes=len( + json.dumps( + plan_manifest, + ensure_ascii=True, + sort_keys=True, + separators=(",", ":"), + ).encode("utf-8") + ), + resource_version_id=DEFINITION_ID, + manifest=plan_manifest, + created_by="workload:metadata-plan-compiler", + created_at=REHEARSAL_TIME - timedelta(hours=1), + ) + subject = SubjectContext( + tenant_id=TENANT, + subject_id="metadata-projection-runner", + subject_type="workload", + roles=("metadata_projector",), + purpose="project active metadata change", + ) + decision = PolicyDecision( + tenant_id=TENANT, + run_id=RUN_ID, + subject_context=subject, + action="dolphinscheduler.dispatch", + definition_version_id=DEFINITION_ID, + resource_version_ids=(DEFINITION_ID, SOURCE_ID), + execution_plan_artifact_id=PLAN_ID, + effect="allow", + policy_version_ref=f"gda://{TENANT}/policy/metadata-dispatch-v1", + evaluator_subject="workload:metadata-policy-evaluator", + requires_approval=True, + decided_at=REHEARSAL_TIME - timedelta(minutes=30), + expires_at=REHEARSAL_TIME + timedelta(days=365), + ) + policy_decision = build_policy_decision_artifact(decision) + approval_record = ApprovalRecord( + tenant_id=TENANT, + run_id=RUN_ID, + definition_version_id=DEFINITION_ID, + policy_decision_artifact_id=policy_decision.artifact_id, + policy_decision_sha256=policy_decision.content_sha256, + verdict="approved", + approver_subject="human:metadata-governance-approver", + reason="approved bounded metadata projection", + decided_at=REHEARSAL_TIME - timedelta(minutes=20), + expires_at=REHEARSAL_TIME + timedelta(days=180), + ) + approval = build_approval_artifact(approval_record) + run = PlatformRun( + tenant_id=TENANT, + run_id=RUN_ID, + definition_version_id=DEFINITION_ID, + orchestration_class="dataops", + subject_context=subject, + input_bindings=( + { + "binding_name": "metadata_change", + "resource_version_id": SOURCE_ID, + "semantic_type": "gis.cultural_districts", + }, + ), + idempotency_key="metadata-projection:cultural-districts:v1", + policy_refs=RunPolicyReferences( + policy_decision_artifact_id=policy_decision.artifact_id, + approval_artifact_id=approval.artifact_id, + ), + submitted_at=REHEARSAL_TIME - timedelta(minutes=10), + ) + authorization = build_metadata_activation_authorization( + request, + source_version, + definition, + run, + execution_plan, + policy_decision, + approval, + authorized_by=AUTHORIZER, + authorized_at=REHEARSAL_TIME, + ) + return AuthorizationBundle( + source_resource=source_resource, + registration=registration, + request=request, + definition_registration=definition_registration, + execution_plan=execution_plan, + policy_decision=policy_decision, + approval=approval, + run=run, + authorization=authorization, + ) + + +def _load_json_object(path: Path) -> dict[str, Any]: + value = json.loads(path.read_text(encoding="utf-8")) + if not isinstance(value, dict): + raise ActiveMetadataAuthorizationEvidenceError( + f"{path.name} must contain an object" + ) + return value + + +def build_contract_report() -> dict[str, Any]: + errors: list[str] = [] + paths = { + "authorization_contract": Path(__file__).resolve().parent + / "active_metadata_authorization.py", + "dataset_bundle": Path(__file__).resolve().parent + / "spatial_dataset_bundle.py", + "gateway": Path(__file__).resolve().parent / "platform_gateway.py", + "migration": Path(__file__).resolve().parent + / "migrations/101_active_metadata_authorization.sql", + "rehearsal": Path(__file__).resolve(), + "wrapper": DEFAULT_WRAPPER_PATH, + } + required = { + "authorization_contract": ( + "class MetadataActivationAuthorization", + "build_metadata_activation_authorization", + "Active Metadata dispatch requires approval evidence", + "scheduler_command_enqueued: Literal[True]", + "production_ready: Literal[False]", + ), + "dataset_bundle": ( + "build_shapefile_bundle_inventory", + "validate_shapefile_bundle_inventory", + "content_sha256", + ), + "gateway": ( + "def authorize_metadata_activation(", + "Active Metadata dispatch requires activation authorization", + "metadata_activation_authorization_id", + ), + "migration": ( + "CREATE TABLE IF NOT EXISTS gda_control.metadata_activation_authorization", + "DEFERRABLE INITIALLY DEFERRED", + "authorize_metadata_activation", + "guard_active_metadata_dispatch", + "Active Metadata dispatch requires exact authorization", + "FORCE ROW LEVEL SECURITY", + ), + "rehearsal": ( + "def run_local_rehearsal(", + "orphan_authorization_rollback_verified", + "real_dataset_resource_version_bound", + ), + "wrapper": ( + "data_agent.metadata_fabric_active_metadata_authorization", + '"$@"', + ), + } + files: dict[str, dict[str, Any]] = {} + for name, path in paths.items(): + if not path.is_file(): + errors.append(f"{name} is missing") + files[name] = { + "path": path.resolve().relative_to(REPO_ROOT).as_posix(), + "sha256": None, + } + continue + raw = path.read_bytes() + source = raw.decode("utf-8") + files[name] = { + "path": path.resolve().relative_to(REPO_ROOT).as_posix(), + "sha256": hashlib.sha256(raw).hexdigest(), + } + if any(marker not in source for marker in required[name]): + errors.append(f"{name} is missing required authorization markers") + stable = { + "schema": CONTRACT_SCHEMA, + "activation_route": "metadata_fabric.projection_plan", + "promotion_boundary": "authorization_and_dispatch_same_transaction", + "approval_required": True, + "real_data_role": "acceptance_input_and_resource_version_fingerprint", + "files": files, + "errors": errors, + } + return { + **stable, + "status": "valid" if not errors else "invalid", + "contract_sha256": canonical_json_fingerprint(stable), + "provider_apply_authorized": False, + "provider_mutations_executed": False, + "production_scheduler_submission_verified": False, + "production_ingestion_verified": False, + "production_ready": False, + } + + +def _apply_migrations(engine: Any) -> 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 ActiveMetadataAuthorizationEvidenceError( + "local authorization rehearsal requires a fresh superuser database" + ) + connection.exec_driver_sql( + """ + CREATE TABLE IF NOT EXISTS agent_app_users ( + id SERIAL PRIMARY KEY, + username VARCHAR(100) UNIQUE NOT NULL + ) + """ + ) + for migration in MIGRATIONS: + connection.execute(text(migration.read_text(encoding="utf-8"))) + + +def run_local_rehearsal( + database_url: str, + dataset_inventory: dict[str, Any], +) -> dict[str, Any]: + if validate_shapefile_bundle_inventory(dataset_inventory): + raise ActiveMetadataAuthorizationEvidenceError( + "real dataset bundle inventory is invalid" + ) + bundle = build_authorization_bundle(dataset_inventory["content_sha256"]) + engine = create_engine(database_url) + try: + _apply_migrations(engine) + gateway = PlatformGateway(engine) + gateway.register_resource(bundle.source_resource) + gateway.register_resource_version_with_metadata_event(bundle.registration) + claimed = gateway.claim_metadata_changes( + TENANT, + WORKER, + consumer_subject=CONSUMER_SUBJECT, + ) + gateway.stage_metadata_activation_request( + TENANT, + claimed[0].event.event_id, + worker_id=WORKER, + request=bundle.request, + ) + gateway.register_definition(bundle.definition_registration) + for artifact in ( + bundle.execution_plan, + bundle.policy_decision, + bundle.approval, + ): + gateway.record_artifact(artifact) + gateway.submit_run(bundle.run) + + ordinary_dispatch_blocked = False + try: + gateway.submit_run(bundle.run, request_dispatch=True) + except GatewayValidationError: + ordinary_dispatch_blocked = True + + orphan_authorization_rollback_verified = False + try: + with gateway._transaction(TENANT) as connection: + connection.execute( + text( + """ + SELECT * FROM gda_control.authorize_metadata_activation( + :tenant_id, CAST(:authorization AS jsonb) + ) + """ + ), + { + "tenant_id": TENANT, + "authorization": json.dumps( + bundle.authorization.model_dump( + mode="json", by_alias=True + ), + ensure_ascii=True, + sort_keys=True, + separators=(",", ":"), + ), + }, + ).one() + except GatewayValidationError: + orphan_authorization_rollback_verified = True + + authorization_absent_after_rollback = False + try: + gateway.get_metadata_activation_authorization( + TENANT, bundle.authorization.authorization_id + ) + except GatewayNotFoundError: + authorization_absent_after_rollback = True + + first = gateway.authorize_metadata_activation(bundle.authorization) + replay = gateway.authorize_metadata_activation(bundle.authorization) + stored = gateway.get_metadata_activation_authorization( + TENANT, bundle.authorization.authorization_id + ) + command = gateway.get_command(TENANT, bundle.authorization.command_id) + + direct_mutation_blocked = True + with gateway._transaction(TENANT) as connection: + for statement in ( + """ + UPDATE gda_control.metadata_activation_authorization + SET status = 'authorized_for_dispatch' + WHERE authorization_id = :authorization_id + """, + """ + DELETE FROM gda_control.metadata_activation_authorization + WHERE authorization_id = :authorization_id + """, + ): + try: + with connection.begin_nested(): + connection.execute( + text(statement), + { + "authorization_id": ( + bundle.authorization.authorization_id + ) + }, + ) + except DBAPIError: + continue + direct_mutation_blocked = False + + with engine.connect() as connection: + privileges = connection.exec_driver_sql( + """ + SELECT + has_table_privilege( + 'gda_control_gateway', + 'gda_control.metadata_activation_authorization', 'SELECT' + ), + NOT has_table_privilege( + 'gda_control_gateway', + 'gda_control.metadata_activation_authorization', 'INSERT' + ), + NOT has_table_privilege( + 'gda_control_gateway', + 'gda_control.metadata_activation_authorization', 'UPDATE' + ), + NOT has_table_privilege( + 'gda_control_gateway', + 'gda_control.metadata_activation_authorization', 'DELETE' + ), + has_function_privilege( + 'gda_control_gateway', + 'gda_control.authorize_metadata_activation(text,jsonb)', + 'EXECUTE' + ) + """ + ).one() + force_rls = connection.exec_driver_sql( + """ + SELECT relforcerowsecurity + FROM pg_class + WHERE oid = + 'gda_control.metadata_activation_authorization'::regclass + """ + ).scalar_one() + counts = connection.execute( + text( + """ + SELECT + ( + SELECT count(*) + FROM gda_control.metadata_activation_authorization + WHERE tenant_id = :tenant_id + ) AS authorizations, + ( + SELECT count(*) + FROM gda_control.platform_command_outbox + WHERE tenant_id = :tenant_id + AND command_type = 'dolphinscheduler.dispatch' + ) AS dispatch_commands + """ + ), + {"tenant_id": TENANT}, + ).one() + finally: + engine.dispose() + + real_dataset_resource_version_bound = ( + bundle.registration.resource_version.content_sha256 + == dataset_inventory["content_sha256"] + == stored.content_sha256 + ) + verified = ( + len(claimed) == 1 + and ordinary_dispatch_blocked + and orphan_authorization_rollback_verified + and authorization_absent_after_rollback + and first.created + and not replay.created + and replay.value == stored == bundle.authorization + and counts.authorizations == counts.dispatch_commands == 1 + and command.status.value == "pending" + and command.payload.get("metadata_activation_authorization_id") + == str(bundle.authorization.authorization_id) + and direct_mutation_blocked + and privileges == (True, True, True, True, True) + and bool(force_rls) + and real_dataset_resource_version_bound + ) + contract = build_contract_report() + stable = { + "schema": EVIDENCE_SCHEMA, + "status": ( + "local_real_data_authorization_dispatch_verified" + if verified + else "blocked" + ), + "contract_sha256": contract["contract_sha256"], + "dataset_bundle": dataset_inventory, + "dataset_source_committed": False, + "dataset_absolute_path_committed": False, + "dataset_required_in_ci": False, + "real_dataset_inspected": dataset_inventory.get("spatial_inventory") + is not None, + "real_dataset_resource_version_bound": real_dataset_resource_version_bound, + "request_id": str(bundle.request.request_id), + "resource_version_id": str(SOURCE_ID), + "resource_version_content_sha256": ( + bundle.registration.resource_version.content_sha256 + ), + "definition_version_id": str(DEFINITION_ID), + "run_id": str(RUN_ID), + "execution_plan_artifact_id": str(PLAN_ID), + "policy_decision_artifact_id": str(bundle.policy_decision.artifact_id), + "approval_artifact_id": str(bundle.approval.artifact_id), + "authorization_id": str(bundle.authorization.authorization_id), + "authorization_sha256": bundle.authorization.authorization_sha256, + "command_id": str(bundle.authorization.command_id), + "authorization_count": counts.authorizations, + "dispatch_command_count": counts.dispatch_commands, + "dispatch_command_status": command.status.value, + "ordinary_dispatch_without_activation_authorization_blocked": ( + ordinary_dispatch_blocked + ), + "orphan_authorization_rollback_verified": ( + orphan_authorization_rollback_verified + ), + "authorization_absent_after_rollback": ( + authorization_absent_after_rollback + ), + "exact_authorization_replay_created": replay.created, + "gateway_function_only_insert_verified": privileges + == (True, True, True, True, True), + "direct_authorization_mutation_blocked": direct_mutation_blocked, + "force_rls_verified": bool(force_rls), + "local_postgresql_authorization_dispatch_verified": verified, + "deployment_applied": False, + "production_workload_identity_verified": False, + "provider_apply_authorized": False, + "provider_mutations_executed": False, + "production_scheduler_submission_verified": False, + "production_ingestion_verified": False, + "production_ready": False, + "errors": [] if verified else ["local authorization rehearsal failed"], + } + return {**stable, "evidence_sha256": canonical_json_fingerprint(stable)} + + +def validate_rehearsal_evidence(evidence: dict[str, Any]) -> list[str]: + errors: list[str] = [] + stable = {key: value for key, value in evidence.items() if key != "evidence_sha256"} + if evidence.get("schema") != EVIDENCE_SCHEMA: + errors.append("Active Metadata authorization evidence schema does not match") + if evidence.get("evidence_sha256") != canonical_json_fingerprint(stable): + errors.append("Active Metadata authorization evidence SHA-256 does not match") + contract = build_contract_report() + if evidence.get("contract_sha256") != contract.get("contract_sha256"): + errors.append("Active Metadata authorization contract fingerprint is stale") + dataset = evidence.get("dataset_bundle") + if not isinstance(dataset, dict): + errors.append("Active Metadata authorization dataset bundle is missing") + else: + errors.extend(validate_shapefile_bundle_inventory(dataset)) + if evidence.get("resource_version_content_sha256") != dataset.get( + "content_sha256" + ): + errors.append("real dataset fingerprint is not bound to ResourceVersion") + for claim in ( + "dataset_source_committed", + "dataset_absolute_path_committed", + "dataset_required_in_ci", + "deployment_applied", + "production_workload_identity_verified", + "provider_apply_authorized", + "provider_mutations_executed", + "production_scheduler_submission_verified", + "production_ingestion_verified", + "production_ready", + ): + if evidence.get(claim) is not False: + errors.append(f"local authorization evidence may not claim {claim}") + for claim in ( + "real_dataset_inspected", + "real_dataset_resource_version_bound", + "ordinary_dispatch_without_activation_authorization_blocked", + "orphan_authorization_rollback_verified", + "authorization_absent_after_rollback", + "gateway_function_only_insert_verified", + "direct_authorization_mutation_blocked", + "force_rls_verified", + "local_postgresql_authorization_dispatch_verified", + ): + if evidence.get(claim) is not True: + errors.append(f"Active Metadata authorization did not verify {claim}") + if evidence.get("authorization_count") != 1: + errors.append("authorization evidence must contain one authorization") + if evidence.get("dispatch_command_count") != 1: + errors.append("authorization evidence must contain one dispatch command") + if evidence.get("dispatch_command_status") != "pending": + errors.append("local dispatch command must remain pending") + if evidence.get("exact_authorization_replay_created") is not False: + errors.append("authorization replay must not create a row") + return errors + + +def main(argv: list[str] | None = None) -> int: + parser = argparse.ArgumentParser(description=__doc__) + subparsers = parser.add_subparsers(dest="command", required=True) + validate = subparsers.add_parser("validate") + validate.add_argument("--evidence", type=Path, default=DEFAULT_EVIDENCE_PATH) + rehearse = subparsers.add_parser("rehearse") + rehearse.add_argument("--database-url", required=True) + rehearse.add_argument("--shapefile", type=Path, required=True) + rehearse.add_argument("--source-label", required=True) + rehearse.add_argument("--ogrinfo", type=Path, required=True) + rehearse.add_argument("--proj-data", type=Path) + rehearse.add_argument("--evidence-out", type=Path, required=True) + args = parser.parse_args(argv) + + if args.command == "validate": + report = build_contract_report() + try: + report["errors"].extend( + validate_rehearsal_evidence(_load_json_object(args.evidence)) + ) + except (OSError, ValueError) as exc: + report["errors"].append( + f"Active Metadata authorization evidence is invalid: {type(exc).__name__}" + ) + report["status"] = "valid" if not report["errors"] else "invalid" + print(json.dumps(report, ensure_ascii=True, indent=2, sort_keys=True)) + return 0 if not report["errors"] else 1 + + inventory = build_shapefile_bundle_inventory( + args.shapefile, + source_label=args.source_label, + ogrinfo_path=args.ogrinfo, + proj_data_path=args.proj_data, + ) + evidence = run_local_rehearsal(args.database_url, inventory) + args.evidence_out.write_text( + json.dumps(evidence, ensure_ascii=True, indent=2, sort_keys=True) + "\n", + encoding="utf-8", + ) + print(json.dumps(evidence, ensure_ascii=True, indent=2, sort_keys=True)) + return 0 if not evidence["errors"] else 1 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/data_agent/migrations/101_active_metadata_authorization.sql b/data_agent/migrations/101_active_metadata_authorization.sql new file mode 100644 index 00000000..b1b47228 --- /dev/null +++ b/data_agent/migrations/101_active_metadata_authorization.sql @@ -0,0 +1,540 @@ +-- 101: Evidence-bound promotion of inert Active Metadata activation requests. +-- +-- The authorization ledger is append-only. For the metadata projection +-- capability, a dispatch command can exist only when the exact request, +-- ResourceVersion, DefinitionVersion, accepted Run, execution plan, +-- PolicyDecision and independent Approval are bound in the same transaction. + +CREATE TABLE IF NOT EXISTS gda_control.metadata_activation_authorization ( + tenant_id TEXT NOT NULL, + authorization_id UUID PRIMARY KEY, + request_id UUID NOT NULL, + request_sha256 CHAR(64) NOT NULL, + resource_urn TEXT NOT NULL, + resource_version_id UUID NOT NULL, + content_sha256 CHAR(64) NOT NULL, + definition_version_id UUID NOT NULL, + definition_sha256 CHAR(64) NOT NULL, + run_id UUID NOT NULL, + execution_plan_artifact_id UUID NOT NULL, + execution_plan_sha256 CHAR(64) NOT NULL, + policy_decision_artifact_id UUID NOT NULL, + policy_decision_sha256 CHAR(64) NOT NULL, + approval_artifact_id UUID NOT NULL, + approval_sha256 CHAR(64) NOT NULL, + command_id UUID NOT NULL, + route TEXT NOT NULL, + status TEXT NOT NULL, + authorized_by TEXT NOT NULL, + authorized_at TIMESTAMPTZ NOT NULL, + authorization_document JSONB NOT NULL, + authorization_sha256 CHAR(64) NOT NULL, + created_at TIMESTAMPTZ NOT NULL DEFAULT clock_timestamp(), + CONSTRAINT uq_gda_activation_authorization_tenant_id + UNIQUE (tenant_id, authorization_id), + CONSTRAINT uq_gda_activation_authorization_request + UNIQUE (tenant_id, request_id), + CONSTRAINT uq_gda_activation_authorization_run + UNIQUE (tenant_id, run_id), + CONSTRAINT uq_gda_activation_authorization_command + UNIQUE (tenant_id, command_id), + CONSTRAINT fk_gda_activation_authorization_request + FOREIGN KEY (tenant_id, request_id) + REFERENCES gda_control.metadata_activation_request(tenant_id, request_id), + CONSTRAINT fk_gda_activation_authorization_version + FOREIGN KEY ( + tenant_id, resource_urn, resource_version_id, content_sha256 + ) REFERENCES gda_control.resource_version( + tenant_id, resource_urn, resource_version_id, content_sha256 + ), + CONSTRAINT fk_gda_activation_authorization_definition + FOREIGN KEY (tenant_id, definition_version_id) + REFERENCES gda_control.platform_definition_version( + tenant_id, definition_version_id + ), + CONSTRAINT fk_gda_activation_authorization_run + FOREIGN KEY (tenant_id, run_id) + REFERENCES gda_control.platform_run(tenant_id, run_id), + CONSTRAINT fk_gda_activation_authorization_plan + FOREIGN KEY (tenant_id, execution_plan_artifact_id) + REFERENCES gda_control.artifact(tenant_id, artifact_id), + CONSTRAINT fk_gda_activation_authorization_policy + FOREIGN KEY (tenant_id, policy_decision_artifact_id) + REFERENCES gda_control.artifact(tenant_id, artifact_id), + CONSTRAINT fk_gda_activation_authorization_approval + FOREIGN KEY (tenant_id, approval_artifact_id) + REFERENCES gda_control.artifact(tenant_id, artifact_id), + CONSTRAINT fk_gda_activation_authorization_command + FOREIGN KEY (tenant_id, command_id) + REFERENCES gda_control.platform_command_outbox(tenant_id, command_id) + DEFERRABLE INITIALLY DEFERRED, + CONSTRAINT ck_gda_activation_authorization_hashes CHECK ( + request_sha256 ~ '^[0-9a-f]{64}$' + AND content_sha256 ~ '^[0-9a-f]{64}$' + AND definition_sha256 ~ '^[0-9a-f]{64}$' + AND execution_plan_sha256 ~ '^[0-9a-f]{64}$' + AND policy_decision_sha256 ~ '^[0-9a-f]{64}$' + AND approval_sha256 ~ '^[0-9a-f]{64}$' + AND authorization_sha256 ~ '^[0-9a-f]{64}$' + ), + CONSTRAINT ck_gda_activation_authorization_route CHECK ( + route = 'metadata_fabric.projection_plan' + ), + CONSTRAINT ck_gda_activation_authorization_status CHECK ( + status = 'authorized_for_dispatch' + ), + CONSTRAINT ck_gda_activation_authorizer CHECK ( + authorized_by ~ '^workload:.+' + ), + CONSTRAINT ck_gda_activation_authorization_document CHECK ( + jsonb_typeof(authorization_document) = 'object' + AND authorization_document ?& ARRAY[ + 'schema', 'authorization_id', 'tenant_id', 'request_id', + 'request_sha256', 'resource_urn', 'resource_version_id', + 'content_sha256', 'definition_version_id', 'definition_sha256', + 'run_id', 'execution_plan_artifact_id', 'execution_plan_sha256', + 'policy_decision_artifact_id', 'policy_decision_sha256', + 'approval_artifact_id', 'approval_sha256', 'command_id', 'route', + 'status', 'authorized_by', 'authorized_at', + 'scheduler_command_enqueued', 'provider_apply_authorized', + 'provider_mutations_executed', + 'production_scheduler_submission_verified', + 'production_ingestion_verified', 'production_ready', + 'authorization_sha256' + ] + AND authorization_document - ARRAY[ + 'schema', 'authorization_id', 'tenant_id', 'request_id', + 'request_sha256', 'resource_urn', 'resource_version_id', + 'content_sha256', 'definition_version_id', 'definition_sha256', + 'run_id', 'execution_plan_artifact_id', 'execution_plan_sha256', + 'policy_decision_artifact_id', 'policy_decision_sha256', + 'approval_artifact_id', 'approval_sha256', 'command_id', 'route', + 'status', 'authorized_by', 'authorized_at', + 'scheduler_command_enqueued', 'provider_apply_authorized', + 'provider_mutations_executed', + 'production_scheduler_submission_verified', + 'production_ingestion_verified', 'production_ready', + 'authorization_sha256' + ] = '{}'::jsonb + AND authorization_document->>'schema' = 'gda.metadata_activation_authorization.v1' + AND authorization_document->>'authorization_id' = authorization_id::text + AND authorization_document->>'tenant_id' = tenant_id + AND authorization_document->>'request_id' = request_id::text + AND authorization_document->>'request_sha256' = request_sha256 + AND authorization_document->>'resource_urn' = resource_urn + AND authorization_document->>'resource_version_id' = resource_version_id::text + AND authorization_document->>'content_sha256' = content_sha256 + AND authorization_document->>'definition_version_id' = definition_version_id::text + AND authorization_document->>'definition_sha256' = definition_sha256 + AND authorization_document->>'run_id' = run_id::text + AND authorization_document->>'execution_plan_artifact_id' = + execution_plan_artifact_id::text + AND authorization_document->>'execution_plan_sha256' = execution_plan_sha256 + AND authorization_document->>'policy_decision_artifact_id' = + policy_decision_artifact_id::text + AND authorization_document->>'policy_decision_sha256' = policy_decision_sha256 + AND authorization_document->>'approval_artifact_id' = approval_artifact_id::text + AND authorization_document->>'approval_sha256' = approval_sha256 + AND authorization_document->>'command_id' = command_id::text + AND authorization_document->>'route' = route + AND authorization_document->>'status' = status + AND authorization_document->>'authorized_by' = authorized_by + AND (authorization_document->>'authorized_at')::timestamptz = authorized_at + AND (authorization_document->>'scheduler_command_enqueued')::boolean = true + AND (authorization_document->>'provider_apply_authorized')::boolean = false + AND (authorization_document->>'provider_mutations_executed')::boolean = false + AND ( + authorization_document->>'production_scheduler_submission_verified' + )::boolean = false + AND (authorization_document->>'production_ingestion_verified')::boolean = false + AND (authorization_document->>'production_ready')::boolean = false + AND authorization_document->>'authorization_sha256' = authorization_sha256 + ) +); + +CREATE INDEX IF NOT EXISTS idx_gda_activation_authorization_resource + ON gda_control.metadata_activation_authorization( + tenant_id, resource_version_id, authorized_at DESC + ); + +ALTER TABLE gda_control.metadata_activation_authorization + ENABLE ROW LEVEL SECURITY; +ALTER TABLE gda_control.metadata_activation_authorization + FORCE ROW LEVEL SECURITY; +DROP POLICY IF EXISTS gda_activation_authorization_tenant_isolation + ON gda_control.metadata_activation_authorization; +CREATE POLICY gda_activation_authorization_tenant_isolation + ON gda_control.metadata_activation_authorization + USING (tenant_id = gda_control.current_tenant()) + WITH CHECK (tenant_id = gda_control.current_tenant()); + +CREATE OR REPLACE FUNCTION gda_control.authorize_metadata_activation( + p_tenant_id TEXT, + p_authorization JSONB +) +RETURNS TABLE(activation_authorization JSONB, created BOOLEAN) +LANGUAGE plpgsql +SECURITY DEFINER +SET search_path = pg_catalog, gda_control +SET row_security = on +AS $$ +DECLARE + request_row gda_control.metadata_activation_request%ROWTYPE; + version_row gda_control.resource_version%ROWTYPE; + definition_row gda_control.platform_definition_version%ROWTYPE; + run_row gda_control.platform_run%ROWTYPE; + plan_row gda_control.artifact%ROWTYPE; + policy_row gda_control.artifact%ROWTYPE; + approval_row gda_control.artifact%ROWTYPE; + stored gda_control.metadata_activation_authorization%ROWTYPE; + decision JSONB; + approval JSONB; + expected_resources TEXT[]; + decision_resources TEXT[]; + decision_resource_count INTEGER; + decision_distinct_count INTEGER; + inserted_rows INTEGER := 0; + p_authorized_at TIMESTAMPTZ; +BEGIN + IF gda_control.current_tenant() IS DISTINCT FROM p_tenant_id THEN + RAISE EXCEPTION 'tenant context mismatch' USING ERRCODE = '42501'; + END IF; + IF jsonb_typeof(p_authorization) IS DISTINCT FROM 'object' THEN + RAISE EXCEPTION 'activation authorization must be an object' + USING ERRCODE = '22023'; + END IF; + IF p_authorization->>'tenant_id' IS DISTINCT FROM p_tenant_id + OR p_authorization->>'schema' IS DISTINCT FROM + 'gda.metadata_activation_authorization.v1' + OR p_authorization->>'route' IS DISTINCT FROM + 'metadata_fabric.projection_plan' + OR p_authorization->>'status' IS DISTINCT FROM + 'authorized_for_dispatch' + OR COALESCE((p_authorization->>'scheduler_command_enqueued')::boolean, false) + IS DISTINCT FROM true + OR COALESCE((p_authorization->>'provider_apply_authorized')::boolean, true) + IS DISTINCT FROM false + OR COALESCE((p_authorization->>'provider_mutations_executed')::boolean, true) + IS DISTINCT FROM false + OR COALESCE( + (p_authorization->>'production_scheduler_submission_verified')::boolean, + true + ) IS DISTINCT FROM false + OR COALESCE((p_authorization->>'production_ingestion_verified')::boolean, true) + IS DISTINCT FROM false + OR COALESCE((p_authorization->>'production_ready')::boolean, true) + IS DISTINCT FROM false THEN + RAISE EXCEPTION 'activation authorization safety claims are invalid' + USING ERRCODE = '22023'; + END IF; + IF p_authorization->>'authorized_by' !~ '^workload:.+' THEN + RAISE EXCEPTION 'activation authorizer must use workload identity' + USING ERRCODE = '22023'; + END IF; + p_authorized_at := (p_authorization->>'authorized_at')::timestamptz; + IF p_authorized_at > clock_timestamp() + interval '5 minutes' THEN + RAISE EXCEPTION 'activation authorization time is in the future' + USING ERRCODE = '22023'; + END IF; + + SELECT * INTO request_row + FROM gda_control.metadata_activation_request + WHERE tenant_id = p_tenant_id + AND request_id = (p_authorization->>'request_id')::uuid + FOR UPDATE; + IF NOT FOUND + OR request_row.status <> 'awaiting_authorization' + OR request_row.request_sha256 IS DISTINCT FROM + p_authorization->>'request_sha256' + OR request_row.resource_urn IS DISTINCT FROM + p_authorization->>'resource_urn' + OR request_row.resource_version_id IS DISTINCT FROM + (p_authorization->>'resource_version_id')::uuid + OR request_row.content_sha256 IS DISTINCT FROM + p_authorization->>'content_sha256' + OR request_row.route IS DISTINCT FROM p_authorization->>'route' THEN + RAISE EXCEPTION 'activation authorization does not match the durable request' + USING ERRCODE = '23514'; + END IF; + + SELECT * INTO version_row + FROM gda_control.resource_version + WHERE tenant_id = p_tenant_id + AND resource_version_id = request_row.resource_version_id; + IF NOT FOUND + OR version_row.resource_urn IS DISTINCT FROM request_row.resource_urn + OR version_row.content_sha256 IS DISTINCT FROM request_row.content_sha256 THEN + RAISE EXCEPTION 'activation ResourceVersion binding is invalid' + USING ERRCODE = '23514'; + END IF; + + SELECT * INTO definition_row + FROM gda_control.platform_definition_version + WHERE tenant_id = p_tenant_id + AND definition_version_id = + (p_authorization->>'definition_version_id')::uuid; + IF NOT FOUND + OR definition_row.definition_sha256 IS DISTINCT FROM + p_authorization->>'definition_sha256' + OR definition_row.orchestration_class <> 'dataops' + OR definition_row.capability_id <> 'metadata_fabric.projection_plan' THEN + RAISE EXCEPTION 'activation DefinitionVersion binding is invalid' + USING ERRCODE = '23514'; + END IF; + + SELECT * INTO run_row + FROM gda_control.platform_run + WHERE tenant_id = p_tenant_id + AND run_id = (p_authorization->>'run_id')::uuid + FOR UPDATE; + IF NOT FOUND + OR run_row.definition_version_id IS DISTINCT FROM + definition_row.definition_version_id + OR run_row.orchestration_class <> 'dataops' + OR run_row.status <> 'accepted' + OR run_row.subject_context->>'subject_type' <> 'workload' + OR run_row.submitted_at > p_authorized_at + OR run_row.policy_refs->>'policy_decision_artifact_id' IS DISTINCT FROM + p_authorization->>'policy_decision_artifact_id' + OR run_row.policy_refs->>'approval_artifact_id' IS DISTINCT FROM + p_authorization->>'approval_artifact_id' + OR NOT EXISTS ( + SELECT 1 FROM gda_control.platform_run_input_binding binding + WHERE binding.tenant_id = p_tenant_id + AND binding.run_id = run_row.run_id + AND binding.resource_version_id = request_row.resource_version_id + ) THEN + RAISE EXCEPTION 'activation Run binding is invalid' + USING ERRCODE = '23514'; + END IF; + + SELECT * INTO plan_row + FROM gda_control.artifact + WHERE tenant_id = p_tenant_id + AND artifact_id = + (p_authorization->>'execution_plan_artifact_id')::uuid; + IF NOT FOUND + OR plan_row.artifact_role <> 'execution_plan' + OR plan_row.run_id IS NOT NULL + OR plan_row.resource_version_id IS DISTINCT FROM + definition_row.definition_version_id + OR plan_row.content_sha256 IS DISTINCT FROM + p_authorization->>'execution_plan_sha256' + OR plan_row.created_at > p_authorized_at THEN + RAISE EXCEPTION 'activation execution plan binding is invalid' + USING ERRCODE = '23514'; + END IF; + + SELECT * INTO policy_row + FROM gda_control.artifact + WHERE tenant_id = p_tenant_id + AND artifact_id = + (p_authorization->>'policy_decision_artifact_id')::uuid; + decision := policy_row.manifest->'decision'; + IF NOT FOUND + OR policy_row.artifact_role <> 'evidence' + OR policy_row.run_id IS NOT NULL + OR policy_row.resource_version_id IS DISTINCT FROM + definition_row.definition_version_id + OR policy_row.content_sha256 IS DISTINCT FROM + p_authorization->>'policy_decision_sha256' + OR policy_row.manifest->>'schema' <> + 'gda.policy_decision_artifact.v1' + OR jsonb_typeof(decision) IS DISTINCT FROM 'object' + OR decision->>'tenant_id' IS DISTINCT FROM p_tenant_id + OR decision->>'run_id' IS DISTINCT FROM run_row.run_id::text + OR decision->'subject_context' IS DISTINCT FROM run_row.subject_context + OR decision->>'definition_version_id' IS DISTINCT FROM + definition_row.definition_version_id::text + OR decision->>'execution_plan_artifact_id' IS DISTINCT FROM + plan_row.artifact_id::text + OR decision->>'action' <> 'dolphinscheduler.dispatch' + OR decision->>'effect' <> 'allow' + OR decision->'obligations' <> '[]'::jsonb + OR jsonb_typeof(decision->'resource_version_ids') IS DISTINCT FROM 'array' + OR COALESCE((decision->>'requires_approval')::boolean, false) <> true + OR decision->>'evaluator_subject' !~ '^workload:.+' + OR decision->>'evaluator_subject' = run_row.submitted_by + OR (decision->>'decided_at')::timestamptz > p_authorized_at + OR p_authorized_at >= (decision->>'expires_at')::timestamptz + OR clock_timestamp() >= (decision->>'expires_at')::timestamptz THEN + RAISE EXCEPTION 'activation PolicyDecision binding is invalid' + USING ERRCODE = '23514'; + END IF; + + SELECT array_agg(resource_id ORDER BY resource_id) + INTO expected_resources + FROM ( + SELECT definition_row.definition_version_id::text AS resource_id + UNION + SELECT binding.resource_version_id::text + FROM gda_control.platform_run_input_binding binding + WHERE binding.tenant_id = p_tenant_id + AND binding.run_id = run_row.run_id + ) resources; + SELECT array_agg(value ORDER BY value), count(*), count(DISTINCT value) + INTO decision_resources, decision_resource_count, decision_distinct_count + FROM jsonb_array_elements_text(decision->'resource_version_ids') items(value); + IF decision_resources IS DISTINCT FROM expected_resources + OR decision_resource_count IS DISTINCT FROM decision_distinct_count THEN + RAISE EXCEPTION 'activation PolicyDecision resource scope is invalid' + USING ERRCODE = '23514'; + END IF; + + SELECT * INTO approval_row + FROM gda_control.artifact + WHERE tenant_id = p_tenant_id + AND artifact_id = (p_authorization->>'approval_artifact_id')::uuid; + approval := approval_row.manifest->'approval'; + IF NOT FOUND + OR approval_row.artifact_role <> 'evidence' + OR approval_row.run_id IS NOT NULL + OR approval_row.resource_version_id IS DISTINCT FROM + definition_row.definition_version_id + OR approval_row.content_sha256 IS DISTINCT FROM + p_authorization->>'approval_sha256' + OR approval_row.manifest->>'schema' <> 'gda.approval_artifact.v1' + OR jsonb_typeof(approval) IS DISTINCT FROM 'object' + OR approval->>'tenant_id' IS DISTINCT FROM p_tenant_id + OR approval->>'run_id' IS DISTINCT FROM run_row.run_id::text + OR approval->>'definition_version_id' IS DISTINCT FROM + definition_row.definition_version_id::text + OR approval->>'policy_decision_artifact_id' IS DISTINCT FROM + policy_row.artifact_id::text + OR approval->>'policy_decision_sha256' IS DISTINCT FROM + policy_row.content_sha256 + OR approval->>'verdict' <> 'approved' + OR approval->>'approver_subject' !~ '^human:.+' + OR approval->>'approver_subject' IN ( + run_row.submitted_by, decision->>'evaluator_subject' + ) + OR (approval->>'decided_at')::timestamptz < + (decision->>'decided_at')::timestamptz + OR (approval->>'expires_at')::timestamptz > + (decision->>'expires_at')::timestamptz + OR (approval->>'decided_at')::timestamptz > p_authorized_at + OR p_authorized_at >= (approval->>'expires_at')::timestamptz + OR clock_timestamp() >= (approval->>'expires_at')::timestamptz THEN + RAISE EXCEPTION 'activation Approval binding is invalid' + USING ERRCODE = '23514'; + END IF; + IF p_authorization->>'authorized_by' IN ( + run_row.submitted_by, + decision->>'evaluator_subject', + approval->>'approver_subject' + ) THEN + RAISE EXCEPTION 'activation authorizer is not independent' + USING ERRCODE = '23514'; + END IF; + + INSERT INTO gda_control.metadata_activation_authorization ( + tenant_id, authorization_id, request_id, request_sha256, + resource_urn, resource_version_id, content_sha256, + definition_version_id, definition_sha256, run_id, + execution_plan_artifact_id, execution_plan_sha256, + policy_decision_artifact_id, policy_decision_sha256, + approval_artifact_id, approval_sha256, command_id, route, status, + authorized_by, authorized_at, authorization_document, + authorization_sha256 + ) VALUES ( + p_tenant_id, + (p_authorization->>'authorization_id')::uuid, + request_row.request_id, + request_row.request_sha256, + request_row.resource_urn, + request_row.resource_version_id, + request_row.content_sha256, + definition_row.definition_version_id, + definition_row.definition_sha256, + run_row.run_id, + plan_row.artifact_id, + plan_row.content_sha256, + policy_row.artifact_id, + policy_row.content_sha256, + approval_row.artifact_id, + approval_row.content_sha256, + (p_authorization->>'command_id')::uuid, + p_authorization->>'route', + p_authorization->>'status', + p_authorization->>'authorized_by', + p_authorized_at, + p_authorization, + p_authorization->>'authorization_sha256' + ) + ON CONFLICT DO NOTHING; + GET DIAGNOSTICS inserted_rows = ROW_COUNT; + + SELECT * INTO stored + FROM gda_control.metadata_activation_authorization + WHERE tenant_id = p_tenant_id + AND request_id = request_row.request_id; + IF NOT FOUND OR stored.authorization_document IS DISTINCT FROM p_authorization THEN + RAISE EXCEPTION 'activation authorization identity has different content' + USING ERRCODE = '23505'; + END IF; + RETURN QUERY SELECT stored.authorization_document, inserted_rows = 1; +END; +$$; + +CREATE OR REPLACE FUNCTION gda_control.guard_active_metadata_dispatch() +RETURNS TRIGGER +LANGUAGE plpgsql +SET search_path = pg_catalog, gda_control +AS $$ +DECLARE + capability TEXT; + executor_subject TEXT; + authorization_row gda_control.metadata_activation_authorization%ROWTYPE; +BEGIN + IF NEW.command_type <> 'dolphinscheduler.dispatch' THEN + RETURN NEW; + END IF; + SELECT definition.capability_id, run.submitted_by + INTO capability, executor_subject + FROM gda_control.platform_run run + JOIN gda_control.platform_definition_version definition + ON definition.tenant_id = run.tenant_id + AND definition.definition_version_id = run.definition_version_id + WHERE run.tenant_id = NEW.tenant_id AND run.run_id = NEW.run_id; + IF capability IS DISTINCT FROM 'metadata_fabric.projection_plan' THEN + RETURN NEW; + END IF; + SELECT * INTO authorization_row + FROM gda_control.metadata_activation_authorization + WHERE tenant_id = NEW.tenant_id AND command_id = NEW.command_id; + IF NOT FOUND + OR authorization_row.run_id IS DISTINCT FROM NEW.run_id + OR authorization_row.execution_plan_artifact_id IS DISTINCT FROM + NEW.execution_plan_artifact_id + OR executor_subject IS DISTINCT FROM NEW.actor_subject + OR NEW.payload->>'schema' <> + 'gda.dolphinscheduler_dispatch_command.v1' + OR NEW.payload->>'policy_decision_artifact_id' IS DISTINCT FROM + authorization_row.policy_decision_artifact_id::text + OR NEW.payload->>'metadata_activation_authorization_id' IS DISTINCT FROM + authorization_row.authorization_id::text + OR NEW.payload->>'metadata_activation_request_id' IS DISTINCT FROM + authorization_row.request_id::text THEN + RAISE EXCEPTION 'Active Metadata dispatch requires exact authorization' + USING ERRCODE = '23514'; + END IF; + RETURN NEW; +END; +$$; + +DROP TRIGGER IF EXISTS trg_guard_active_metadata_dispatch + ON gda_control.platform_command_outbox; +CREATE TRIGGER trg_guard_active_metadata_dispatch +BEFORE INSERT ON gda_control.platform_command_outbox +FOR EACH ROW EXECUTE FUNCTION gda_control.guard_active_metadata_dispatch(); + +REVOKE ALL ON TABLE gda_control.metadata_activation_authorization FROM PUBLIC; +REVOKE ALL ON TABLE gda_control.metadata_activation_authorization + FROM gda_control_gateway; +GRANT SELECT ON gda_control.metadata_activation_authorization + TO gda_control_gateway; + +REVOKE ALL ON FUNCTION gda_control.authorize_metadata_activation(text, jsonb) + FROM PUBLIC; +GRANT EXECUTE ON FUNCTION gda_control.authorize_metadata_activation(text, jsonb) + TO gda_control_gateway; diff --git a/data_agent/platform_gateway.py b/data_agent/platform_gateway.py index 886dff67..0cb39d8d 100644 --- a/data_agent/platform_gateway.py +++ b/data_agent/platform_gateway.py @@ -17,7 +17,13 @@ from sqlalchemy import text from sqlalchemy.exc import DBAPIError, SQLAlchemyError +from .active_metadata_authorization import ( + MetadataActivationAuthorization, + MetadataActivationAuthorizationError, + build_metadata_activation_authorization, +) from .active_metadata_change_contract import ( + METADATA_PROJECTION_ROUTE, ActiveMetadataRegistration, MetadataActivationIntent, MetadataActivationRequest, @@ -67,12 +73,11 @@ QualityResult, Resource, ResourceVersion, - RunSuccessEvidence, RunStatus, + RunSuccessEvidence, TenantId, ) - GATEWAY_DATABASE_ROLE = "gda_control_gateway" GATEWAY_SCHEMA_VERSION = "gda.platform_gateway.v1" GATEWAY_ROLE_MIGRATION = ( @@ -110,6 +115,11 @@ / "migrations" / "100_active_metadata_activation_request.sql" ) +ACTIVE_METADATA_AUTHORIZATION_MIGRATION = ( + Path(__file__).resolve().parent + / "migrations" + / "101_active_metadata_authorization.sql" +) USER_TENANT_MIGRATION = ( Path(__file__).resolve().parent / "migrations" @@ -624,23 +634,34 @@ def get_metadata_activation_request( request_id: UUID, ) -> MetadataActivationRequest: with self._transaction(tenant_id) as connection: - row = connection.execute( - text( - """ - SELECT request - FROM gda_control.metadata_activation_request - WHERE tenant_id = :tenant_id AND request_id = :request_id - """ - ), - {"tenant_id": tenant_id, "request_id": request_id}, - ).mappings().one_or_none() - if row is None: + request = self._load_metadata_activation_request( + connection, tenant_id, request_id + ) + if request is None: raise GatewayNotFoundError( "MetadataActivationRequest was not found" ) - return MetadataActivationRequest.model_validate( - _as_json(row["request"]) - ) + return request + + @staticmethod + def _load_metadata_activation_request( + connection, tenant_id: str, request_id: UUID + ) -> MetadataActivationRequest | None: + row = connection.execute( + text( + """ + SELECT request + FROM gda_control.metadata_activation_request + WHERE tenant_id = :tenant_id AND request_id = :request_id + """ + ), + {"tenant_id": tenant_id, "request_id": request_id}, + ).mappings().one_or_none() + return ( + MetadataActivationRequest.model_validate(_as_json(row["request"])) + if row is not None + else None + ) def stage_metadata_activation_request( self, @@ -999,6 +1020,7 @@ def _dispatch_command( run: PlatformRun, decision: PolicyDecision | None, execution_plan: Artifact | None, + activation_authorization: MetadataActivationAuthorization | None = None, ) -> PlatformCommand: if decision is None or execution_plan is None: raise GatewayValidationError( @@ -1019,6 +1041,23 @@ def _dispatch_command( f"{execution_plan.artifact_id}" ) enqueued_at = datetime.now(UTC) + payload = { + "schema": "gda.dolphinscheduler_dispatch_command.v1", + "policy_decision_artifact_id": str( + run.policy_refs.policy_decision_artifact_id + ), + } + if activation_authorization is not None: + payload.update( + { + "metadata_activation_authorization_id": str( + activation_authorization.authorization_id + ), + "metadata_activation_request_id": str( + activation_authorization.request_id + ), + } + ) return PlatformCommand( tenant_id=run.tenant_id, command_id=uuid5(run.run_id, dedupe_key), @@ -1027,12 +1066,7 @@ def _dispatch_command( execution_plan_artifact_id=execution_plan.artifact_id, dedupe_key=dedupe_key, actor_subject=cls._run_actor(run), - payload={ - "schema": "gda.dolphinscheduler_dispatch_command.v1", - "policy_decision_artifact_id": str( - run.policy_refs.policy_decision_artifact_id - ), - }, + payload=payload, available_at=enqueued_at, created_at=enqueued_at, ) @@ -1044,6 +1078,17 @@ def submit_run( decision, execution_plan = self._validate_run_policy_references( connection, run ) + definition = self._load_definition( + connection, run.tenant_id, run.definition_version_id + ) + if ( + request_dispatch + and definition is not None + and definition.capability_id == METADATA_PROJECTION_ROUTE + ): + raise GatewayValidationError( + "Active Metadata dispatch requires activation authorization" + ) actor = self._run_actor(run) inserted = connection.execute( text( @@ -1137,6 +1182,143 @@ def submit_run( ) return GatewayWriteResult(stored, inserted is not None) + @staticmethod + def _load_metadata_activation_authorization( + connection, tenant_id: str, authorization_id: UUID + ) -> MetadataActivationAuthorization | None: + row = connection.execute( + text( + """ + SELECT authorization_document + FROM gda_control.metadata_activation_authorization + WHERE tenant_id = :tenant_id + AND authorization_id = :authorization_id + """ + ), + {"tenant_id": tenant_id, "authorization_id": authorization_id}, + ).mappings().one_or_none() + return ( + MetadataActivationAuthorization.model_validate( + _as_json(row["authorization_document"]) + ) + if row is not None + else None + ) + + def get_metadata_activation_authorization( + self, tenant_id: str, authorization_id: UUID + ) -> MetadataActivationAuthorization: + with self._transaction(tenant_id) as connection: + authorization = self._load_metadata_activation_authorization( + connection, tenant_id, authorization_id + ) + if authorization is None: + raise GatewayNotFoundError( + "MetadataActivationAuthorization was not found" + ) + return authorization + + def authorize_metadata_activation( + self, authorization: MetadataActivationAuthorization + ) -> GatewayWriteResult: + """Atomically append exact authorization and its pending dispatch.""" + with self._transaction(authorization.tenant_id) as connection: + request = self._load_metadata_activation_request( + connection, authorization.tenant_id, authorization.request_id + ) + version = self._load_resource_version( + connection, + authorization.tenant_id, + authorization.resource_version_id, + ) + definition = self._load_definition( + connection, + authorization.tenant_id, + authorization.definition_version_id, + ) + run = self._load_run( + connection, authorization.tenant_id, authorization.run_id + ) + plan = self._load_artifact( + connection, + authorization.tenant_id, + authorization.execution_plan_artifact_id, + ) + policy = self._load_artifact( + connection, + authorization.tenant_id, + authorization.policy_decision_artifact_id, + ) + approval = self._load_artifact( + connection, + authorization.tenant_id, + authorization.approval_artifact_id, + ) + if any( + item is None + for item in ( + request, + version, + definition, + run, + plan, + policy, + approval, + ) + ): + raise GatewayValidationError( + "activation authorization evidence was not found" + ) + try: + expected = build_metadata_activation_authorization( + request, + version, + definition, + run, + plan, + policy, + approval, + authorized_by=authorization.authorized_by, + authorized_at=authorization.authorized_at, + ) + except MetadataActivationAuthorizationError as exc: + raise GatewayValidationError(str(exc)) from exc + if expected != authorization: + raise GatewayValidationError( + "activation authorization does not match stored evidence" + ) + decision = parse_policy_decision_artifact(policy) + command = self._dispatch_command( + run, + decision, + plan, + activation_authorization=authorization, + ) + row = connection.execute( + text( + """ + SELECT * FROM gda_control.authorize_metadata_activation( + :tenant_id, CAST(:authorization AS jsonb) + ) + """ + ), + { + "tenant_id": authorization.tenant_id, + "authorization": _json( + authorization.model_dump(mode="json", by_alias=True) + ), + }, + ).mappings().one() + stored = MetadataActivationAuthorization.model_validate( + _as_json(row["activation_authorization"]) + ) + if stored != authorization: + raise GatewayConflictError( + "stored activation authorization differs from input" + ) + self._put_command(connection, command) + return GatewayWriteResult(stored, bool(row["created"])) + def get_run(self, tenant_id: str, run_id: UUID) -> PlatformRun: with self._transaction(tenant_id) as connection: run = self._load_run(connection, tenant_id, run_id) @@ -2270,6 +2452,7 @@ def build_gateway_report( lineage_migration: Path | None = None, active_metadata_migration: Path | None = None, activation_request_migration: Path | None = None, + activation_authorization_migration: Path | None = None, gateway_source: Path | None = None, routes_source: Path | None = None, command_consumer_source: Path | None = None, @@ -2298,6 +2481,10 @@ def build_gateway_report( activation_request_migration or ACTIVE_METADATA_ACTIVATION_MIGRATION ).resolve(), + "activation_authorization_migration": ( + activation_authorization_migration + or ACTIVE_METADATA_AUTHORIZATION_MIGRATION + ).resolve(), "gateway_source": (gateway_source or Path(__file__)).resolve(), "routes_source": (routes_source or GATEWAY_ROUTES_SOURCE).resolve(), "command_consumer_source": ( @@ -2397,6 +2584,14 @@ def build_gateway_report( "ALTER TABLE gda_control.metadata_activation_request FORCE ROW LEVEL SECURITY", "GRANT SELECT, INSERT ON gda_control.metadata_activation_request", ), + "activation_authorization_migration": ( + "CREATE TABLE IF NOT EXISTS gda_control.metadata_activation_authorization", + "authorize_metadata_activation", + "Active Metadata dispatch requires exact authorization", + "DEFERRABLE INITIALLY DEFERRED", + "FORCE ROW LEVEL SECURITY", + "GRANT SELECT ON gda_control.metadata_activation_authorization", + ), "gateway_source": ( 'SET LOCAL ROLE "{GATEWAY_DATABASE_ROLE}"', "SELECT set_config('app.current_tenant', :tenant, true)", @@ -2419,6 +2614,8 @@ def build_gateway_report( "def fail_metadata_change(", "def get_metadata_activation_request(", "def stage_metadata_activation_request(", + "def authorize_metadata_activation(", + "def get_metadata_activation_authorization(", ), "routes_source": ( 'base = "/api/platform/v1"', @@ -2472,6 +2669,7 @@ def build_gateway_report( or forbidden in texts.get("binding_migration", "") or forbidden in texts.get("lineage_migration", "") or forbidden in texts.get("active_metadata_migration", "") + or forbidden in texts.get("activation_authorization_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 dbd36176..a33653d1 100644 --- a/data_agent/platform_truth.py +++ b/data_agent/platform_truth.py @@ -878,6 +878,26 @@ def _config( ), "Protected authorization and scheduler promotion of durable requests", ), + RuntimeSpec( + "metadata_active_metadata_authorization_rehearsal", + "active_metadata_authorization_rehearsal", + "governed", + "evidence_durable", + "temporary PostgreSQL authorization/dispatch + committed local evidence", + "metadata-platform", + "local_verification_only", + ( + "data_agent/metadata_fabric_active_metadata_authorization.py", + "scripts/metadata-fabric-active-metadata-authorization.sh", + ), + ( + ( + "data_agent/metadata_fabric_active_metadata_authorization.py", + "def run_local_rehearsal", + ), + ), + "Protected workload identity and real scheduler submission/read-back", + ), RuntimeSpec( "datalake_monitor", "monitor_loop", diff --git a/data_agent/spatial_dataset_bundle.py b/data_agent/spatial_dataset_bundle.py new file mode 100644 index 00000000..40d966ca --- /dev/null +++ b/data_agent/spatial_dataset_bundle.py @@ -0,0 +1,153 @@ +"""Path-free content fingerprints and spatial inventory for local datasets.""" + +from __future__ import annotations + +import hashlib +import json +import os +import subprocess +from pathlib import Path +from typing import Any + +from .platform_contracts import canonical_json_fingerprint + +SHAPEFILE_BUNDLE_SCHEMA = "gda.spatial_dataset_bundle.v1" +_REQUIRED_SHAPEFILE_SUFFIXES = frozenset({".shp", ".shx", ".dbf", ".prj"}) + + +class SpatialDatasetBundleError(RuntimeError): + """A local dataset cannot produce a complete, path-free inventory.""" + + +def _sha256(path: Path) -> str: + digest = hashlib.sha256() + with path.open("rb") as handle: + for chunk in iter(lambda: handle.read(1024 * 1024), b""): + digest.update(chunk) + return digest.hexdigest() + + +def _component_suffix(path: Path, stem_name: str) -> str: + suffix = path.name[len(stem_name) :].lower() + if not suffix.startswith("."): + raise SpatialDatasetBundleError("dataset component has no stable suffix") + return suffix + + +def _ogr_inventory( + shapefile_path: Path, + *, + ogrinfo_path: Path, + proj_data_path: Path | None, +) -> dict[str, Any]: + environment = os.environ.copy() + if proj_data_path is not None: + if not (proj_data_path / "proj.db").is_file(): + raise SpatialDatasetBundleError("PROJ data directory has no proj.db") + environment["PROJ_DATA"] = str(proj_data_path) + result = subprocess.run( + [str(ogrinfo_path), "-json", "-so", str(shapefile_path)], + check=False, + capture_output=True, + text=True, + env=environment, + ) + if result.returncode != 0: + raise SpatialDatasetBundleError("ogrinfo could not inspect the dataset") + try: + payload = json.loads(result.stdout) + layer = payload["layers"][0] + geometry = layer["geometryFields"][0] + coordinate_system = geometry["coordinateSystem"]["projjson"] + except (KeyError, IndexError, TypeError, ValueError) as exc: + raise SpatialDatasetBundleError("ogrinfo JSON is incomplete") from exc + authority = coordinate_system.get("id") or {} + fields = layer.get("fields") or [] + return { + "driver": payload.get("driverShortName"), + "geometry_type": geometry.get("type"), + "feature_count": layer.get("featureCount"), + "field_count": len(fields), + "crs": { + "authority": authority.get("authority"), + "code": authority.get("code"), + "name": coordinate_system.get("name"), + }, + "bounds": geometry.get("extent"), + } + + +def build_shapefile_bundle_inventory( + shapefile_path: Path, + *, + source_label: str, + ogrinfo_path: Path | None = None, + proj_data_path: Path | None = None, +) -> dict[str, Any]: + """Hash every same-stem sidecar and omit all source filesystem paths.""" + shapefile_path = shapefile_path.resolve() + if not shapefile_path.is_file() or shapefile_path.suffix.lower() != ".shp": + raise SpatialDatasetBundleError("a readable .shp file is required") + if not source_label.strip() or "/" in source_label or "\\" in source_label: + raise SpatialDatasetBundleError("source label must be non-path text") + stem_name = shapefile_path.stem + components = sorted( + ( + path + for path in shapefile_path.parent.iterdir() + if path.is_file() and path.name.startswith(f"{stem_name}.") + ), + key=lambda item: item.name.lower(), + ) + suffixes = {_component_suffix(path, stem_name) for path in components} + missing = sorted(_REQUIRED_SHAPEFILE_SUFFIXES - suffixes) + if missing: + raise SpatialDatasetBundleError( + f"required shapefile components are missing: {', '.join(missing)}" + ) + component_inventory = [ + { + "component": _component_suffix(path, stem_name), + "size_bytes": path.stat().st_size, + "sha256": _sha256(path), + } + for path in components + ] + stable: dict[str, Any] = { + "schema": SHAPEFILE_BUNDLE_SCHEMA, + "source_label": source_label.strip(), + "format": "ESRI Shapefile", + "components": component_inventory, + "spatial_inventory": ( + _ogr_inventory( + shapefile_path, + ogrinfo_path=ogrinfo_path, + proj_data_path=proj_data_path, + ) + if ogrinfo_path is not None + else None + ), + } + return {**stable, "content_sha256": canonical_json_fingerprint(stable)} + + +def validate_shapefile_bundle_inventory(inventory: dict[str, Any]) -> list[str]: + """Validate checked evidence without requiring the local source dataset.""" + errors: list[str] = [] + stable = {key: value for key, value in inventory.items() if key != "content_sha256"} + if inventory.get("schema") != SHAPEFILE_BUNDLE_SCHEMA: + errors.append("spatial dataset bundle schema does not match") + if inventory.get("content_sha256") != canonical_json_fingerprint(stable): + errors.append("spatial dataset bundle SHA-256 does not match") + if any( + "/" in str(value) or "\\" in str(value) + for value in ( + inventory.get("source_label", ""), + *(item.get("component", "") for item in inventory.get("components", [])), + ) + ): + errors.append("spatial dataset bundle must not contain source paths") + suffixes = {item.get("component") for item in inventory.get("components", [])} + if not _REQUIRED_SHAPEFILE_SUFFIXES.issubset(suffixes): + errors.append("spatial dataset bundle is missing required components") + return errors diff --git a/data_agent/test_active_metadata_authorization.py b/data_agent/test_active_metadata_authorization.py new file mode 100644 index 00000000..10b066ac --- /dev/null +++ b/data_agent/test_active_metadata_authorization.py @@ -0,0 +1,265 @@ +from datetime import UTC, datetime, timedelta +from uuid import UUID + +import pytest +from pydantic import ValidationError + +from data_agent.active_metadata_authorization import ( + MetadataActivationAuthorization, + MetadataActivationAuthorizationError, + build_metadata_activation_authorization, + dispatch_command_id, +) +from data_agent.active_metadata_change_contract import ( + build_active_metadata_registration, + build_metadata_activation_intent, + build_metadata_activation_request, +) +from data_agent.platform_authorization import ( + build_approval_artifact, + build_policy_decision_artifact, +) +from data_agent.platform_contracts import ( + ApprovalRecord, + Artifact, + PlatformDefinitionVersion, + PlatformRun, + PolicyDecision, + ResourceVersion, + RunPolicyReferences, + SubjectContext, + canonical_json_fingerprint, + platform_definition_fingerprint, +) + +TENANT = "metadata-auth" +SOURCE_ID = UUID("a5000000-0000-4000-8000-000000000001") +DEFINITION_ID = UUID("a5000000-0000-4000-8000-000000000002") +RUN_ID = UUID("a5000000-0000-4000-8000-000000000003") +PLAN_ID = UUID("a5000000-0000-4000-8000-000000000004") +NOW = datetime(2026, 7, 30, 8, 0, tzinfo=UTC) + + +def _source() -> ResourceVersion: + return ResourceVersion( + tenant_id=TENANT, + resource_urn=f"gda://{TENANT}/dataset/cultural-districts", + resource_version_id=SOURCE_ID, + version_key="bundle-v1", + content_sha256="a" * 64, + authority_version_ref={"bundle": "local-acceptance"}, + created_by="workload:metadata-registrar", + created_at=NOW - timedelta(hours=2), + ) + + +def _request(): + registration = build_active_metadata_registration( + _source(), consumer_subject="workload:active-metadata-consumer" + ) + intent = build_metadata_activation_intent( + registration.event, + routed_by="workload:active-metadata-consumer", + ) + return build_metadata_activation_request(intent) + + +def _definition() -> PlatformDefinitionVersion: + document = {"tasks": ["project-governance-metadata"]} + input_contract = {"metadata_change": "dataset"} + output_contract = {"projection_plan": "artifact"} + fingerprint = platform_definition_fingerprint( + orchestration_class="dataops", + capability_id="metadata_fabric.projection_plan", + portability_class="portable", + definition_document=document, + input_contract=input_contract, + output_contract=output_contract, + ) + return PlatformDefinitionVersion( + tenant_id=TENANT, + definition_urn=f"gda://{TENANT}/definition/metadata-projection", + definition_version_id=DEFINITION_ID, + orchestration_class="dataops", + capability_id="metadata_fabric.projection_plan", + portability_class="portable", + definition_document=document, + input_contract=input_contract, + output_contract=output_contract, + definition_sha256=fingerprint, + ) + + +def _plan() -> Artifact: + manifest = {"schema": "gda.metadata_projection_execution_plan.v1"} + return Artifact( + tenant_id=TENANT, + artifact_id=PLAN_ID, + artifact_key="metadata-projection-plan", + artifact_role="execution_plan", + storage_uri=f"postgresql://gda-control/execution-plans/{TENANT}/{PLAN_ID}", + media_type="application/vnd.gda.metadata-projection-plan+json", + content_sha256=canonical_json_fingerprint(manifest), + size_bytes=len(b'{"schema":"gda.metadata_projection_execution_plan.v1"}'), + resource_version_id=DEFINITION_ID, + manifest=manifest, + created_by="workload:metadata-plan-compiler", + created_at=NOW - timedelta(minutes=30), + ) + + +def _chain(): + subject = SubjectContext( + tenant_id=TENANT, + subject_id="metadata-projection-runner", + subject_type="workload", + roles=("metadata_projector",), + purpose="project active metadata change", + ) + provisional = PlatformRun( + tenant_id=TENANT, + run_id=RUN_ID, + definition_version_id=DEFINITION_ID, + orchestration_class="dataops", + subject_context=subject, + input_bindings=( + { + "binding_name": "metadata_change", + "resource_version_id": SOURCE_ID, + "semantic_type": "gis.cultural_districts", + }, + ), + idempotency_key="metadata-projection:cultural-districts:v1", + submitted_at=NOW - timedelta(minutes=20), + ) + decision = PolicyDecision( + tenant_id=TENANT, + run_id=RUN_ID, + subject_context=subject, + action="dolphinscheduler.dispatch", + definition_version_id=DEFINITION_ID, + resource_version_ids=(DEFINITION_ID, SOURCE_ID), + execution_plan_artifact_id=PLAN_ID, + effect="allow", + policy_version_ref=f"gda://{TENANT}/policy/metadata-dispatch-v1", + evaluator_subject="workload:metadata-policy-evaluator", + requires_approval=True, + decided_at=NOW - timedelta(minutes=15), + expires_at=NOW + timedelta(hours=1), + ) + policy_artifact = build_policy_decision_artifact(decision) + approval = ApprovalRecord( + tenant_id=TENANT, + run_id=RUN_ID, + definition_version_id=DEFINITION_ID, + policy_decision_artifact_id=policy_artifact.artifact_id, + policy_decision_sha256=policy_artifact.content_sha256, + verdict="approved", + approver_subject="human:metadata-governance-approver", + reason="approved bounded metadata projection", + decided_at=NOW - timedelta(minutes=10), + expires_at=NOW + timedelta(minutes=45), + ) + approval_artifact = build_approval_artifact(approval) + run = provisional.model_copy( + update={ + "policy_refs": RunPolicyReferences( + policy_decision_artifact_id=policy_artifact.artifact_id, + approval_artifact_id=approval_artifact.artifact_id, + ) + } + ) + return run, policy_artifact, approval_artifact + + +def _authorization(**overrides): + run, policy, approval = _chain() + values = { + "request": _request(), + "resource_version": _source(), + "definition": _definition(), + "run": run, + "execution_plan_artifact": _plan(), + "policy_decision_artifact": policy, + "approval_artifact": approval, + "authorized_by": "workload:metadata-activation-authorizer", + "authorized_at": NOW, + } + values.update(overrides) + return build_metadata_activation_authorization(**values) + + +def test_authorization_binds_complete_chain_with_stable_identity(): + first = _authorization() + second = _authorization() + + assert first == second + assert first.command_id == dispatch_command_id(RUN_ID, PLAN_ID) + assert first.scheduler_command_enqueued is True + assert first.provider_apply_authorized is False + assert first.production_scheduler_submission_verified is False + + +def test_authorization_rejects_unbound_source_definition_and_authorizer(): + with pytest.raises( + MetadataActivationAuthorizationError, match="ResourceVersion" + ): + _authorization( + resource_version=_source().model_copy( + update={"content_sha256": "b" * 64} + ) + ) + + with pytest.raises( + MetadataActivationAuthorizationError, match="DefinitionVersion" + ): + _authorization( + definition=_definition().model_copy( + update={"capability_id": "metadata_fabric.apply"} + ) + ) + + with pytest.raises(MetadataActivationAuthorizationError, match="independent"): + _authorization(authorized_by="workload:metadata-policy-evaluator") + + +def test_authorization_requires_active_independent_approval(): + run, policy, approval = _chain() + expired_approval = build_approval_artifact( + ApprovalRecord.model_validate( + approval.manifest["approval"] + | { + "decided_at": (NOW - timedelta(minutes=10)).isoformat(), + "expires_at": (NOW - timedelta(minutes=1)).isoformat(), + } + ) + ) + run = run.model_copy( + update={ + "policy_refs": RunPolicyReferences( + policy_decision_artifact_id=policy.artifact_id, + approval_artifact_id=expired_approval.artifact_id, + ) + } + ) + with pytest.raises(MetadataActivationAuthorizationError, match="not active"): + build_metadata_activation_authorization( + _request(), + _source(), + _definition(), + run, + _plan(), + policy, + expired_approval, + authorized_by="workload:metadata-activation-authorizer", + authorized_at=NOW, + ) + + +def test_authorization_document_tampering_fails_closed(): + authorization = _authorization() + payload = authorization.model_dump(mode="json", by_alias=True) + payload["content_sha256"] = "f" * 64 + + with pytest.raises(ValidationError, match="authorization"): + MetadataActivationAuthorization.model_validate(payload) diff --git a/data_agent/test_active_metadata_authorization_postgres.py b/data_agent/test_active_metadata_authorization_postgres.py new file mode 100644 index 00000000..a690b0d8 --- /dev/null +++ b/data_agent/test_active_metadata_authorization_postgres.py @@ -0,0 +1,87 @@ +import os +from pathlib import Path +from uuid import uuid4 + +import pytest +from sqlalchemy import create_engine, text +from sqlalchemy.engine import make_url + +from data_agent.metadata_fabric_active_metadata_authorization import ( + run_local_rehearsal, +) +from data_agent.spatial_dataset_bundle import build_shapefile_bundle_inventory + +DATABASE_URL = os.environ.get("DATABASE_URL") + + +def _temporary_database_url() -> tuple[object, str, str]: + admin_url = make_url(DATABASE_URL) + admin_engine = create_engine(admin_url, isolation_level="AUTOCOMMIT") + with admin_engine.connect() as connection: + is_superuser = connection.exec_driver_sql( + "SELECT rolsuper FROM pg_roles WHERE rolname = current_user" + ).scalar_one() + if not is_superuser: + admin_engine.dispose() + pytest.skip("Active Metadata authorization test requires a superuser") + database_name = f"gda_active_authorization_{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() + + +def _inventory(tmp_path: Path) -> dict: + stem = tmp_path / "districts" + for suffix, content in ( + (".shp", b"shape"), + (".shx", b"index"), + (".dbf", b"attributes"), + (".prj", b"crs"), + (".cpg", b"UTF-8"), + ): + stem.with_suffix(suffix).write_bytes(content) + return build_shapefile_bundle_inventory( + stem.with_suffix(".shp"), source_label="postgres-golden-slice" + ) + + +@pytest.mark.skipif(not DATABASE_URL, reason="DATABASE_URL is not configured") +def test_authorization_and_dispatch_are_atomic_fail_closed_and_idempotent(tmp_path): + admin_engine, database_name, database_url = _temporary_database_url() + try: + evidence = run_local_rehearsal(database_url, _inventory(tmp_path)) + + assert evidence["local_postgresql_authorization_dispatch_verified"] is True + assert evidence["ordinary_dispatch_without_activation_authorization_blocked"] + assert evidence["orphan_authorization_rollback_verified"] is True + assert evidence["authorization_absent_after_rollback"] is True + assert evidence["authorization_count"] == 1 + assert evidence["dispatch_command_count"] == 1 + assert evidence["dispatch_command_status"] == "pending" + assert evidence["exact_authorization_replay_created"] is False + assert evidence["gateway_function_only_insert_verified"] is True + assert evidence["direct_authorization_mutation_blocked"] is True + assert evidence["force_rls_verified"] is True + assert evidence["provider_apply_authorized"] is False + assert evidence["production_scheduler_submission_verified"] is False + finally: + _drop_temporary_database(admin_engine, database_name) diff --git a/data_agent/test_metadata_fabric_active_metadata_authorization.py b/data_agent/test_metadata_fabric_active_metadata_authorization.py new file mode 100644 index 00000000..6e61ba2e --- /dev/null +++ b/data_agent/test_metadata_fabric_active_metadata_authorization.py @@ -0,0 +1,78 @@ +import json +from copy import deepcopy + +from data_agent import metadata_fabric_active_metadata_authorization as authorization + + +def test_static_contract_requires_atomic_evidence_bound_dispatch(): + report = authorization.build_contract_report() + + assert report["status"] == "valid" + assert report["errors"] == [] + assert report["approval_required"] is True + assert report["promotion_boundary"] == ( + "authorization_and_dispatch_same_transaction" + ) + assert report["real_data_role"] == ( + "acceptance_input_and_resource_version_fingerprint" + ) + assert report["provider_apply_authorized"] is False + assert report["production_scheduler_submission_verified"] is False + assert report["production_ready"] is False + assert all( + not item["path"].startswith("/") for item in report["files"].values() + ) + + +def test_checked_real_data_evidence_is_current_path_free_and_fail_closed(): + evidence = json.loads( + authorization.DEFAULT_EVIDENCE_PATH.read_text(encoding="utf-8") + ) + + assert authorization.validate_rehearsal_evidence(evidence) == [] + assert evidence["real_dataset_inspected"] is True + assert evidence["dataset_bundle"]["spatial_inventory"] == { + "bounds": [ + 106.37987914500007, + 29.558008447000077, + 106.59532712300008, + 29.877271985000025, + ], + "crs": { + "authority": "EPSG", + "code": 4490, + "name": "China Geodetic Coordinate System 2000", + }, + "driver": "ESRI Shapefile", + "feature_count": 20, + "field_count": 33, + "geometry_type": "PolygonZ", + } + assert evidence["resource_version_content_sha256"] == ( + evidence["dataset_bundle"]["content_sha256"] + ) + assert evidence["dataset_source_committed"] is False + assert evidence["dataset_absolute_path_committed"] is False + assert evidence["dataset_required_in_ci"] is False + assert "/Users/" not in json.dumps(evidence) + assert evidence["authorization_count"] == 1 + assert evidence["dispatch_command_count"] == 1 + assert evidence["dispatch_command_status"] == "pending" + assert evidence["provider_apply_authorized"] is False + assert evidence["production_scheduler_submission_verified"] is False + assert evidence["production_ready"] is False + + +def test_evidence_validation_rejects_tampering_and_production_overclaim(): + evidence = json.loads( + authorization.DEFAULT_EVIDENCE_PATH.read_text(encoding="utf-8") + ) + tampered = deepcopy(evidence) + tampered["dispatch_command_count"] = 2 + tampered["production_ready"] = True + + errors = authorization.validate_rehearsal_evidence(tampered) + + assert "Active Metadata authorization evidence SHA-256 does not match" in errors + assert "local authorization evidence may not claim production_ready" in errors + assert "authorization evidence must contain one dispatch command" in errors diff --git a/data_agent/test_platform_contracts.py b/data_agent/test_platform_contracts.py index fb6eae6c..1711e9cc 100644 --- a/data_agent/test_platform_contracts.py +++ b/data_agent/test_platform_contracts.py @@ -472,7 +472,7 @@ def test_control_ledger_contract_and_migration_catalog_are_valid(): assert report["contract_count"] == 16 assert report["migration"]["sha256"] == migration["checksum"] assert migrations[-1]["migration_id"] == ( - "100_active_metadata_activation_request" + "101_active_metadata_authorization" ) diff --git a/data_agent/test_platform_truth.py b/data_agent/test_platform_truth.py index 5a130abd..f4e2c1ad 100644 --- a/data_agent/test_platform_truth.py +++ b/data_agent/test_platform_truth.py @@ -257,6 +257,11 @@ def test_repository_source_access_and_runtime_baselines_match(): and item["production_role"] == "local_verification_only" for item in static_report["runtime"]["inventory"] ) + assert any( + item["runtime_id"] == "metadata_active_metadata_authorization_rehearsal" + and item["production_role"] == "local_verification_only" + for item in static_report["runtime"]["inventory"] + ) def test_runtime_report_detects_unregistered_background_mechanism(tmp_path): diff --git a/data_agent/test_spatial_dataset_bundle.py b/data_agent/test_spatial_dataset_bundle.py new file mode 100644 index 00000000..181cbc62 --- /dev/null +++ b/data_agent/test_spatial_dataset_bundle.py @@ -0,0 +1,51 @@ +from copy import deepcopy + +from data_agent.spatial_dataset_bundle import ( + build_shapefile_bundle_inventory, + validate_shapefile_bundle_inventory, +) + + +def test_shapefile_bundle_hashes_all_sidecars_without_source_paths(tmp_path): + stem = tmp_path / "districts" + for suffix, content in ( + (".shp", b"shape"), + (".shx", b"index"), + (".dbf", b"attributes"), + (".prj", b"crs"), + (".cpg", b"UTF-8"), + (".shp.xml", b"metadata"), + ): + stem.with_suffix(suffix).write_bytes(content) + + inventory = build_shapefile_bundle_inventory( + stem.with_suffix(".shp"), + source_label="chongqing-cultural-districts", + ) + + assert validate_shapefile_bundle_inventory(inventory) == [] + assert inventory["spatial_inventory"] is None + assert {item["component"] for item in inventory["components"]} == { + ".cpg", + ".dbf", + ".prj", + ".shp", + ".shp.xml", + ".shx", + } + assert str(tmp_path) not in str(inventory) + + +def test_checked_bundle_validation_detects_tampering(tmp_path): + stem = tmp_path / "districts" + for suffix in (".shp", ".shx", ".dbf", ".prj"): + stem.with_suffix(suffix).write_bytes(suffix.encode()) + inventory = build_shapefile_bundle_inventory( + stem.with_suffix(".shp"), source_label="golden-slice" + ) + tampered = deepcopy(inventory) + tampered["components"][0]["size_bytes"] += 1 + + assert "spatial dataset bundle SHA-256 does not match" in ( + validate_shapefile_bundle_inventory(tampered) + ) diff --git a/docs/architecture-decisions/adr-060-transactional-active-metadata-change-outbox.md b/docs/architecture-decisions/adr-060-transactional-active-metadata-change-outbox.md index 2f5a691d..d3fb95af 100644 --- a/docs/architecture-decisions/adr-060-transactional-active-metadata-change-outbox.md +++ b/docs/architecture-decisions/adr-060-transactional-active-metadata-change-outbox.md @@ -54,9 +54,9 @@ consumer scoping, wrong-worker rejection, retry, lease-expiry reclaim, processed finality, tenant isolation, forced RLS, direct update/delete denial, and rollback of a legacy-version event attempt. It records one authoritative event, three delivery attempts, contract fingerprint -`b6162ed687c6554859b0056e34dd5cf9f051818deb533daf708eeea67b18eca5`, +`c3d94228456aff7e9b134fa6bc746bbe6b7485950c16f8ec41ec34fc7a5ae567`, and evidence fingerprint -`807444bb1adb2ea633bf98d845ba8a62fee608a4c6eff18c17c0c94f406e8a2`. +`2b8a408e078cec44fde9a6d63e4b94f988dc820668c87ac0d2ef0a434d1a16a3`. This is local transactional evidence only and does not establish production Active Metadata readiness. diff --git a/docs/architecture-decisions/adr-061-durable-inert-active-metadata-activation-request.md b/docs/architecture-decisions/adr-061-durable-inert-active-metadata-activation-request.md index 68d7e68c..0e923e82 100644 --- a/docs/architecture-decisions/adr-061-durable-inert-active-metadata-activation-request.md +++ b/docs/architecture-decisions/adr-061-durable-inert-active-metadata-activation-request.md @@ -86,9 +86,9 @@ scheduler secret, disabled Kubernetes token mounting, and admission of the consumer selector to PostgreSQL ingress. - Contract fingerprint: - `7a0bc0e3aaa53509443f031e7c40d7c5ff0f304510ddfeceedbad00147b505f5` + `2256daa97c1f3a2e71f4d7026592171daea802b371f5616fbbd72a63939ee6b5` - Evidence fingerprint: - `480a2d651c1389feed471cb6ff082537ddfcd2f108e2652b9168b16ea4d02977` + `029aaf7de476115dcf6385ca4a0e05bb84492ebea8ddecf6c3edf36edd76dbef` `deployment_applied`, `production_workload_identity_verified`, `production_scheduler_submission_verified`, `production_ingestion_verified`, diff --git a/docs/architecture-decisions/adr-062-atomic-active-metadata-authorization-and-dispatch.md b/docs/architecture-decisions/adr-062-atomic-active-metadata-authorization-and-dispatch.md new file mode 100644 index 00000000..381ff991 --- /dev/null +++ b/docs/architecture-decisions/adr-062-atomic-active-metadata-authorization-and-dispatch.md @@ -0,0 +1,115 @@ +# ADR-062: Atomic Active Metadata authorization and dispatch + +- Status: Accepted for local AR-1 verification +- Date: 2026-07-30 +- Owners: Metadata Platform / Data Platform / Security +- Related decisions: ADR-024, ADR-025, ADR-027, ADR-047, ADR-060, ADR-061 + +## Context + +ADR-061 deliberately stops an Active Metadata change at an immutable +`awaiting_authorization` request. That request binds the changed +ResourceVersion but is not a DefinitionVersion, Run, execution plan, +PolicyDecision, Approval, or scheduler command. + +The next boundary must prevent two unsafe partial states: a dispatch command +without complete authorization, and an authorization row committed without +its dispatch command. It must also distinguish realistic acceptance input from +authority. Local Chongqing files can prove content identity, spatial metadata, +and conformance, but their presence cannot authorize execution or establish +production provenance. + +## Options considered + +| Option | Benefit | Cost / risk | Decision | +|---|---|---|---| +| Reuse ordinary `submit_run(request_dispatch=true)` | No new ledger | The activation request is not bound and can be bypassed | Rejected | +| Write authorization, then enqueue in a later transaction | Simple components | Leaves orphan authorization and retry ambiguity | Rejected | +| Append authorization and dispatch in one transaction with a database guard | Exact replay, fail-closed evidence chain, no partial commit | Adds migration, deferred FK, trigger, and dedicated gateway API | Accepted | + +## Decision + +Migration 101 adds tenant-scoped, append-only +`gda_control.metadata_activation_authorization`. Its security-definer function +accepts only a content-bound `MetadataActivationAuthorization` and validates +the exact durable request, ResourceVersion/content hash, DataOps +DefinitionVersion with capability `metadata_fabric.projection_plan`, accepted +workload Run and input binding, execution-plan Artifact, allow PolicyDecision, +and approved independent ApprovalRecord. The authorizer is a fourth independent +workload identity and must differ from executor, evaluator, and approver. + +The gateway role has SELECT but no direct INSERT/UPDATE/DELETE on this ledger. +The authorization row has a deferred foreign key to its command. The dedicated +gateway transaction appends the authorization first, then inserts one pending +DolphinScheduler dispatch. A command trigger rejects activation dispatch when +the exact authorization is absent or its request/policy/plan/payload binding +differs. Calling the authorization function without inserting the command +therefore fails at commit and rolls back the authorization. + +Ordinary `submit_run(request_dispatch=true)` explicitly rejects the metadata +projection capability. Other DataOps capabilities retain the existing generic +command path. Exact replay returns the existing authorization and command; +conflicting identity or content fails closed. + +Creation of a pending command proves local durable enqueue only. It does not +prove DolphinScheduler submission, provider apply authorization, provider +mutation, production ingestion, or production readiness. All those claims +remain `false`. + +## Real data boundary + +The local golden slice is the Chongqing central-city historical cultural +district Shapefile bundle. Every same-stem component, including spatial index +and XML sidecars, is hashed. Checked evidence retains only component suffixes, +sizes, SHA-256 values, format, feature/field counts, geometry type, CRS, bounds, +and the canonical bundle fingerprint. It contains no absolute source path and +does not commit source data. + +The bundle fingerprint becomes the source ResourceVersion `content_sha256`. +This makes real data useful as acceptance input and immutable content evidence, +while policy and approval remain the only authorization evidence. CI validates +the checked inventory and evidence fingerprints without requiring the local +dataset. DEM and TAP remain candidates for later raster and scale/time-series +conformance; TAP is not policy-outcome authority. + +## Consequences and trade-offs + +- A dedicated authorizer identity and approval are mandatory, increasing + operational setup but preserving separation of duties. +- The Run and all evidence Artifacts must exist before promotion; malformed or + expired evidence cannot be partially accepted. +- Deferred constraints make authorization and enqueue atomic, but actual + scheduler delivery remains the existing leased outbox worker's responsibility. +- One request maps to one authorization, Run, and dispatch command. A future + re-plan requires a new immutable activation request rather than rewriting the + ledger. +- Local files improve realism but do not establish license, vintage, + authoritative custody, protected provenance, or production availability. + +## Local verification + +The checked PostgreSQL rehearsal proves ordinary dispatch rejection, rollback +of authorization without a command, exact authorization/command replay, forced +RLS, function-only INSERT, direct UPDATE/DELETE denial, one authorization and +one pending dispatch. The real bundle contains 20 `PolygonZ` features and 33 +fields in EPSG:4490, with the recorded Chongqing central-city bounds. + +- Dataset bundle fingerprint: + `fd474fd65c8e4a71da241eb3fd07748ca3b972fbd2d3c32833376dbe71104007` +- Contract fingerprint: + `cef78f91058a8529f4e86330790e714b52b73725a45ffe4dc9eded35bc8ccfa4` +- Evidence fingerprint: + `6ae387240e3bcebaafe2ad7acc73f4e09d53df2e73b2ec63cd92edbc262d831e` + +`deployment_applied`, `production_workload_identity_verified`, +`provider_apply_authorized`, `provider_mutations_executed`, +`production_scheduler_submission_verified`, `production_ingestion_verified`, +and `production_ready` remain `false`. + +## Revisit triggers + +Revisit when a protected production authorization service can provide an +equivalent atomic append/enqueue proof, when one request legitimately needs +multiple independently governed Runs, when policy obligations gain an +executable and audited implementation, or when a real scheduler submission and +provider read-back slice is ready to replace the local pending-command proof. diff --git a/docs/evidence/metadata-fabric-active-metadata-authorization-2026-07-30.json b/docs/evidence/metadata-fabric-active-metadata-authorization-2026-07-30.json new file mode 100644 index 00000000..52ee98fd --- /dev/null +++ b/docs/evidence/metadata-fabric-active-metadata-authorization-2026-07-30.json @@ -0,0 +1,106 @@ +{ + "approval_artifact_id": "f3e9fac1-d67f-5943-942f-c92a7a7e2f63", + "authorization_absent_after_rollback": true, + "authorization_count": 1, + "authorization_id": "8c2738e2-848c-5069-ac89-d47057a1bff9", + "authorization_sha256": "6cdf1def214d503da39a786cd92bef3a2875e10bd4157eeb8d52ba10d8e97f28", + "command_id": "bea71c2e-ac58-52df-8942-8740c8590df0", + "contract_sha256": "cef78f91058a8529f4e86330790e714b52b73725a45ffe4dc9eded35bc8ccfa4", + "dataset_absolute_path_committed": false, + "dataset_bundle": { + "components": [ + { + "component": ".cpg", + "sha256": "3ad3031f5503a4404af825262ee8232cc04d4ea6683d42c5dd0a2f2a27ac9824", + "size_bytes": 5 + }, + { + "component": ".dbf", + "sha256": "ee7c6c4c6957aea296b69d62118d416e5ee989aa77f7b98cf0fe580874ce5127", + "size_bytes": 44990 + }, + { + "component": ".prj", + "sha256": "b10dbe4d6d1de908d340f892c90b3d31a552630af3742bb515bfe1bd26124f2c", + "size_bytes": 176 + }, + { + "component": ".sbn", + "sha256": "7d0279465b18beec40308717e0ef0ea5701a586bc5c84e6a9aa309d5bc0ec99a", + "size_bytes": 308 + }, + { + "component": ".sbx", + "sha256": "019156149b2c7771ec0dd249c757dd7aa01a98075e246a77f77b080860c57333", + "size_bytes": 124 + }, + { + "component": ".shp", + "sha256": "6ac0d5c8c8db66fc0e2a74d8232b7779bd2454257df14efa2930e3dbc181aed0", + "size_bytes": 283640 + }, + { + "component": ".shp.xml", + "sha256": "8ef222ce1952552b366acf14a996e1c8cbdbe3eed0bb829dacfe7eafd068d948", + "size_bytes": 43100 + }, + { + "component": ".shx", + "sha256": "f3fbb6a7775ca833c066e3a3f2a332f99f979840045909ae187d08dac126a119", + "size_bytes": 260 + } + ], + "content_sha256": "fd474fd65c8e4a71da241eb3fd07748ca3b972fbd2d3c32833376dbe71104007", + "format": "ESRI Shapefile", + "schema": "gda.spatial_dataset_bundle.v1", + "source_label": "chongqing-central-cultural-districts", + "spatial_inventory": { + "bounds": [ + 106.37987914500007, + 29.558008447000077, + 106.59532712300008, + 29.877271985000025 + ], + "crs": { + "authority": "EPSG", + "code": 4490, + "name": "China Geodetic Coordinate System 2000" + }, + "driver": "ESRI Shapefile", + "feature_count": 20, + "field_count": 33, + "geometry_type": "PolygonZ" + } + }, + "dataset_required_in_ci": false, + "dataset_source_committed": false, + "definition_version_id": "a6000000-0000-4000-8000-000000000002", + "deployment_applied": false, + "direct_authorization_mutation_blocked": true, + "dispatch_command_count": 1, + "dispatch_command_status": "pending", + "errors": [], + "evidence_sha256": "6ae387240e3bcebaafe2ad7acc73f4e09d53df2e73b2ec63cd92edbc262d831e", + "exact_authorization_replay_created": false, + "execution_plan_artifact_id": "a6000000-0000-4000-8000-000000000004", + "force_rls_verified": true, + "gateway_function_only_insert_verified": true, + "local_postgresql_authorization_dispatch_verified": true, + "ordinary_dispatch_without_activation_authorization_blocked": true, + "orphan_authorization_rollback_verified": true, + "policy_decision_artifact_id": "3b978b60-6e4c-5b0f-b885-a87243233d9b", + "production_ingestion_verified": false, + "production_ready": false, + "production_scheduler_submission_verified": false, + "production_workload_identity_verified": false, + "provider_apply_authorized": false, + "provider_mutations_executed": false, + "real_dataset_inspected": true, + "real_dataset_resource_version_bound": true, + "request_id": "e52b84fc-b39d-534b-9445-e627836fd607", + "resource_version_content_sha256": "fd474fd65c8e4a71da241eb3fd07748ca3b972fbd2d3c32833376dbe71104007", + "resource_version_id": "a6000000-0000-4000-8000-000000000001", + "run_id": "a6000000-0000-4000-8000-000000000003", + "schema": "gda.active_metadata_authorization_evidence.v1", + "status": "local_real_data_authorization_dispatch_verified" +} diff --git a/docs/evidence/metadata-fabric-active-metadata-consumer-2026-07-30.json b/docs/evidence/metadata-fabric-active-metadata-consumer-2026-07-30.json index 5b2c63ac..894f6ffc 100644 --- a/docs/evidence/metadata-fabric-active-metadata-consumer-2026-07-30.json +++ b/docs/evidence/metadata-fabric-active-metadata-consumer-2026-07-30.json @@ -2,7 +2,7 @@ "activation_request_count": 2, "activation_route": "metadata_fabric.projection_plan", "atomic_completion_guard_verified": true, - "contract_sha256": "7a0bc0e3aaa53509443f031e7c40d7c5ff0f304510ddfeceedbad00147b505f5", + "contract_sha256": "2256daa97c1f3a2e71f4d7026592171daea802b371f5616fbbd72a63939ee6b5", "cross_tenant_read_blocked": true, "deployment_applied": false, "deployment_contract_verified": true, @@ -13,7 +13,7 @@ "3a7303a0-06a8-5f60-aa16-1abe23a050d5", "52073ce1-db1f-526c-a6a7-51976c2ec8cd" ], - "evidence_sha256": "480a2d651c1389feed471cb6ff082537ddfcd2f108e2652b9168b16ea4d02977", + "evidence_sha256": "029aaf7de476115dcf6385ca4a0e05bb84492ebea8ddecf6c3edf36edd76dbef", "exact_request_replay_created": false, "force_rls_verified": true, "gateway_select_insert_only_verified": true, diff --git a/docs/evidence/metadata-fabric-active-metadata-outbox-2026-07-30.json b/docs/evidence/metadata-fabric-active-metadata-outbox-2026-07-30.json index 22f57fb7..b5438283 100644 --- a/docs/evidence/metadata-fabric-active-metadata-outbox-2026-07-30.json +++ b/docs/evidence/metadata-fabric-active-metadata-outbox-2026-07-30.json @@ -2,12 +2,12 @@ "activation_intent_sha256": "0380e458ac706d6e1f66f0a320177f1ea7956c15a93faa654d8e01323258fce9", "activation_route": "metadata_fabric.projection_plan", "authoritative_event_count": 1, - "contract_sha256": "b6162ed687c6554859b0056e34dd5cf9f051818deb533daf708eeea67b18eca5", + "contract_sha256": "c3d94228456aff7e9b134fa6bc746bbe6b7485950c16f8ec41ec34fc7a5ae567", "cross_tenant_read_blocked": true, "errors": [], "event_id": "52073ce1-db1f-526c-a6a7-51976c2ec8cd", "event_sha256": "08087307bdd6694a3b2de2176fdb02ab7db720728b65b6da641e0b25dfa5edcf", - "evidence_sha256": "807444bb1adb2ea633bf98d845ba8a62fee608a4c6eff18c17c0c94f406e8a2c", + "evidence_sha256": "2b8a408e078cec44fde9a6d63e4b94f988dc820668c87ac0d2ef0a434d1a16a3", "exact_replay_created": false, "final_attempt_count": 3, "first_registration_created": true, diff --git a/docs/roadmap-ar0-platform-truth-2026-07-24.md b/docs/roadmap-ar0-platform-truth-2026-07-24.md index 654e612b..90091163 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-15(本地 Active Metadata durable request consumer 已验证,生产验证待执行) +### 4.8 Metadata Fabric Bridge M1 + M2 + M3-16(本地真实数据授权/dispatch 原子提升已验证,生产验证待执行) 第八块回到 AR-1 的 metadata control plane,以 [ADR-036](architecture-decisions/adr-036-read-only-metadata-fabric-bridge-contract.md) 固定 OpenMetadata + Gravitino + GDA Control Ledger 的首条 table slice: @@ -233,10 +233,11 @@ Temporal 继续保持目标组件状态,不在这一包并行接入。OpenMeta 28. [ADR-057](architecture-decisions/adr-057-production-object-store-readiness-gate.md) 已将 M3-10 evidence、S3-compatible provider/account/region/bucket、独立 failure domain、OIDC workload federation、精确八项 S3 permission、TLS/private path、KMS、versioning、cross-region replication、strong read/list consistency、tenant isolation、owner/SLO/runbook 和 26 项 protected attestation check 冻结为 fail-closed profile。当前 profile fingerprint 为 `668e194b3c688307014148391e7f389c9d6e9ca69c95d7b4cc92b4acae93181a`,report fingerprint 为 `85362dd10b7dc565f9fa567673d90b774cdec714bd1e70fb2c3c83c1af48b5ea`,合同有效但 43 项生产输入仍 blocked,全部 production claims 为 `false`。这只是 provider-neutral 决策和验收合同:没有选择、部署或验证 AWS S3、华为云 OBS 或其他生产对象存储;原生非 S3 provider 必须进入新的 conformance slice。 29. [ADR-058](architecture-decisions/adr-058-local-spark-commit-failure-recovery.md) 已在 Spark driver 的 loopback Iceberg REST proxy 中于 provider 转发前注入 HTTP 503。baseline 为 1 个 append snapshot、2 行和 1 个 referenced Parquet;失败调用经过精确 2 次 503 后,snapshot/row/file 均零漂移;对同一 `spark-recovery` 行做一次显式重试后为父子相连的 2 个 append snapshots、3 行和 2 个 referenced Parquet。直接 MinIO inventory 精确为 2 data + 3 metadata + 4 manifest = 9 objects,没有孤儿 data file;namespace、两块 PV 和 port-forward 均清理。contract fingerprint 为 `6d8944ab80246dc65891aa81118cb8b73f7ecad699be9a2af5e62d8260c41002`,evidence fingerprint 为 `39571cdac1e4043bcfc2d03a73b2b12ff925210daf8ae36bc640b8cb14d89401`。该结果只证明已知 pre-forward 失败下的本地原子性和一次显式重试,不证明 uncertain commit reconciliation、网络 exactly-once、生产对象存储或完整 engine conformance。 30. [ADR-059](architecture-decisions/adr-059-local-spark-uncertain-commit-reconciliation.md) 已将一个 armed commit 转发给 Gravitino,并在 provider 返回 200 后丢弃成功响应、向 Spark 返回 Iceberg `CommitStateUnknownException` 所需的 HTTP 504;一次传输重试被抑制。Spark 不重提逻辑写,而是 readback 得到父子相连的 2 个 append snapshots、3 行和 2 个 referenced Parquet,决策为 `committed_do_not_resubmit`、`write_resubmitted=false`。MinIO inventory 为 2 data + 3 metadata + 4 manifest = 9 objects;Job `Complete 1/1`,namespace、两块 PV 和 port-forward 均清理。contract fingerprint 为 `7a8d75a1d6b4558b982c6c3242d8d356c5046955f8aae7a45e5c297b6f4d4132`,evidence fingerprint 为 `d6462fff78d07047311b1f715d5f2c7f08c0ce8fbdd5c8b26a3d95ddc3474786`。该结果只证明一个本地 append 的确定性 readback/no-resubmit,不证明持久 reconcile controller、并发写、进程崩溃恢复、网络 exactly-once 或生产能力。 -31. [ADR-060](architecture-decisions/adr-060-transactional-active-metadata-change-outbox.md) 已新增 migration 099、内容绑定 `MetadataChangeEvent`、deterministic activation intent 与 PlatformGateway 原子注册/claim/fail/complete API。真实 PostgreSQL 16 演练中,ResourceVersion 与事件同事务创建,精确 replay 与 processed replay 均不新增事件;错误 consumer/worker 被拒绝,一次 retry 和一次强制租约过期后由第三个 worker 完成,旧 ResourceVersion 的补事件尝试整笔回滚。最终只有 1 条权威事件、3 次 attempt,FORCE RLS、跨租户拒绝与 gateway 无直接 UPDATE/DELETE 均通过。因共享 contract/gateway 源码演进后已在 fresh database 重跑,当前 contract fingerprint 为 `b6162ed687c6554859b0056e34dd5cf9f051818deb533daf708eeea67b18eca5`,evidence fingerprint 为 `807444bb1adb2ea633bf98d845ba8a62fee608a4c6eff18c17c0c94f406e8a2`。该切片只生成 `metadata_fabric.projection_plan` 意图,不新增常驻 consumer、不提交 DolphinScheduler、不授权或执行 provider mutation,也不证明 production ingestion。 -32. [ADR-061](architecture-decisions/adr-061-durable-inert-active-metadata-activation-request.md) 拒绝把缺少 Definition、Run、execution plan、PolicyDecision 与 Approval 的 metadata event 直接转换为 DolphinScheduler command。migration 100 新增 tenant-scoped `MetadataActivationRequest`;managed consumer 只持 PostgreSQL 权限,并在同一事务写入 `awaiting_authorization` request 与完成 event。真实 PostgreSQL 演练得到 2 个 processed events、2 个精确 durable requests 和 0 个 `platform_command_outbox` rows;旧 no-request complete 被阻断,精确 request replay 不新增行,跨租户、FORCE RLS、直接 UPDATE/DELETE 拒绝均通过。base Deployment 为 0 replicas、无 provider/scheduler Secret、禁用 Kubernetes token mount。contract fingerprint 为 `7a0bc0e3aaa53509443f031e7c40d7c5ff0f304510ddfeceedbad00147b505f5`,evidence fingerprint 为 `480a2d651c1389feed471cb6ff082537ddfcd2f108e2652b9168b16ea4d02977`;deployment 未 apply,production workload identity、scheduler submission、provider mutation/ingestion 与 readiness 全部仍为 `false`。 +31. [ADR-060](architecture-decisions/adr-060-transactional-active-metadata-change-outbox.md) 已新增 migration 099、内容绑定 `MetadataChangeEvent`、deterministic activation intent 与 PlatformGateway 原子注册/claim/fail/complete API。真实 PostgreSQL 16 演练中,ResourceVersion 与事件同事务创建,精确 replay 与 processed replay 均不新增事件;错误 consumer/worker 被拒绝,一次 retry 和一次强制租约过期后由第三个 worker 完成,旧 ResourceVersion 的补事件尝试整笔回滚。最终只有 1 条权威事件、3 次 attempt,FORCE RLS、跨租户拒绝与 gateway 无直接 UPDATE/DELETE 均通过。因共享 contract/gateway 源码演进后已在 fresh database 重跑,当前 contract fingerprint 为 `c3d94228456aff7e9b134fa6bc746bbe6b7485950c16f8ec41ec34fc7a5ae567`,evidence fingerprint 为 `2b8a408e078cec44fde9a6d63e4b94f988dc820668c87ac0d2ef0a434d1a16a3`。该切片只生成 `metadata_fabric.projection_plan` 意图,不新增常驻 consumer、不提交 DolphinScheduler、不授权或执行 provider mutation,也不证明 production ingestion。 +32. [ADR-061](architecture-decisions/adr-061-durable-inert-active-metadata-activation-request.md) 拒绝把缺少 Definition、Run、execution plan、PolicyDecision 与 Approval 的 metadata event 直接转换为 DolphinScheduler command。migration 100 新增 tenant-scoped `MetadataActivationRequest`;managed consumer 只持 PostgreSQL 权限,并在同一事务写入 `awaiting_authorization` request 与完成 event。真实 PostgreSQL 演练得到 2 个 processed events、2 个精确 durable requests 和 0 个 `platform_command_outbox` rows;旧 no-request complete 被阻断,精确 request replay 不新增行,跨租户、FORCE RLS、直接 UPDATE/DELETE 拒绝均通过。base Deployment 为 0 replicas、无 provider/scheduler Secret、禁用 Kubernetes token mount。contract fingerprint 为 `2256daa97c1f3a2e71f4d7026592171daea802b371f5616fbbd72a63939ee6b5`,evidence fingerprint 为 `029aaf7de476115dcf6385ca4a0e05bb84492ebea8ddecf6c3edf36edd76dbef`;deployment 未 apply,production workload identity、scheduler submission、provider mutation/ingestion 与 readiness 全部仍为 `false`。 +33. [ADR-062](architecture-decisions/adr-062-atomic-active-metadata-authorization-and-dispatch.md) 新增 migration 101、内容绑定 `MetadataActivationAuthorization` 与专用 PlatformGateway 提升 API。`awaiting_authorization` request 只有在同租户真实 ResourceVersion/content hash、`metadata_fabric.projection_plan` DefinitionVersion、accepted workload Run/input、execution-plan Artifact、allow PolicyDecision、独立 approved Approval 与第四方 authorizer 完整匹配时,才能与一个 pending DolphinScheduler dispatch 同事务提交;普通 `request_dispatch` 绕过和无 command 的孤立授权均回滚。真实重庆中心城区历史文化街区 Shapefile 8 组件被规范化为不含路径的 inventory,20 个 `PolygonZ`、33 字段、EPSG:4490,其 bundle SHA `fd474fd65c8e4a71da241eb3fd07748ca3b972fbd2d3c32833376dbe71104007` 精确成为 ResourceVersion content hash。PostgreSQL 演练最终只有 1 个 authorization、1 个 pending command,精确 replay 不新增,FORCE RLS、function-only INSERT、直接 UPDATE/DELETE 拒绝均通过。contract fingerprint 为 `cef78f91058a8529f4e86330790e714b52b73725a45ffe4dc9eded35bc8ccfa4`,evidence fingerprint 为 `6ae387240e3bcebaafe2ad7acc73f4e09d53df2e73b2ec63cd92edbc262d831e`;源数据/绝对路径不入 Git、CI 不依赖本机文件,scheduler submission、provider apply/mutation/ingestion 与 production readiness 仍为 `false`。 -此处 M1 只证明静态合同和只读 HTTP 边界;M2a 只证明本地 live foundation 与 PVC 重挂载连续性;M2b-1/M2b-2 分别限定在同集群新 PVC 和同集群隔离 repository;M2b-3 的 `local_cross_cluster_recovery_verified=true` 只限定在 `local_same_host_distinct_kubernetes_clusters_external_s3_repository`;M2c-1/M2c-2/M2c-3 分别限定本地 provider metrics、临时双周期 OTel 和单 job scrape recovery;M2c-4/M2d-2 只证明 production observability/NetworkPolicy profile 与 attestation 合同可校验;M2d-1 只证明本地两节点 kindnet 的隔离合成流量;M3-1 的 terminal evidence 与 M3-2 的 PolicyDecision/Approval 仍是 deterministic local fixtures。M3-2 只把 projection 写入本地 provider 并证明 retained target 的单次零写入 replay;M3-3 只把该本地 evidence 对应的 binding 写入临时 GDA Control 账本;M3-4 只向无认证 loopback receiver 发送精确 candidate 并验证 503 后幂等恢复;M3-5 只证明 OpenMetadata 在 provider 强制默认 role 之上的项目新增 grant 限定为 `table/Create`,以及本地 JWT 轮换/吊销和越权拒绝;M3-6 只证明隔离 Gravitino Basic IdP 的 bounded table-create、catalog-create 拒绝、登录轮换/吊销和完整清理;M3-7 只证明 pending production identity profile、profile-bound attestation 和派生 claim 的 fail-closed 合同可校验,没有部署或证明真实身份路径;M3-8 只证明同一 Docker Desktop 集群内 Basic 用户、JDBC metadata 与 file warehouse PVC 在受控 Pod restart 后连续;M3-9 只证明同节点共享 RWO PVC 的 Spark interoperability;M3-10 移除了该共享 PVC,并证明同一 Docker Desktop 主机/集群内 Spark 与 MinIO 的跨节点 S3-compatible 互操作,但不证明生产云对象存储、独立 failure domain、持久 identity binding、Flink 或完整 engine conformance;M3-11 只冻结 provider-neutral production object-store profile、精确 attestation binding 与 fail-closed claims,没有选择 provider、部署 bucket/KMS/policy 或提交真实 attestation;M3-12 只证明同一本地路径的 pre-forward commit failure 不改变可见 table state,随后一次显式重试产生一个新 snapshot/row,且无孤儿 data file;M3-13 只证明单次本地 append 在 provider 200 响应丢失并映射为 commit-state-unknown 后,可以由即时 table readback 判定 committed 且不重提,不覆盖持久 controller、进程崩溃、并发写或任意 mutation;M3-14 只证明 ResourceVersion 注册与 Active Metadata 事件在本地 PostgreSQL 同事务创建,并验证租户/workload scoped claim/retry/complete;M3-15 只证明默认零副本 managed consumer 的代码/部署边界,以及本地 PostgreSQL 中 inert activation request 与 event completion 的原子性,不包含已部署常驻 consumer、受保护 workload identity、DolphinScheduler submission、provider policy/approval 或 provider mutation。M3-2 ingestion 仍使用 bootstrap admin,生产持久 binding、ResourceVersion 和 legacy authority 都未写入;生产对象存储、双 provider/生产最小权限、protected workload identity、OIDC、TLS、生产持久 catalog、tenant isolation、真实 receiver/alert/SLO、受保护 provider policy、生产故障注入、source-loss recovery、cancel/reconcile/lineage、完整 Spark/Flink conformance、生产 ingest、四项 production gate 和 `production_ready` 仍为 `false`。 +此处 M1 只证明静态合同和只读 HTTP 边界;M2a 只证明本地 live foundation 与 PVC 重挂载连续性;M2b-1/M2b-2 分别限定在同集群新 PVC 和同集群隔离 repository;M2b-3 的 `local_cross_cluster_recovery_verified=true` 只限定在 `local_same_host_distinct_kubernetes_clusters_external_s3_repository`;M2c-1/M2c-2/M2c-3 分别限定本地 provider metrics、临时双周期 OTel 和单 job scrape recovery;M2c-4/M2d-2 只证明 production observability/NetworkPolicy profile 与 attestation 合同可校验;M2d-1 只证明本地两节点 kindnet 的隔离合成流量;M3-1 的 terminal evidence 与 M3-2 的 PolicyDecision/Approval 仍是 deterministic local fixtures。M3-2 只把 projection 写入本地 provider 并证明 retained target 的单次零写入 replay;M3-3 只把该本地 evidence 对应的 binding 写入临时 GDA Control 账本;M3-4 只向无认证 loopback receiver 发送精确 candidate 并验证 503 后幂等恢复;M3-5 只证明 OpenMetadata 在 provider 强制默认 role 之上的项目新增 grant 限定为 `table/Create`,以及本地 JWT 轮换/吊销和越权拒绝;M3-6 只证明隔离 Gravitino Basic IdP 的 bounded table-create、catalog-create 拒绝、登录轮换/吊销和完整清理;M3-7 只证明 pending production identity profile、profile-bound attestation 和派生 claim 的 fail-closed 合同可校验,没有部署或证明真实身份路径;M3-8 只证明同一 Docker Desktop 集群内 Basic 用户、JDBC metadata 与 file warehouse PVC 在受控 Pod restart 后连续;M3-9 只证明同节点共享 RWO PVC 的 Spark interoperability;M3-10 移除了该共享 PVC,并证明同一 Docker Desktop 主机/集群内 Spark 与 MinIO 的跨节点 S3-compatible 互操作,但不证明生产云对象存储、独立 failure domain、持久 identity binding、Flink 或完整 engine conformance;M3-11 只冻结 provider-neutral production object-store profile、精确 attestation binding 与 fail-closed claims,没有选择 provider、部署 bucket/KMS/policy 或提交真实 attestation;M3-12 只证明同一本地路径的 pre-forward commit failure 不改变可见 table state,随后一次显式重试产生一个新 snapshot/row,且无孤儿 data file;M3-13 只证明单次本地 append 在 provider 200 响应丢失并映射为 commit-state-unknown 后,可以由即时 table readback 判定 committed 且不重提,不覆盖持久 controller、进程崩溃、并发写或任意 mutation;M3-14 只证明 ResourceVersion 注册与 Active Metadata 事件在本地 PostgreSQL 同事务创建,并验证租户/workload scoped claim/retry/complete;M3-15 只证明默认零副本 managed consumer 的代码/部署边界,以及本地 PostgreSQL 中 inert activation request 与 event completion 的原子性;M3-16 只证明本地真实数据 content fingerprint、证据绑定授权与 pending command 的 PostgreSQL 原子性,不包含受保护 workload identity、已部署 authorization controller、真实 DolphinScheduler submission/read-back 或 provider mutation。M3-2 ingestion 仍使用 bootstrap admin,生产持久 binding、ResourceVersion 和 legacy authority 都未写入;生产对象存储、双 provider/生产最小权限、protected workload identity、OIDC、TLS、生产持久 catalog、tenant isolation、真实 receiver/alert/SLO、受保护 provider policy、生产故障注入、source-loss recovery、cancel/reconcile/lineage、完整 Spark/Flink conformance、生产 ingest、四项 production gate 和 `production_ready` 仍为 `false`。 ## 5. 重新评估条件 diff --git a/docs/system-of-record-matrix-2026-07-24.md b/docs/system-of-record-matrix-2026-07-24.md index 055a8657..2fcb368a 100644 --- a/docs/system-of-record-matrix-2026-07-24.md +++ b/docs/system-of-record-matrix-2026-07-24.md @@ -2,9 +2,9 @@ 日期:2026-07-30 -阶段:AR-0 `in_progress`;AR-1 gateway、成功终局 evidence gate、DolphinScheduler adapter sandbox POC、Metadata Fabric M1/M2、M2c-4/M2d-2 production readiness contracts、M3-1/M3-2、M3-3 local binding ledger、M3-4 local OpenLineage wire delivery、M3-5 local OpenMetadata bounded identity、M3-6 local Gravitino Basic bounded identity、M3-7 production identity readiness contract、M3-8 local Gravitino JDBC restart continuity、M3-9 local Spark/Iceberg REST interoperability、M3-10 local cross-node Spark/object-store interoperability、M3-11 production object-store readiness contract、M3-12 local Spark commit-failure recovery、M3-13 local uncertain-commit reconciliation、M3-14 local Active Metadata transactional outbox 与 M3-15 local durable activation request consumer 已验证,生产 provider ingestion、生产观测、生产 policy/tenant isolation、生产 identity/object-store attestation、生产 consumer/scheduler 和生产切换仍 `in_progress` +阶段:AR-0 `in_progress`;AR-1 gateway、成功终局 evidence gate、DolphinScheduler adapter sandbox POC、Metadata Fabric M1/M2、M2c-4/M2d-2 production readiness contracts、M3-1/M3-2、M3-3 local binding ledger、M3-4 local OpenLineage wire delivery、M3-5 local OpenMetadata bounded identity、M3-6 local Gravitino Basic bounded identity、M3-7 production identity readiness contract、M3-8 local Gravitino JDBC restart continuity、M3-9 local Spark/Iceberg REST interoperability、M3-10 local cross-node Spark/object-store interoperability、M3-11 production object-store readiness contract、M3-12 local Spark commit-failure recovery、M3-13 local uncertain-commit reconciliation、M3-14 local Active Metadata transactional outbox、M3-15 local durable activation request consumer 与 M3-16 local real-data authorization/dispatch promotion 已验证,生产 provider ingestion、生产观测、生产 policy/tenant isolation、生产 identity/object-store attestation、生产 consumer/scheduler 和生产切换仍 `in_progress` -适用分支:`feat/ar1-metadata-fabric-active-metadata-consumer` +适用分支:`feat/ar1-metadata-fabric-active-metadata-authorization` ## 判定规则 @@ -20,7 +20,7 @@ | SQL schema 历史 | PostgreSQL `schema_migrations`,以完整 migration ID + checksum 为权威 | migration CLI 的 JSON 报告 | 保持现有 ledger;任何 drift fail closed | Data Platform | AR-0,已验证 | | 部署配置策略 | Compose/K8s/进程环境;`platform_truth.CONFIG_SPECS` 定义关键类型与策略;DolphinScheduler worker 与 Active Metadata consumer 均有默认零副本、外部 ConfigMap/Secret 驱动的 Kustomize 模板和静态 validator,前者另有 staging activation preflight | `.env` 仅补默认;脱敏 snapshot、Secret key attestation、未扩容 Deployment 和 `ready_for_activation` 都是观测/模板 | 版本化 DeploymentProfile + secret reference;部署环境始终优先;模板或 preflight 通过都不等于环境已启用 | Platform/SRE/Security | AR-0,部分实现;worker 模板/preflight 本地已验证 | | 环境发布与晋级 | 本地 candidate/registry/provenance/release/live 合同已绑定 publisher、verifier、OCI 和 manifest identity;canonical `main@0182406`、archive refs、三组 active ruleset 与 `staging-provenance` protected environment 已建立,但尚无成功 publisher/verifier 或 deployment | 旧 mainline、feature branch、CI artifact、JSON、离线 report 和合成 `verified_for_staging_apply` 都不能单独成为发布权威;publisher SHA、verifier SHA 与 branch lineage 必须分别验证 | 由受保护 environment 的 DeploymentRevision 绑定 OCI、provenance artifact、release manifest 与全部 live verdict | Platform/SRE/Security/Repository Owner | AR-1 mainline 治理已恢复 -> 首次 GHCR publish/verify -> 真实 staging | -| 后台运行时清单 | `platform_truth.RUNTIME_INVENTORY` 是代码层登记;`gda_control` 已有受控 PlatformRun 写入口;DolphinScheduler managed worker 已登记但尚无生产调用方;Active Metadata consumer 登记为 `activation_request_staging_only`,其 deployment 默认为 0 replicas;Metadata Fabric recovery/metrics/policy/catalog/interoperability/failure/uncertain-commit/outbox/consumer rehearsal 均为 `local_verification_only` | AST primitive report、worker status JSON、FrameworkAttemptObservation、DolphinScheduler instance state、本地 recovery/metrics/network-policy/catalog/interoperability/failure/outbox/consumer evidence | PlatformRun ledger 唯一登记最终状态;activation request 只拥有待授权意图,不是 Run、command 或 provider mutation 权威;本地演练进程与 evidence 不得变成生产控制器、监控后端、catalog authority 或 tenant-isolation 权威 | Platform Architecture | AR-1 adapter/worker 与 M3-15 consumer 本地验证;metadata runner 仅本地验证 -> staging 控制链待接入 | +| 后台运行时清单 | `platform_truth.RUNTIME_INVENTORY` 是代码层登记;`gda_control` 已有受控 PlatformRun 写入口;DolphinScheduler managed worker 已登记但尚无生产调用方;Active Metadata consumer 登记为 `activation_request_staging_only`,其 deployment 默认为 0 replicas;Metadata Fabric recovery/metrics/policy/catalog/interoperability/failure/uncertain-commit/outbox/consumer/authorization rehearsal 均为 `local_verification_only` | AST primitive report、worker status JSON、FrameworkAttemptObservation、DolphinScheduler instance state、本地 recovery/metrics/network-policy/catalog/interoperability/failure/outbox/consumer/authorization evidence | PlatformRun ledger 唯一登记最终状态;activation request 只拥有待授权意图;M3-16 authorization 只授权本地 pending dispatch,不是 provider mutation 或生产 scheduler submission 权威;本地演练进程与 evidence 不得变成生产控制器、监控后端、catalog authority 或 tenant-isolation 权威 | Platform Architecture | AR-1 adapter/worker、M3-15 consumer 与 M3-16 authorization 本地验证;metadata runner 仅本地验证 -> staging 控制链待接入 | | 原始文件/对象 | 当前 local uploads、S3/MinIO/OBS 均可能被直接写入,权威边界未统一 | 临时上传、下载缓存、预览文件 | Landing object 以 immutable URI + checksum + retention 为权威;本地 scratch 可删除 | Data Platform | AR-2 | | 湖仓表与 snapshot | Iceberg/STAC/S3A 有局部实现,尚无通用发布权威 | STAC item、GeoParquet export | Iceberg catalog snapshot 是分析表版本权威;对象是物理内容,STAC 是发现投影 | Data Platform | AR-2 | | 在线空间数据 | PostGIS 业务表是当前编辑/查询事实,部分临时表混入 | Martin MVT、API JSON、导出文件 | 已批准 DataProductVersion 物化到 PostGIS;不能由瓦片或临时表反向定义产品版本 | GIS/Data Platform | AR-2 -> AR-4 | @@ -31,7 +31,7 @@ | Definition | `gda_control.platform_definition_version` 已绑定 definition ResourceVersion、完整逻辑 hash 和原子 gateway registration;3.4.2 adapter 可编译、创建并上线 provider DAG;binding 已以 append-only `execution_plan` Artifact 持久化并可按 tenant + artifact UUID 读取,旧 workflow/template/YAML 仍在写入 | 编辑器状态、DolphinScheduler DAG/definition | 旧 workflow 必须规范化并完整 hash 后才可形成 PlatformDefinitionVersion;provider binding 作为 ExecutionPlanArtifact/evidence,不可反写 definition | DataOps | AR-1 binding persistence 代码已验证 -> staging 调用链待验收 | | Run 最终状态 | `gda_control.platform_run/event` 已实现受控 submit/read/CAS;通用 transition 已禁止 `succeeded`,专用数据库 finalizer 只接受精确 workload、DolphinScheduler success observation、内容匹配 output、独立 passed QualityResult/evidence 和 input-to-output lineage;adapter standalone API path 已验证,但端到端 staging 尚未完成,legacy 路径继续运行 | Redis progress、日志、DolphinScheduler state、attempt observation | 旧 run 到 PlatformRun 永久 prohibited;已有 PlatformRun correlation 时才可转为 observation;provider 终态只进入 `reconciling`,ledger 经证据门唯一裁决成功 | DataOps/AgentOps | AR-1 success authority 本地/PostgreSQL 已验证 -> staging/生产切换待验收 | | 调度与补数 | APScheduler、自进化 scheduler 和调用方定时逻辑并存;DolphinScheduler POC 只验证 manual start/list/variables/STOP | UI schedule 列表 | DolphinScheduler 管 DataOps schedule/complement;Temporal 只管需要 durable signal/compensation 的 Agent/GWM workflow | DataOps/AgentOps | AR-1 manual correlation 已验证;schedule/complement/failover 待验收 | -| 事件交付 | Standards outbox 已数据库耐久;`platform_command_outbox` 支持 DolphinScheduler dispatch/reconcile;M3-14 `metadata_change_outbox` 将新 ResourceVersion 与内容绑定事件同事务写入;M3-15 managed consumer 在同一事务创建 `awaiting_authorization` activation request 并完成 event,base 为 0 replicas | command/metadata delivery status、消费者 claim、activation intent/request、worker status JSON、WebSocket 消息 | command/event 与源事实同事务入 outbox,幂等 consumer 交付;Active Metadata consumer 只能耐久化 inert request,必须在后续绑定真实 Definition/Run/execution plan/PolicyDecision/Approval 后才能创建 scheduler command,不能自行执行 provider mutation | Platform/Integrations/Metadata Platform | AR-1 command worker、M3-14 outbox 与 M3-15 consumer 本地已验证 -> protected authorization/scheduler 与 production scale-up 待执行 | +| 事件交付 | Standards outbox 已数据库耐久;`platform_command_outbox` 支持 DolphinScheduler dispatch/reconcile;M3-14 `metadata_change_outbox` 将新 ResourceVersion 与内容绑定事件同事务写入;M3-15 managed consumer 同事务创建 inert request;M3-16 将真实 ResourceVersion、Definition/Run/plan/PolicyDecision/Approval/authorizer 绑定后与一个 pending dispatch 同事务提交 | command/metadata delivery status、消费者 claim、activation intent/request/authorization、worker status JSON、WebSocket 消息 | command/event 与源事实同事务入 outbox,幂等 consumer 交付;Active Metadata consumer 不能授权或执行;activation capability 的普通 dispatch 被拒绝,只有 append-only authorization + deferred command FK + trigger 可创建命令 | Platform/Integrations/Metadata Platform | AR-1 command worker、M3-14/M3-15/M3-16 本地已验证 -> protected authorizer、真实 scheduler submission/read-back 与 production scale-up 待执行 | | 质量结果 | `gda_control.quality_result` 已提供 tenant RLS、append-only gateway 写入,绑定 Run、output ResourceVersion、rule version、verdict、metrics、evidence Artifact 和独立 evaluator;standards、QC、MMFE 专项结果仍未迁移 | dashboard、OpenMetadata quality summary | GDA ledger 保存产品终局所需的不可变 verdict/evidence;OpenMetadata 与 UI 只作可重建发现投影;旧结果缺稳定版本和证据时不得升级为终局依据 | Governance/DataOps | AR-1 最小成功证据已验证 -> 真实规则/staging 待接入 | | 标准与语义定义 | `std_*`、semantic registry 和 YAML 共同存在,生命周期未统一 | prompt/context、搜索索引 | 版本化 Standard/SemanticDefinition 经审批后为权威;Agent context 只消费批准版本 | Governance | AR-1 -> AR-3 | | 身份与权限 | Chainlit user 可显式绑定 tenant;versioned API 从认证 principal 派生 SubjectContext;`gda_control_gateway` 是 non-login/non-bypass 最小权限角色;Run 可引用强类型 PolicyDecision/Approval Artifact;M3-5/M3-6 分别验证本地 provider scoped grant、越权拒绝和 credential rotation/revocation;M3-7 已冻结生产 OIDC/workload/tenant binding、TLS、持久 catalog 与 attestation contract,但 40 个外部输入仍 blocked;M3-8 证明同一 Gravitino Basic role 在本地 JDBC restart 后连续 | session/cache、前端菜单权限、本地 provider identity/JDBC restart evidence、pending profile 与合成 readiness report | IdP/workload identity 提供真实 service identity;PolicyDecision/Approval 继续绑定不可变资源与 execution plan;只有 fresh protected attestation 可派生双 provider production identity claims,profile、Basic/JWT evidence、restart continuity 或人工批准均不可替代 | Security | AR-1 local identities/persistence + production readiness contract 已验证 -> protected 双 provider IAM/attestation 待执行 | @@ -59,6 +59,7 @@ 14. Metadata Fabric M1 只允许 OpenMetadata/Gravitino GET;M2 只执行本地 foundation/recovery/metrics/policy 演练或验证 production readiness profile;M3-1 只从 synthetic terminal evidence 生成 plan/candidate;M3-2 只允许 exact local PolicyDecision/Approval 后向本地 provider 写 projection;M3-3 只将同一 source evidence 经 PlatformGateway 写入临时 append-only binding ledger;M3-4 只经 tenant-scoped outbox 向无认证 loopback receiver 投递精确 candidate,并验证 at-least-once + receiver idempotency;M3-5 只证明临时 OpenMetadata bot 在 provider 强制 `DefaultBotRole` 之上的项目新增 grant 是 `table/Create`,并验证 policy-create 拒绝与本地 JWT 轮换/吊销;M3-6 只证明隔离 Gravitino Basic user 的 bounded table-create、catalog-create 拒绝、密码轮换/用户吊销和完整清理;M3-7 只冻结 production identity profile、精确 attestation binding 与 fail-closed 派生 claims,既不部署 identity path,也不提交真实 production attestation;M3-8 只证明 Docker Desktop 单集群中 Basic role、PostgreSQL JDBC metadata 与 file warehouse PVC 在受控 Pod restart 后连续;M3-9 只证明同节点共享 RWO PVC 的 Spark/Iceberg REST 互操作;M3-10 移除 Spark/Gravitino 共享 warehouse PVC,且只证明同一 Docker Desktop 主机/集群内跨节点 MinIO 的 read/write/schema evolution/snapshot/time travel 与对象级 metadata 一致;M3-11 只冻结 S3-compatible production profile、精确 attestation binding 与 fail-closed claims,既不选择/部署 provider,也不创建 bucket/KMS/policy 或提交真实 attestation;M3-12 只在 Spark driver loopback proxy 中于转发前注入 503,证明失败尝试零可见漂移、随后一次显式重试和直接对象清单无孤儿 data file;M3-13 只在 provider 200 响应丢失后以 HTTP 504 触发 commit-state-unknown,并对一个本地 append 做即时 readback/no-resubmit,不构成持久生产 reconcile controller、crash/concurrency proof 或网络 exactly-once。生产对象存储、identity/TLS、受保护环境故障注入、source-loss recovery、cancel/reconcile/lineage 与 Flink 仍未证明。Gravitino `1.3.0` Basic IdP 不算 OIDC,生产必须明确选择并证明 custom OIDC authenticator 或 identity-aware proxy。本地 bootstrap provisioner、Basic IdP、loopback/cluster HTTP、memory/file-backed JDBC catalog、同节点 PVC/MinIO、临时 identity/ledger/outbox、loopback receiver、pending profile、合成 attestation 和 local evidence 都不等于双 provider/生产最小权限、protected workload identity/OIDC、生产持久 catalog/binding、TLS、受保护 OpenLineage receiver、tenant isolation、alert/SLO、生产 ingestion/conformance 或生产写权威。 15. M3-14 只证明本地 PostgreSQL 16 中 ResourceVersion 与 `resource_version.registered` 事件同事务创建,以及 tenant/workload scoped claim、retry、lease reclaim 和 complete;它没有部署常驻 consumer、没有提交 DolphinScheduler、没有取得 provider apply 授权、没有执行 provider mutation,也不构成 production Active Metadata readiness。 16. M3-15 只证明 managed consumer 代码、默认 0 replicas/无 provider 或 scheduler credential 的部署边界,以及本地 PostgreSQL 中 durable inert request 与 event completion 的同事务原子性;`awaiting_authorization` request 不是 Definition、Run、execution plan、PolicyDecision、Approval 或 PlatformCommand,不能授权调度或 provider mutation。受保护 workload identity、实际 scale-up、scheduler submission、provider read-back、告警/SLO 与 production readiness 仍未验证。 +17. M3-16 只证明本地重庆 Shapefile bundle content fingerprint 与 ResourceVersion 精确绑定,以及 authorization + pending dispatch 的 PostgreSQL 原子性;真实数据不自动形成 authority 或授权,源文件/绝对路径不入 Git、CI 不依赖本机路径。受保护 authorizer identity、常驻 promotion controller、真实 DolphinScheduler submission/read-back、provider apply/mutation/ingestion、告警/SLO 与 production readiness 仍未验证。 ## 已建立的 AR-0/AR-1 entry 证据 @@ -95,8 +96,9 @@ - Metadata Fabric M3-11 已建立 production object-store profile/attestation gate;checked-in profile fingerprint 为 `668e194b3c688307014148391e7f389c9d6e9ca69c95d7b4cc92b4acae93181a`,report fingerprint 为 `85362dd10b7dc565f9fa567673d90b774cdec714bd1e70fb2c3c83c1af48b5ea`,`profile_valid=true`,43 项 provider/identity/transport/encryption/durability/consistency/tenancy/operations 外部输入以 blockers 暴露,`ready_for_protected_verification=false`、`attestation_valid=false`、`production_object_store_gate_passed=false`、`production_ready=false`。该合同绑定 M3-10 evidence,但没有选择或部署 provider;合成完整 attestation 只验证门禁逻辑,不计入生产证据。下一项真实证据是经 owner 批准并物化的 provider profile,以及来自 `production-object-store` 受保护环境、绑定当前 source/profile 并通过全部 26 项检查的 attestation。 - Metadata Fabric M3-12 已在本地 Spark driver loopback Iceberg REST proxy 中于 provider 转发前注入 HTTP 503。baseline 为 1 个 append snapshot、2 行和 1 个 referenced Parquet;失败写经 2 次 503 后 snapshot/row/file 零漂移;对同一逻辑行一次显式重试后为父子相连的 2 个 append snapshots、3 行和 2 个 referenced Parquet。直接 MinIO inventory 精确为 2 data + 3 metadata + 4 manifest = 9 objects,没有孤儿 data file;namespace、两块 PV 和 port-forward 均清理。contract fingerprint 为 `6d8944ab80246dc65891aa81118cb8b73f7ecad699be9a2af5e62d8260c41002`,evidence fingerprint 为 `39571cdac1e4043bcfc2d03a73b2b12ff925210daf8ae36bc640b8cb14d89401`。这只证明已知 pre-forward failure 的本地原子性与一次显式重试,不证明 provider uncertain outcome reconcile、网络 exactly-once、生产对象存储、cancel/lineage、Flink 或完整 Spark conformance。 - Metadata Fabric M3-13 已在本地 Spark driver loopback proxy 将 armed commit 转发给 Gravitino,并在 provider 200 后丢弃成功响应、返回 HTTP 504;Iceberg 将其映射为 commit-state-unknown,一次传输重试被抑制。Spark 只读 readback 后输出 `committed_do_not_resubmit` 和 `write_resubmitted=false`,最终为父子相连的 2 个 append snapshots、3 行、2 个 referenced Parquet;MinIO 为 2 data + 3 metadata + 4 manifest = 9 objects。Job `Complete 1/1`,namespace、两块 PV 和 port-forward 均清理。contract fingerprint 为 `7a8d75a1d6b4558b982c6c3242d8d356c5046955f8aae7a45e5c297b6f4d4132`,evidence fingerprint 为 `d6462fff78d07047311b1f715d5f2c7f08c0ce8fbdd5c8b26a3d95ddc3474786`。这只证明一个本地 append 的确定性 readback/no-resubmit;持久 controller、crash/concurrency、网络 exactly-once、生产对象存储、cancel/lineage、Flink 和完整 Spark conformance 仍未证明。 -- Metadata Fabric M3-14 已新增内容绑定 `MetadataChangeEvent`、migration 099 transactional outbox 与 PlatformGateway 原子注册/claim/fail/complete。真实 PostgreSQL 16 演练验证首次注册 `created=true`、pending/processed 精确 replay 均 `created=false`、错误 consumer/worker 拒绝、retry、lease-expiry reclaim、第三次 attempt 完成、processed 不再认领、旧 ResourceVersion 补事件整笔回滚、跨租户不可见、FORCE RLS 和 gateway 无直接 UPDATE/DELETE;最终只有 1 条权威事件。共享 contract/gateway 源码演进后已在 fresh database 重跑,当前 contract fingerprint 为 `b6162ed687c6554859b0056e34dd5cf9f051818deb533daf708eeea67b18eca5`,evidence fingerprint 为 `807444bb1adb2ea633bf98d845ba8a62fee608a4c6eff18c17c0c94f406e8a2`。激活意图只路由到 `metadata_fabric.projection_plan`,`provider_apply_authorized=false`、`provider_mutations_executed=false`、`production_ingestion_verified=false`、`production_scheduler_submission_verified=false`、`production_ready=false`。 -- Metadata Fabric M3-15 已新增 migration 100、内容绑定 `MetadataActivationRequest`、PlatformGateway atomic stage/read API、managed consumer/worker、默认 0 replicas Kustomize manifest 与静态部署 validator。真实 PostgreSQL 16 演练得到 2 个 processed events、2 个 `awaiting_authorization` requests、0 个 platform commands;no-request legacy complete 被阻断,request 精确 replay 为 `created=false`,consumer staging、跨租户拒绝、FORCE RLS、gateway 无直接 UPDATE/DELETE 均通过。contract fingerprint 为 `7a0bc0e3aaa53509443f031e7c40d7c5ff0f304510ddfeceedbad00147b505f5`,evidence fingerprint 为 `480a2d651c1389feed471cb6ff082537ddfcd2f108e2652b9168b16ea4d02977`。consumer 没有 provider/scheduler credential 或 Kubernetes token;`deployment_applied=false`、`production_workload_identity_verified=false`、`provider_apply_authorized=false`、`provider_mutations_executed=false`、`production_scheduler_submission_verified=false`、`production_ingestion_verified=false`、`production_ready=false`。 +- Metadata Fabric M3-14 已新增内容绑定 `MetadataChangeEvent`、migration 099 transactional outbox 与 PlatformGateway 原子注册/claim/fail/complete。真实 PostgreSQL 16 演练验证首次注册 `created=true`、pending/processed 精确 replay 均 `created=false`、错误 consumer/worker 拒绝、retry、lease-expiry reclaim、第三次 attempt 完成、processed 不再认领、旧 ResourceVersion 补事件整笔回滚、跨租户不可见、FORCE RLS 和 gateway 无直接 UPDATE/DELETE;最终只有 1 条权威事件。共享 contract/gateway 源码演进后已在 fresh database 重跑,当前 contract fingerprint 为 `c3d94228456aff7e9b134fa6bc746bbe6b7485950c16f8ec41ec34fc7a5ae567`,evidence fingerprint 为 `2b8a408e078cec44fde9a6d63e4b94f988dc820668c87ac0d2ef0a434d1a16a3`。激活意图只路由到 `metadata_fabric.projection_plan`,`provider_apply_authorized=false`、`provider_mutations_executed=false`、`production_ingestion_verified=false`、`production_scheduler_submission_verified=false`、`production_ready=false`。 +- Metadata Fabric M3-15 已新增 migration 100、内容绑定 `MetadataActivationRequest`、PlatformGateway atomic stage/read API、managed consumer/worker、默认 0 replicas Kustomize manifest 与静态部署 validator。真实 PostgreSQL 16 演练得到 2 个 processed events、2 个 `awaiting_authorization` requests、0 个 platform commands;no-request legacy complete 被阻断,request 精确 replay 为 `created=false`,consumer staging、跨租户拒绝、FORCE RLS、gateway 无直接 UPDATE/DELETE 均通过。contract fingerprint 为 `2256daa97c1f3a2e71f4d7026592171daea802b371f5616fbbd72a63939ee6b5`,evidence fingerprint 为 `029aaf7de476115dcf6385ca4a0e05bb84492ebea8ddecf6c3edf36edd76dbef`。consumer 没有 provider/scheduler credential 或 Kubernetes token;`deployment_applied=false`、`production_workload_identity_verified=false`、`provider_apply_authorized=false`、`provider_mutations_executed=false`、`production_scheduler_submission_verified=false`、`production_ingestion_verified=false`、`production_ready=false`。 +- Metadata Fabric M3-16 已新增 migration 101、内容绑定 `MetadataActivationAuthorization`、PlatformGateway atomic authorize/dispatch API 与 activation dispatch 数据库 guard。重庆中心城区历史文化街区 Shapefile 8 组件被转换为 path-free inventory,20 个 `PolygonZ`、33 字段、EPSG:4490,bundle SHA `fd474fd65c8e4a71da241eb3fd07748ca3b972fbd2d3c32833376dbe71104007` 精确成为 ResourceVersion content hash。真实 PostgreSQL 16 演练证明普通 dispatch 绕过被拒绝、无 command 的授权因 deferred FK 回滚,最终只有 1 个 authorization、1 个 pending command,精确 replay 不新增,FORCE RLS、function-only INSERT 与直接 UPDATE/DELETE 拒绝均通过。contract fingerprint 为 `cef78f91058a8529f4e86330790e714b52b73725a45ffe4dc9eded35bc8ccfa4`,evidence fingerprint 为 `6ae387240e3bcebaafe2ad7acc73f4e09d53df2e73b2ec63cd92edbc262d831e`。`deployment_applied=false`、`production_workload_identity_verified=false`、`provider_apply_authorized=false`、`provider_mutations_executed=false`、`production_scheduler_submission_verified=false`、`production_ingestion_verified=false`、`production_ready=false`。 ## 下一验收证据 @@ -104,7 +106,7 @@ - 真实 provenance artifact verify、受保护 overlay 的 `verified_for_staging_apply` release report,以及 staging/production 的 schema、config/runtime snapshot、registry/live DeploymentRevision 绑定、release/live artifact attestation 和环境 compare 报告; - staging 的 migration role、应用 login membership、连接池 role/tenant 复位、双租户 API 和 success finalization 运行产物; - DolphinScheduler adapter 的真实 IAM/OIDC、service token provisioning/轮换、provider 最小权限、binding artifact staging 接入、managed outbox worker/provider callback 实际扩容部署、唯一 worker ID、status/lease 故障恢复和无双写证据; -- Active Metadata consumer 的 production scale-up、受保护 workload identity,以及 request 绑定 Definition/Run/execution plan/PolicyDecision/Approval 后的 DolphinScheduler submission、幂等 projection execution、provider read-back、重试/死信告警与生产 SLO 证据; +- Active Metadata consumer/authorizer 的 production scale-up、受保护 workload identity、真实 DolphinScheduler submission/read-back、幂等 projection execution、provider read-back、重试/死信告警与生产 SLO 证据; - 首条真实图斑链对 golden slice 的 output hash、独立质量结果/evidence、血缘、发布 revision 和 rollback 演练; - OpenMetadata/Gravitino 的 source host/cluster 外生产 backup account/bucket、已批准 provider profile 与受保护对象存储 attestation、生产对象存储、KMS/TLS/workload identity、PITR/source-loss recovery、RPO/RTO、OIDC、受保护环境 provider NetworkPolicy/tenant isolation、upgrade/rollback、registry provenance、持续 metrics backend/retention/query、真实 alert delivery/SLO owner/runbook,以及受保护 PolicyDecision/Approval、双 provider 最小权限 ingestion、生产持久 binding、受保护 production OpenLineage receiver、无双写 read-back、受保护环境 commit failure injection、provider uncertain outcome reconcile、cancel/reconcile/lineage 和完整 Spark/Flink conformance;M1 fixture、M2 本地 evidence/readiness contracts、M3-1 projection candidate、M3-2 local replay、M3-3 临时 binding ledger、M3-4 loopback delivery、M3-5/M3-6 本地临时 provider identity、M3-7 pending profile/合成 attestation、M3-8 本地 JDBC restart continuity、M3-9 本地同节点 Spark interoperability、M3-10 本地同主机跨节点 MinIO interoperability、M3-11 pending object-store profile/合成 attestation、M3-12 local pre-forward commit-failure recovery 与 M3-13 local uncertain-commit readback/no-resubmit 均不计入生产退出门; - DolphinScheduler/Temporal sandbox 的独立数据库、备份恢复、身份、版本和升级责任证明;DolphinScheduler standalone/H2 不计入此退出门。 diff --git a/scripts/metadata-fabric-active-metadata-authorization.sh b/scripts/metadata-fabric-active-metadata-authorization.sh new file mode 100755 index 00000000..abc65c02 --- /dev/null +++ b/scripts/metadata-fabric-active-metadata-authorization.sh @@ -0,0 +1,4 @@ +#!/usr/bin/env bash +set -euo pipefail + +python -m data_agent.metadata_fabric_active_metadata_authorization "$@"