diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 1c04c3b7..894d2971 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -30,6 +30,7 @@ on: - feat/ar1-metadata-fabric-spark-commit-failure-recovery - feat/ar1-metadata-fabric-spark-uncertain-commit-reconciliation - feat/ar1-metadata-fabric-active-metadata-outbox + - feat/ar1-metadata-fabric-active-metadata-consumer env: PYTHON_VERSION: "3.13" @@ -169,6 +170,12 @@ jobs: - name: Validate metadata fabric Active Metadata outbox evidence run: python -m data_agent.metadata_fabric_active_metadata_outbox validate + - name: Validate metadata fabric Active Metadata consumer evidence + run: python -m data_agent.metadata_fabric_active_metadata_consumer validate + + - name: Validate Active Metadata consumer deployment boundary + run: python -m data_agent.active_metadata_consumer_deployment validate + - name: Validate DolphinScheduler adapter boundary run: python -m data_agent.dolphinscheduler_adapter validate @@ -200,6 +207,11 @@ jobs: DATABASE_URL: postgresql://postgres:postgres@localhost:5432/gis_agent_test run: python -m pytest data_agent/test_metadata_fabric_active_metadata_outbox_postgres.py -q + - name: Verify Active Metadata consumer on PostgreSQL + env: + DATABASE_URL: postgresql://postgres:postgres@localhost:5432/gis_agent_test + run: python -m pytest data_agent/test_active_metadata_consumer_postgres.py -q + - name: Run required platform tests env: DATABASE_URL: postgresql://postgres:postgres@localhost:5432/gis_agent_test @@ -229,7 +241,11 @@ 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_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_lineage_delivery.py \ data_agent/test_metadata_fabric_provider_identity.py \ data_agent/test_metadata_fabric_gravitino_identity.py \ diff --git a/data_agent/active_metadata_change_contract.py b/data_agent/active_metadata_change_contract.py index 2cee2389..d67ebacb 100644 --- a/data_agent/active_metadata_change_contract.py +++ b/data_agent/active_metadata_change_contract.py @@ -29,6 +29,7 @@ DELIVERY_SCHEMA = "gda.metadata_change_delivery.v1" REGISTRATION_SCHEMA = "gda.active_metadata_registration.v1" ACTIVATION_INTENT_SCHEMA = "gda.metadata_activation_intent.v1" +ACTIVATION_REQUEST_SCHEMA = "gda.metadata_activation_request.v1" RESOURCE_VERSION_REGISTERED = "resource_version.registered" METADATA_PROJECTION_ROUTE = "metadata_fabric.projection_plan" @@ -284,6 +285,64 @@ def build_metadata_activation_intent( ) +class MetadataActivationRequest(_FrozenModel): + request_schema: Literal["gda.metadata_activation_request.v1"] = Field( + default=ACTIVATION_REQUEST_SCHEMA, + alias="schema", + ) + request_id: UUID + intent: MetadataActivationIntent + status: Literal["awaiting_authorization"] = "awaiting_authorization" + 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 + request_sha256: Sha256 + + @model_validator(mode="after") + def _content_bound(self) -> Self: + expected_id = uuid5( + self.intent.event_id, + f"metadata-activation-request:{self.intent.intent_sha256}", + ) + if self.request_id != expected_id: + raise ValueError("MetadataActivationRequest ID does not match its intent") + stable = self.model_dump( + mode="json", + by_alias=True, + exclude={"request_sha256"}, + ) + if self.request_sha256 != canonical_json_fingerprint(stable): + raise ValueError("MetadataActivationRequest SHA-256 does not match") + return self + + +def build_metadata_activation_request( + intent: MetadataActivationIntent, +) -> MetadataActivationRequest: + request_id = uuid5( + intent.event_id, + f"metadata-activation-request:{intent.intent_sha256}", + ) + stable = { + "schema": ACTIVATION_REQUEST_SCHEMA, + "request_id": str(request_id), + "intent": intent.model_dump(mode="json", by_alias=True), + "status": "awaiting_authorization", + "provider_apply_authorized": False, + "provider_mutations_executed": False, + "production_scheduler_submission_verified": False, + "production_ingestion_verified": False, + "production_ready": False, + } + return MetadataActivationRequest( + request_id=request_id, + intent=intent, + request_sha256=canonical_json_fingerprint(stable), + ) + + class MetadataChangeDelivery(_FrozenModel): delivery_schema: Literal["gda.metadata_change_delivery.v1"] = Field( default=DELIVERY_SCHEMA, diff --git a/data_agent/active_metadata_consumer.py b/data_agent/active_metadata_consumer.py new file mode 100644 index 00000000..7344a29e --- /dev/null +++ b/data_agent/active_metadata_consumer.py @@ -0,0 +1,98 @@ +"""Thin consumer that durably stages Active Metadata activation requests.""" + +from __future__ import annotations + +from dataclasses import dataclass +from uuid import UUID + +from .active_metadata_change_contract import ( + build_metadata_activation_intent, + build_metadata_activation_request, +) +from .platform_gateway import ( + GatewayConflictError, + GatewayValidationError, + PlatformGateway, +) + + +@dataclass(frozen=True) +class ActiveMetadataBatchResult: + claimed: int + staged: int + replayed: int + retry_pending: int + failed: int + request_ids: tuple[UUID, ...] + + +class ActiveMetadataConsumer: + """Route changes to inert requests without owning authorization or execution.""" + + def __init__(self, gateway: PlatformGateway, *, consumer_subject: str): + if not consumer_subject.startswith("workload:"): + raise ValueError("Active Metadata consumer must use workload identity") + self.gateway = gateway + self.consumer_subject = consumer_subject + + def run_once( + self, + tenant_id: str, + *, + worker_id: str, + limit: int = 10, + lease_seconds: int = 60, + ) -> ActiveMetadataBatchResult: + deliveries = self.gateway.claim_metadata_changes( + tenant_id, + worker_id, + consumer_subject=self.consumer_subject, + limit=limit, + lease_seconds=lease_seconds, + ) + staged = 0 + replayed = 0 + retry_pending = 0 + failed = 0 + request_ids: list[UUID] = [] + for delivery in deliveries: + intent = build_metadata_activation_intent( + delivery.event, + routed_by=self.consumer_subject, + ) + request = build_metadata_activation_request(intent) + try: + result = self.gateway.stage_metadata_activation_request( + tenant_id, + delivery.event.event_id, + worker_id=worker_id, + request=request, + ) + except GatewayConflictError: + # The commit outcome or lease owner is uncertain. Leave the + # claim for deterministic reclaim instead of declaring failure. + retry_pending += 1 + continue + except GatewayValidationError: + self.gateway.fail_metadata_change( + tenant_id, + delivery.event.event_id, + worker_id=worker_id, + error_code="activation_contract_rejected", + retryable=False, + ) + failed += 1 + continue + request_ids.append(request.request_id) + if result.created: + staged += 1 + else: + replayed += 1 + return ActiveMetadataBatchResult( + claimed=len(deliveries), + staged=staged, + replayed=replayed, + retry_pending=retry_pending, + failed=failed, + request_ids=tuple(request_ids), + ) diff --git a/data_agent/active_metadata_consumer_deployment.py b/data_agent/active_metadata_consumer_deployment.py new file mode 100644 index 00000000..d3df9c70 --- /dev/null +++ b/data_agent/active_metadata_consumer_deployment.py @@ -0,0 +1,278 @@ +"""Validate the inert Kubernetes contract for the Active Metadata consumer.""" + +from __future__ import annotations + +import argparse +import hashlib +import json +import sys +from pathlib import Path +from typing import Any + +import yaml + +REPORT_SCHEMA = "gda.active_metadata_consumer_deployment.v1" +REPO_ROOT = Path(__file__).resolve().parent.parent +DEFAULT_MANIFEST = REPO_ROOT / "k8s/base/active-metadata-consumer.yaml" +DEFAULT_KUSTOMIZATION = REPO_ROOT / "k8s/base/kustomization.yaml" +DEFAULT_NETWORK_POLICY = REPO_ROOT / "k8s/base/networkpolicy.yaml" +DEPLOYMENT_NAME = "gis-agent-active-metadata-consumer" +STATUS_FILE = "/var/run/gis-agent/active-metadata/status.json" + +REQUIRED_LITERAL_ENV = { + "PYTHONDONTWRITEBYTECODE": "1", + "ACTIVE_METADATA_CONSUMER_ENABLED": "true", + "ACTIVE_METADATA_CONSUMER_STATUS_FILE": STATUS_FILE, + "ACTIVE_METADATA_CONSUMER_BATCH_SIZE": "1", + "ACTIVE_METADATA_CONSUMER_LEASE_SECONDS": "60", + "ACTIVE_METADATA_CONSUMER_POLL_INTERVAL_SECONDS": "5", + "ACTIVE_METADATA_CONSUMER_HEALTH_MAX_AGE_SECONDS": "30", +} +REQUIRED_CONFIG_ENV = { + "ACTIVE_METADATA_CONSUMER_TENANT_ID": "tenant-id", + "ACTIVE_METADATA_CONSUMER_SUBJECT": "consumer-subject", +} +FORBIDDEN_ENV_MARKERS = ( + "DOLPHINSCHEDULER", + "OPENMETADATA", + "GRAVITINO", + "PROVIDER_TOKEN", + "ACCESS_TOKEN", + "PASSWORD", +) + + +def _load_documents(path: Path) -> list[dict[str, Any]]: + return [ + value + for value in yaml.safe_load_all(path.read_text(encoding="utf-8")) + if isinstance(value, dict) + ] + + +def _resource( + documents: list[dict[str, Any]], kind: str, name: str +) -> dict[str, Any] | None: + for document in documents: + if document.get("kind") != kind: + continue + if (document.get("metadata") or {}).get("name") == name: + return document + return None + + +def _named(items: Any, name: str) -> dict[str, Any] | None: + if not isinstance(items, list): + return None + for item in items: + if isinstance(item, dict) and item.get("name") == name: + return item + return None + + +def _command_text(container: dict[str, Any]) -> str: + values = list(container.get("command") or []) + list(container.get("args") or []) + return "\n".join(str(value) for value in values) + + +def _probe_valid(container: dict[str, Any], name: str, command: str) -> bool: + values = (((container.get(name) or {}).get("exec") or {}).get("command") or []) + prefix = ( + "python", + "-m", + "data_agent.active_metadata_consumer_worker", + command, + "--status-file", + STATUS_FILE, + "--max-age-seconds", + ) + try: + max_age = float(values[len(prefix)]) + except (IndexError, TypeError, ValueError): + return False + return tuple(values[: len(prefix)]) == prefix and max_age > 0 + + +def build_deployment_report( + manifest_path: Path | None = None, + kustomization_path: Path | None = None, + network_policy_path: Path | None = None, + *, + expected_replicas: int = 0, +) -> dict[str, Any]: + manifest = (manifest_path or DEFAULT_MANIFEST).resolve() + kustomization = (kustomization_path or DEFAULT_KUSTOMIZATION).resolve() + network_policy = (network_policy_path or DEFAULT_NETWORK_POLICY).resolve() + errors: list[str] = [] + try: + documents = _load_documents(manifest) + except (OSError, yaml.YAMLError) as exc: + documents = [] + errors.append(f"consumer manifest is invalid: {type(exc).__name__}") + try: + kustom = yaml.safe_load(kustomization.read_text(encoding="utf-8")) or {} + except (OSError, yaml.YAMLError) as exc: + kustom = {} + errors.append(f"Kustomization is invalid: {type(exc).__name__}") + try: + policies = _load_documents(network_policy) + except (OSError, yaml.YAMLError) as exc: + policies = [] + errors.append(f"NetworkPolicy is invalid: {type(exc).__name__}") + + deployment = _resource(documents, "Deployment", DEPLOYMENT_NAME) + service_account = _resource(documents, "ServiceAccount", DEPLOYMENT_NAME) + if deployment is None: + errors.append("manifest must define the managed consumer Deployment") + if service_account is None: + errors.append("manifest must define a dedicated ServiceAccount") + if _resource(documents, "Role", DEPLOYMENT_NAME) is not None or _resource( + documents, "RoleBinding", DEPLOYMENT_NAME + ) is not None: + errors.append("consumer must not receive Kubernetes API RBAC") + if "active-metadata-consumer.yaml" not in (kustom.get("resources") or []): + errors.append("base Kustomization must register the consumer manifest") + + postgres_policy = _resource(policies, "NetworkPolicy", "postgres-access") + allowed_names = set() + if postgres_policy is not None: + ingress = (postgres_policy.get("spec") or {}).get("ingress") or [] + for rule in ingress: + for source in (rule or {}).get("from") or []: + labels = ((source or {}).get("podSelector") or {}).get( + "matchLabels" + ) or {} + if labels.get("app.kubernetes.io/name"): + allowed_names.add(labels["app.kubernetes.io/name"]) + if DEPLOYMENT_NAME not in allowed_names: + errors.append("PostgreSQL NetworkPolicy must admit the consumer selector") + + if deployment is not None: + spec = deployment.get("spec") or {} + pod_spec = ((spec.get("template") or {}).get("spec") or {}) + if spec.get("replicas") != expected_replicas: + errors.append("consumer replicas do not match the inert deployment gate") + if pod_spec.get("serviceAccountName") != DEPLOYMENT_NAME: + errors.append("consumer must use its dedicated ServiceAccount") + if pod_spec.get("automountServiceAccountToken") is not False: + errors.append("consumer must disable Kubernetes API token mounting") + if pod_spec.get("initContainers"): + errors.append("consumer must not prepare provider credentials") + containers = pod_spec.get("containers") or [] + container = _named(containers, "worker") + if container is None: + errors.append("consumer Deployment must contain the worker container") + else: + command_text = _command_text(container) + if ( + "data_agent.active_metadata_consumer_worker run" not in command_text + or "worker:active-metadata:${POD_UID}" not in command_text + ): + errors.append("consumer command must derive identity and run worker") + if container.get("envFrom"): + errors.append("consumer must not import bulk environment sources") + entries = container.get("env") or [] + names = [item.get("name") for item in entries if isinstance(item, dict)] + if len(names) != len(set(names)): + errors.append("consumer environment names must be unique") + env = { + item["name"]: item + for item in entries + if isinstance(item, dict) and isinstance(item.get("name"), str) + } + for name, value in REQUIRED_LITERAL_ENV.items(): + if (env.get(name) or {}).get("value") != value: + errors.append(f"{name} must use the fixed safe value") + for name, key in REQUIRED_CONFIG_ENV.items(): + ref = ((env.get(name) or {}).get("valueFrom") or {}).get( + "configMapKeyRef" + ) or {} + if ref != {"name": DEPLOYMENT_NAME, "key": key}: + errors.append(f"{name} must use the dedicated ConfigMap key") + database_ref = ( + ((env.get("DATABASE_URL") or {}).get("valueFrom") or {}).get( + "secretKeyRef" + ) + or {} + ) + if database_ref != {"name": DEPLOYMENT_NAME, "key": "database-url"}: + errors.append("DATABASE_URL must use the dedicated Secret key") + pod_uid_ref = ( + ((env.get("POD_UID") or {}).get("valueFrom") or {}).get("fieldRef") + or {} + ) + if pod_uid_ref.get("fieldPath") != "metadata.uid": + errors.append("worker identity must bind the immutable Pod UID") + forbidden = [ + name + for name in names + if isinstance(name, str) + and any(marker in name.upper() for marker in FORBIDDEN_ENV_MARKERS) + ] + if forbidden: + errors.append("consumer must not receive provider or scheduler secrets") + if not _probe_valid(container, "readinessProbe", "health"): + errors.append("readinessProbe must execute the worker health command") + for probe in ("startupProbe", "livenessProbe"): + if not _probe_valid(container, probe, "liveness"): + errors.append(f"{probe} must execute the worker liveness command") + security = container.get("securityContext") or {} + if ( + security.get("allowPrivilegeEscalation") is not False + or security.get("readOnlyRootFilesystem") is not True + or ((security.get("capabilities") or {}).get("drop") or []) + != ["ALL"] + ): + errors.append("consumer container security context is incomplete") + volumes = pod_spec.get("volumes") or [] + if len(volumes) != 1 or "emptyDir" not in (volumes[0] if volumes else {}): + errors.append("consumer may mount only its ephemeral status volume") + + files = {} + for name, path in ( + ("manifest", manifest), + ("kustomization", kustomization), + ("network_policy", network_policy), + ): + files[name] = { + "path": path.relative_to(REPO_ROOT).as_posix() + if path.is_relative_to(REPO_ROOT) + else path.as_posix(), + "sha256": hashlib.sha256(path.read_bytes()).hexdigest() + if path.is_file() + else None, + } + return { + "schema": REPORT_SCHEMA, + "status": "valid" if not errors else "invalid", + "errors": errors, + "expected_replicas": expected_replicas, + "files": files, + "provider_credentials_present": False, + "scheduler_credentials_present": False, + "deployment_applied": False, + "production_scheduler_submission_verified": False, + "production_ready": False, + } + + +def main(argv: list[str] | None = None) -> int: + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("command", choices=("validate",)) + parser.add_argument("--manifest", type=Path, default=DEFAULT_MANIFEST) + parser.add_argument("--kustomization", type=Path, default=DEFAULT_KUSTOMIZATION) + parser.add_argument("--network-policy", type=Path, default=DEFAULT_NETWORK_POLICY) + parser.add_argument("--expected-replicas", type=int, default=0) + args = parser.parse_args(argv) + report = build_deployment_report( + args.manifest, + args.kustomization, + args.network_policy, + expected_replicas=args.expected_replicas, + ) + print(json.dumps(report, ensure_ascii=False, indent=2, sort_keys=True)) + return 0 if report["status"] == "valid" else 1 + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/data_agent/active_metadata_consumer_worker.py b/data_agent/active_metadata_consumer_worker.py new file mode 100644 index 00000000..252d79db --- /dev/null +++ b/data_agent/active_metadata_consumer_worker.py @@ -0,0 +1,478 @@ +"""Managed process for the tenant-scoped Active Metadata consumer.""" + +from __future__ import annotations + +import argparse +import json +import math +import os +import signal +import sys +import threading +from collections.abc import Callable +from datetime import UTC, datetime +from pathlib import Path +from typing import Literal + +from dotenv import load_dotenv +from pydantic import BaseModel, ConfigDict, Field, ValidationError, field_validator + +from .active_metadata_consumer import ( + ActiveMetadataBatchResult, + ActiveMetadataConsumer, +) +from .db_engine import get_engine +from .observability import get_logger, setup_logging +from .platform_contracts import TenantId +from .platform_gateway import PlatformGateway, PlatformGatewayError + +WORKER_SCHEMA = "gda.active_metadata_consumer_worker.v1" +DEFAULT_STATUS_FILE = Path("/tmp/gda-active-metadata-consumer.json") + +load_dotenv(Path(__file__).resolve().parent / ".env", override=False) +setup_logging() +logger = get_logger("active_metadata_consumer_worker") + + +class ActiveMetadataWorkerConfigurationError(RuntimeError): + """The worker cannot start safely with its current configuration.""" + + +class _FrozenModel(BaseModel): + model_config = ConfigDict(extra="forbid", frozen=True) + + +class ActiveMetadataConsumerWorkerConfig(_FrozenModel): + enabled: Literal[True] + tenant_id: TenantId + worker_id: str = Field( + pattern=r"^worker:[A-Za-z0-9][A-Za-z0-9._:-]{0,246}$" + ) + consumer_subject: str = Field( + pattern=r"^workload:[A-Za-z0-9][A-Za-z0-9._:@/-]{0,247}$" + ) + batch_size: int = Field(default=10, ge=1, le=100) + lease_seconds: int = Field(default=60, ge=5, le=3600) + poll_interval_seconds: float = Field(default=5.0, ge=0.1, le=3600) + status_file: Path = DEFAULT_STATUS_FILE + health_max_age_seconds: float = Field(default=30.0, ge=1, le=7200) + + @field_validator("status_file") + @classmethod + def _absolute_status_path(cls, value: Path) -> Path: + if not value.is_absolute(): + raise ValueError("worker status file must be an absolute path") + return value + + @classmethod + def from_env(cls) -> ActiveMetadataConsumerWorkerConfig: + try: + config = cls( + enabled=( + str(os.getenv("ACTIVE_METADATA_CONSUMER_ENABLED") or "") + .strip() + .lower() + == "true" + ), + tenant_id=str( + os.getenv("ACTIVE_METADATA_CONSUMER_TENANT_ID") or "" + ).strip(), + worker_id=str( + os.getenv("ACTIVE_METADATA_CONSUMER_WORKER_ID") or "" + ).strip(), + consumer_subject=str( + os.getenv("ACTIVE_METADATA_CONSUMER_SUBJECT") or "" + ).strip(), + batch_size=int( + os.getenv("ACTIVE_METADATA_CONSUMER_BATCH_SIZE") or "10" + ), + lease_seconds=int( + os.getenv("ACTIVE_METADATA_CONSUMER_LEASE_SECONDS") or "60" + ), + poll_interval_seconds=float( + os.getenv("ACTIVE_METADATA_CONSUMER_POLL_INTERVAL_SECONDS") + or "5" + ), + status_file=Path( + os.getenv("ACTIVE_METADATA_CONSUMER_STATUS_FILE") + or DEFAULT_STATUS_FILE + ), + health_max_age_seconds=float( + os.getenv("ACTIVE_METADATA_CONSUMER_HEALTH_MAX_AGE_SECONDS") + or "30" + ), + ) + except (TypeError, ValueError, ValidationError) as exc: + raise ActiveMetadataWorkerConfigurationError( + "Active Metadata consumer worker configuration is invalid" + ) from exc + if config.health_max_age_seconds < config.poll_interval_seconds * 2: + raise ActiveMetadataWorkerConfigurationError( + "worker health max age must cover at least two polling intervals" + ) + return config + + def safe_summary(self) -> dict[str, object]: + return { + "enabled": self.enabled, + "tenant_id": self.tenant_id, + "worker_id": self.worker_id, + "consumer_subject": self.consumer_subject, + "batch_size": self.batch_size, + "lease_seconds": self.lease_seconds, + "poll_interval_seconds": self.poll_interval_seconds, + "status_file": self.status_file.as_posix(), + "provider_credentials_configured": False, + "scheduler_credentials_configured": False, + } + + +class ActiveMetadataConsumerWorkerStatus(_FrozenModel): + schema_name: Literal["gda.active_metadata_consumer_worker.v1"] = Field( + default=WORKER_SCHEMA, + alias="schema", + ) + state: Literal["starting", "ready", "degraded", "stopped"] + tenant_id: TenantId + worker_id: str + started_at: datetime + updated_at: datetime + last_success_at: datetime | None = None + cycles: int = Field(default=0, ge=0) + claimed: int = Field(default=0, ge=0) + staged: int = Field(default=0, ge=0) + replayed: int = Field(default=0, ge=0) + retry_pending: int = Field(default=0, ge=0) + failed: int = Field(default=0, ge=0) + consecutive_gateway_failures: int = Field(default=0, ge=0) + last_error_code: str | None = None + + @field_validator("started_at", "updated_at", "last_success_at") + @classmethod + def _aware_timestamp(cls, value: datetime | None) -> datetime | None: + if value is not None and (value.tzinfo is None or value.utcoffset() is None): + raise ValueError("worker status timestamps must include a timezone") + return value + + +class ActiveMetadataConsumerStatusStore: + def __init__(self, path: Path): + self.path = path + + def write(self, status: ActiveMetadataConsumerWorkerStatus) -> None: + self.path.parent.mkdir(parents=True, exist_ok=True) + temporary = self.path.parent / f".{self.path.name}.{os.getpid()}.tmp" + rendered = json.dumps( + status.model_dump(mode="json", by_alias=True), + ensure_ascii=True, + indent=2, + sort_keys=True, + ) + try: + temporary.write_text(rendered + "\n", encoding="utf-8") + temporary.chmod(0o600) + os.replace(temporary, self.path) + finally: + if temporary.exists(): + temporary.unlink() + + def read(self) -> ActiveMetadataConsumerWorkerStatus: + payload = json.loads(self.path.read_text(encoding="utf-8")) + return ActiveMetadataConsumerWorkerStatus.model_validate(payload) + + +class ActiveMetadataConsumerWorker: + """Own process lifecycle while PostgreSQL retains event/request state.""" + + def __init__( + self, + consumer: ActiveMetadataConsumer, + config: ActiveMetadataConsumerWorkerConfig, + *, + status_store: ActiveMetadataConsumerStatusStore | None = None, + stop_event: threading.Event | None = None, + clock: Callable[[], datetime] | None = None, + ): + self.consumer = consumer + self.config = config + self.status_store = status_store or ActiveMetadataConsumerStatusStore( + config.status_file + ) + self.stop_event = stop_event or threading.Event() + self.clock = clock or (lambda: datetime.now(UTC)) + self.status: ActiveMetadataConsumerWorkerStatus | None = None + + def _start(self) -> None: + now = self.clock() + self.status = ActiveMetadataConsumerWorkerStatus( + state="starting", + tenant_id=self.config.tenant_id, + worker_id=self.config.worker_id, + started_at=now, + updated_at=now, + ) + self.status_store.write(self.status) + + def run_cycle(self) -> ActiveMetadataBatchResult | None: + if self.status is None: + self._start() + assert self.status is not None + try: + result = self.consumer.run_once( + self.config.tenant_id, + worker_id=self.config.worker_id, + limit=self.config.batch_size, + lease_seconds=self.config.lease_seconds, + ) + except PlatformGatewayError as exc: + self.status = self.status.model_copy( + update={ + "state": "degraded", + "updated_at": self.clock(), + "consecutive_gateway_failures": ( + self.status.consecutive_gateway_failures + 1 + ), + "last_error_code": exc.code, + } + ) + self.status_store.write(self.status) + logger.warning("Active Metadata consumer cycle failed: %s", exc.code) + return None + + now = self.clock() + self.status = self.status.model_copy( + update={ + "state": "ready", + "updated_at": now, + "last_success_at": now, + "cycles": self.status.cycles + 1, + "claimed": self.status.claimed + result.claimed, + "staged": self.status.staged + result.staged, + "replayed": self.status.replayed + result.replayed, + "retry_pending": self.status.retry_pending + result.retry_pending, + "failed": self.status.failed + result.failed, + "consecutive_gateway_failures": 0, + "last_error_code": None, + } + ) + self.status_store.write(self.status) + logger.info( + "Active Metadata batch claimed=%d staged=%d replayed=%d retry=%d failed=%d", + result.claimed, + result.staged, + result.replayed, + result.retry_pending, + result.failed, + ) + return result + + def run(self, *, once: bool = False) -> int: + self._start() + logger.info( + "Active Metadata consumer worker starting tenant=%s worker=%s", + self.config.tenant_id, + self.config.worker_id, + ) + exit_code = 0 + try: + while not self.stop_event.is_set(): + result = self.run_cycle() + if once: + exit_code = 0 if result is not None else 1 + break + if result is not None and result.claimed >= self.config.batch_size: + continue + self.stop_event.wait(self.config.poll_interval_seconds) + finally: + assert self.status is not None + self.status = self.status.model_copy( + update={"state": "stopped", "updated_at": self.clock()} + ) + self.status_store.write(self.status) + logger.info("Active Metadata consumer worker stopped") + return exit_code + + +def _evaluate_status( + status_store: ActiveMetadataConsumerStatusStore, + *, + max_age_seconds: float, + liveness: bool, + now: datetime | None = None, +) -> tuple[dict[str, object], bool]: + current = now or datetime.now(UTC) + if ( + current.tzinfo is None + or current.utcoffset() is None + or not math.isfinite(max_age_seconds) + or max_age_seconds <= 0 + ): + return {"status": "unhealthy", "reason": "invalid_health_window"}, False + try: + status = status_store.read() + except (OSError, ValueError, json.JSONDecodeError, ValidationError): + return {"status": "unhealthy", "reason": "status_unavailable"}, False + if liveness: + if status.state == "stopped": + return {"status": "unhealthy", "reason": "worker_stopped"}, False + reference = status.updated_at + else: + if status.state != "ready" or status.last_success_at is None: + return { + "status": "unhealthy", + "reason": f"worker_{status.state}", + }, False + reference = status.last_success_at + age_seconds = (current - reference).total_seconds() + if age_seconds < -5: + return {"status": "unhealthy", "reason": "clock_skew"}, False + if age_seconds > max_age_seconds: + return { + "status": "unhealthy", + "reason": "status_stale", + "age_seconds": round(age_seconds, 3), + }, False + return { + "status": "healthy", + "age_seconds": round(max(0.0, age_seconds), 3), + "tenant_id": status.tenant_id, + "worker_id": status.worker_id, + "cycles": status.cycles, + "failed_requests": status.failed, + }, True + + +def evaluate_worker_health( + status_store: ActiveMetadataConsumerStatusStore, + *, + max_age_seconds: float, + now: datetime | None = None, +) -> tuple[dict[str, object], bool]: + return _evaluate_status( + status_store, + max_age_seconds=max_age_seconds, + liveness=False, + now=now, + ) + + +def evaluate_worker_liveness( + status_store: ActiveMetadataConsumerStatusStore, + *, + max_age_seconds: float, + now: datetime | None = None, +) -> tuple[dict[str, object], bool]: + return _evaluate_status( + status_store, + max_age_seconds=max_age_seconds, + liveness=True, + now=now, + ) + + +def _build_worker( + config: ActiveMetadataConsumerWorkerConfig, +) -> ActiveMetadataConsumerWorker: + engine = get_engine() + if engine is None: + raise ActiveMetadataWorkerConfigurationError( + "platform database is not configured" + ) + gateway = PlatformGateway(engine) + consumer = ActiveMetadataConsumer( + gateway, + consumer_subject=config.consumer_subject, + ) + return ActiveMetadataConsumerWorker(consumer, config) + + +def _install_signal_handlers(stop_event: threading.Event) -> None: + def request_stop(signum: int, _frame: object) -> None: + logger.info("received signal %d; stopping after the current batch", signum) + stop_event.set() + + signal.signal(signal.SIGINT, request_stop) + signal.signal(signal.SIGTERM, request_stop) + + +def _render(value: dict[str, object]) -> None: + print(json.dumps(value, ensure_ascii=False, indent=2, sort_keys=True)) + + +def main(argv: list[str] | None = None) -> int: + parser = argparse.ArgumentParser(description=__doc__) + subparsers = parser.add_subparsers(dest="command", required=True) + run_parser = subparsers.add_parser("run") + run_parser.add_argument("--once", action="store_true") + subparsers.add_parser("validate") + for probe_command in ("health", "liveness"): + probe_parser = subparsers.add_parser(probe_command) + probe_parser.add_argument("--status-file", type=Path) + probe_parser.add_argument("--max-age-seconds", type=float) + args = parser.parse_args(argv) + + if args.command in {"health", "liveness"}: + status_file = args.status_file or Path( + os.getenv("ACTIVE_METADATA_CONSUMER_STATUS_FILE") + or DEFAULT_STATUS_FILE + ) + try: + max_age = ( + args.max_age_seconds + if args.max_age_seconds is not None + else float( + os.getenv("ACTIVE_METADATA_CONSUMER_HEALTH_MAX_AGE_SECONDS") + or "30" + ) + ) + except ValueError: + _render({"status": "unhealthy", "reason": "invalid_max_age"}) + return 1 + evaluator = ( + evaluate_worker_health + if args.command == "health" + else evaluate_worker_liveness + ) + report, healthy = evaluator( + ActiveMetadataConsumerStatusStore(status_file), + max_age_seconds=max_age, + ) + _render(report) + return 0 if healthy else 1 + + try: + config = ActiveMetadataConsumerWorkerConfig.from_env() + if args.command == "validate": + if get_engine() is None: + raise ActiveMetadataWorkerConfigurationError( + "platform database is not configured" + ) + _render( + { + "schema": WORKER_SCHEMA, + "status": "valid", + "config": config.safe_summary(), + "provider_apply_authorized": False, + "production_scheduler_submission_verified": False, + "production_ready": False, + } + ) + return 0 + + worker = _build_worker(config) + _install_signal_handlers(worker.stop_event) + return worker.run(once=args.once) + except (ActiveMetadataWorkerConfigurationError, OSError) as exc: + _render( + { + "schema": WORKER_SCHEMA, + "status": "invalid", + "error": type(exc).__name__, + "message": str(exc), + } + ) + return 2 + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/data_agent/metadata_fabric_active_metadata_consumer.py b/data_agent/metadata_fabric_active_metadata_consumer.py new file mode 100644 index 00000000..a258bdd2 --- /dev/null +++ b/data_agent/metadata_fabric_active_metadata_consumer.py @@ -0,0 +1,572 @@ +"""Validate and rehearse inert Active Metadata activation request staging.""" + +from __future__ import annotations + +import argparse +import hashlib +import json +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_change_contract import ( + MetadataActivationRequest, + build_active_metadata_registration, + build_metadata_activation_intent, + build_metadata_activation_request, +) +from .active_metadata_consumer import ActiveMetadataConsumer +from .active_metadata_consumer_deployment import build_deployment_report +from .metadata_fabric_active_metadata_outbox import ( + CONSUMER_SUBJECT, + ISOLATED_TENANT, + TENANT, + WORKER_1, + WORKER_2, + build_active_metadata_bundle, +) +from .platform_contracts import canonical_json_fingerprint +from .platform_gateway import ( + GatewayNotFoundError, + GatewayValidationError, + PlatformGateway, +) + +CONTRACT_SCHEMA = "gda.active_metadata_consumer_contract.v1" +EVIDENCE_SCHEMA = "gda.active_metadata_consumer_evidence.v1" +REPO_ROOT = Path(__file__).resolve().parent.parent +DEFAULT_EVIDENCE_PATH = ( + REPO_ROOT + / "docs/evidence/metadata-fabric-active-metadata-consumer-2026-07-30.json" +) +DEFAULT_WRAPPER_PATH = ( + REPO_ROOT / "scripts/metadata-fabric-active-metadata-consumer.sh" +) +MIGRATIONS = tuple( + Path(__file__).resolve().parent / "migrations" / filename + for filename in ( + "092_platform_control_ledger.sql", + "093_app_user_tenant_context.sql", + "094_platform_control_gateway.sql", + "095_platform_command_outbox.sql", + "099_active_metadata_change_outbox.sql", + "100_active_metadata_activation_request.sql", + ) +) + + +class ActiveMetadataConsumerEvidenceError(RuntimeError): + """The consumer contract or local rehearsal failed closed.""" + + +def _load_json_object(path: Path) -> dict[str, Any]: + value = json.loads(path.read_text(encoding="utf-8")) + if not isinstance(value, dict): + raise ActiveMetadataConsumerEvidenceError( + f"{path.name} must contain an object" + ) + return value + + +def build_contract_report() -> dict[str, Any]: + errors: list[str] = [] + paths = { + "contract": Path(__file__).resolve().parent + / "active_metadata_change_contract.py", + "consumer": Path(__file__).resolve().parent + / "active_metadata_consumer.py", + "worker": Path(__file__).resolve().parent + / "active_metadata_consumer_worker.py", + "deployment_validator": Path(__file__).resolve().parent + / "active_metadata_consumer_deployment.py", + "gateway": Path(__file__).resolve().parent / "platform_gateway.py", + "migration": Path(__file__).resolve().parent + / "migrations/100_active_metadata_activation_request.sql", + "manifest": REPO_ROOT / "k8s/base/active-metadata-consumer.yaml", + "kustomization": REPO_ROOT / "k8s/base/kustomization.yaml", + "network_policy": REPO_ROOT / "k8s/base/networkpolicy.yaml", + "wrapper": DEFAULT_WRAPPER_PATH, + } + required = { + "contract": ( + "class MetadataActivationRequest", + "status: Literal[\"awaiting_authorization\"]", + "production_scheduler_submission_verified: Literal[False]", + "production_ready: Literal[False]", + ), + "consumer": ( + "class ActiveMetadataConsumer", + "self.gateway.stage_metadata_activation_request(", + "self.gateway.fail_metadata_change(", + ), + "worker": ( + "class ActiveMetadataConsumerWorker", + "ACTIVE_METADATA_CONSUMER_ENABLED", + "provider_credentials_configured\": False", + "scheduler_credentials_configured\": False", + ), + "deployment_validator": ( + "expected_replicas: int = 0", + "consumer must not receive provider or scheduler secrets", + "consumer must disable Kubernetes API token mounting", + ), + "gateway": ( + "def stage_metadata_activation_request(", + "def get_metadata_activation_request(", + ), + "migration": ( + "CREATE TABLE IF NOT EXISTS gda_control.metadata_activation_request", + "stage_metadata_activation_request", + "durable activation request is required before completion", + "FORCE ROW LEVEL SECURITY", + "status = 'awaiting_authorization'", + ), + "manifest": ( + "replicas: 0", + "automountServiceAccountToken: false", + "data_agent.active_metadata_consumer_worker run", + ), + "kustomization": ("active-metadata-consumer.yaml",), + "network_policy": ("gis-agent-active-metadata-consumer",), + "wrapper": ( + "data_agent.metadata_fabric_active_metadata_consumer", + '"$@"', + ), + } + 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(), + } + missing = [marker for marker in required[name] if marker not in source] + if missing: + errors.append(f"{name} is missing required consumer markers") + + deployment = build_deployment_report() + if deployment["status"] != "valid": + errors.append("Active Metadata consumer deployment contract is invalid") + stable = { + "schema": CONTRACT_SCHEMA, + "activation_route": "metadata_fabric.projection_plan", + "activation_boundary": "durable_request_awaiting_authorization", + "consumer_subject": CONSUMER_SUBJECT, + "deployment_expected_replicas": deployment["expected_replicas"], + "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 ActiveMetadataConsumerEvidenceError( + "local consumer 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 _request_is_inert(request: MetadataActivationRequest) -> bool: + return ( + request.status == "awaiting_authorization" + and request.provider_apply_authorized is False + and request.provider_mutations_executed is False + and request.production_scheduler_submission_verified is False + and request.production_ingestion_verified is False + and request.production_ready is False + ) + + +def run_local_rehearsal(database_url: str) -> dict[str, Any]: + if not database_url.startswith( + ("postgresql://", "postgresql+psycopg://", "postgresql+psycopg2://") + ): + raise ActiveMetadataConsumerEvidenceError( + "Active Metadata consumer rehearsal requires PostgreSQL" + ) + bundle = build_active_metadata_bundle() + engine = create_engine(database_url) + try: + _apply_migrations(engine) + gateway = PlatformGateway(engine) + gateway.register_resource(bundle.resource) + registration = gateway.register_resource_version_with_metadata_event( + bundle.registration, + max_attempts=3, + ) + claimed = gateway.claim_metadata_changes( + TENANT, + WORKER_1, + consumer_subject=CONSUMER_SUBJECT, + lease_seconds=60, + ) + intent = build_metadata_activation_intent( + claimed[0].event, + routed_by=CONSUMER_SUBJECT, + ) + request = build_metadata_activation_request(intent) + + legacy_completion_blocked = False + try: + gateway.complete_metadata_change( + TENANT, + claimed[0].event.event_id, + worker_id=WORKER_1, + activation_intent=intent, + ) + except GatewayValidationError: + legacy_completion_blocked = True + after_block = gateway.get_metadata_change_delivery( + TENANT, + claimed[0].event.event_id, + ) + request_absent_before_stage = False + try: + gateway.get_metadata_activation_request(TENANT, request.request_id) + except GatewayNotFoundError: + request_absent_before_stage = True + + first_stage = gateway.stage_metadata_activation_request( + TENANT, + claimed[0].event.event_id, + worker_id=WORKER_1, + request=request, + ) + replay_stage = gateway.stage_metadata_activation_request( + TENANT, + claimed[0].event.event_id, + worker_id=WORKER_1, + request=request, + ) + + next_version = bundle.registration.resource_version.model_copy( + update={ + "resource_version_id": UUID( + "a4000000-0000-4000-8000-000000000003" + ), + "version_key": "snapshot-2", + "predecessor_version_id": ( + bundle.registration.resource_version.resource_version_id + ), + "content_sha256": "c" * 64, + "authority_version_ref": {"snapshot_id": 2}, + } + ) + next_registration = build_active_metadata_registration( + next_version, + consumer_subject=CONSUMER_SUBJECT, + ) + gateway.register_resource_version_with_metadata_event( + next_registration, + max_attempts=3, + ) + consumer_result = ActiveMetadataConsumer( + gateway, + consumer_subject=CONSUMER_SUBJECT, + ).run_once( + TENANT, + worker_id=WORKER_2, + limit=1, + lease_seconds=60, + ) + stored_requests = [ + gateway.get_metadata_activation_request(TENANT, request.request_id), + gateway.get_metadata_activation_request( + TENANT, + consumer_result.request_ids[0], + ), + ] + + cross_tenant_read_blocked = False + try: + gateway.get_metadata_activation_request( + ISOLATED_TENANT, + request.request_id, + ) + except GatewayNotFoundError: + cross_tenant_read_blocked = True + + direct_mutation_blocked = True + with gateway._transaction(TENANT) as connection: + for statement in ( + """ + UPDATE gda_control.metadata_activation_request + SET status = 'awaiting_authorization' + WHERE request_id = :request_id + """, + """ + DELETE FROM gda_control.metadata_activation_request + WHERE request_id = :request_id + """, + ): + try: + with connection.begin_nested(): + connection.execute( + text(statement), + {"request_id": request.request_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_request', + 'SELECT,INSERT' + ), + NOT has_table_privilege( + 'gda_control_gateway', + 'gda_control.metadata_activation_request', 'UPDATE' + ), + NOT has_table_privilege( + 'gda_control_gateway', + 'gda_control.metadata_activation_request', 'DELETE' + ), + has_function_privilege( + 'gda_control_gateway', + 'gda_control.stage_metadata_activation_request(text,uuid,text,jsonb)', + 'EXECUTE' + ) + """ + ).one() + force_rls = connection.exec_driver_sql( + """ + SELECT bool_and(relforcerowsecurity) + FROM pg_class + WHERE oid IN ( + 'gda_control.metadata_change_outbox'::regclass, + 'gda_control.metadata_activation_request'::regclass + ) + """ + ).scalar_one() + counts = connection.execute( + text( + """ + SELECT + ( + SELECT count(*) + FROM gda_control.metadata_change_outbox + WHERE tenant_id = :tenant_id + AND status = 'processed' + ) AS processed_events, + ( + SELECT count(*) + FROM gda_control.metadata_activation_request + WHERE tenant_id = :tenant_id + AND status = 'awaiting_authorization' + ) AS activation_requests, + ( + SELECT count(*) + FROM gda_control.platform_command_outbox + WHERE tenant_id = :tenant_id + ) AS platform_commands + """ + ), + {"tenant_id": TENANT}, + ).one() + + deployment = build_deployment_report() + finally: + engine.dispose() + + requests_inert = all(_request_is_inert(item) for item in stored_requests) + atomic_completion_guard_verified = ( + legacy_completion_blocked + and request_absent_before_stage + and after_block.status.value == "in_flight" + ) + verified = ( + registration.created + and len(claimed) == 1 + and atomic_completion_guard_verified + and first_stage.created + and not replay_stage.created + and replay_stage.value == request + and consumer_result.claimed == consumer_result.staged == 1 + and consumer_result.replayed == 0 + and consumer_result.retry_pending == 0 + and consumer_result.failed == 0 + and len(consumer_result.request_ids) == 1 + and requests_inert + and cross_tenant_read_blocked + and direct_mutation_blocked + and privileges == (True, True, True, True) + and bool(force_rls) + and counts.processed_events == counts.activation_requests == 2 + and counts.platform_commands == 0 + and deployment["status"] == "valid" + and deployment["expected_replicas"] == 0 + ) + contract = build_contract_report() + stable = { + "schema": EVIDENCE_SCHEMA, + "status": ( + "local_postgresql_activation_request_staging_verified" + if verified + else "blocked" + ), + "contract_sha256": contract["contract_sha256"], + "activation_route": intent.route, + "event_ids": sorted(str(item.intent.event_id) for item in stored_requests), + "request_ids": sorted(str(item.request_id) for item in stored_requests), + "request_sha256": sorted(item.request_sha256 for item in stored_requests), + "request_statuses": sorted(item.status for item in stored_requests), + "processed_event_count": counts.processed_events, + "activation_request_count": counts.activation_requests, + "platform_command_count": counts.platform_commands, + "legacy_completion_without_request_blocked": legacy_completion_blocked, + "request_absent_before_atomic_stage": request_absent_before_stage, + "atomic_completion_guard_verified": atomic_completion_guard_verified, + "exact_request_replay_created": replay_stage.created, + "managed_consumer_staged": consumer_result.staged == 1, + "requests_inert": requests_inert, + "cross_tenant_read_blocked": cross_tenant_read_blocked, + "gateway_select_insert_only_verified": privileges + == (True, True, True, True), + "direct_request_mutation_blocked": direct_mutation_blocked, + "force_rls_verified": bool(force_rls), + "deployment_contract_verified": deployment["status"] == "valid", + "deployment_expected_replicas": deployment["expected_replicas"], + "local_postgresql_activation_request_staging_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 consumer rehearsal did not verify"], + } + return {**stable, "evidence_sha256": canonical_json_fingerprint(stable)} + + +def validate_rehearsal_evidence(evidence: dict[str, Any]) -> list[str]: + errors: list[str] = [] + stable = {key: value for key, value in evidence.items() if key != "evidence_sha256"} + if evidence.get("schema") != EVIDENCE_SCHEMA: + errors.append("Active Metadata consumer evidence schema does not match") + if evidence.get("evidence_sha256") != canonical_json_fingerprint(stable): + errors.append("Active Metadata consumer evidence SHA-256 does not match") + contract = build_contract_report() + if evidence.get("contract_sha256") != contract.get("contract_sha256"): + errors.append("Active Metadata consumer contract fingerprint is stale") + for claim in ( + "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 Active Metadata consumer evidence may not claim {claim}" + ) + for claim in ( + "legacy_completion_without_request_blocked", + "request_absent_before_atomic_stage", + "atomic_completion_guard_verified", + "managed_consumer_staged", + "requests_inert", + "cross_tenant_read_blocked", + "gateway_select_insert_only_verified", + "direct_request_mutation_blocked", + "force_rls_verified", + "deployment_contract_verified", + "local_postgresql_activation_request_staging_verified", + ): + if evidence.get(claim) is not True: + errors.append(f"Active Metadata consumer evidence did not verify {claim}") + if evidence.get("activation_route") != "metadata_fabric.projection_plan": + errors.append("Active Metadata consumer activation route is invalid") + if evidence.get("processed_event_count") != 2: + errors.append("Active Metadata consumer evidence must contain two events") + if evidence.get("activation_request_count") != 2: + errors.append("Active Metadata consumer evidence must contain two requests") + if evidence.get("platform_command_count") != 0: + errors.append("Active Metadata consumer must not create platform commands") + if evidence.get("exact_request_replay_created") is not False: + errors.append("Active Metadata request replay must not create a row") + if evidence.get("request_statuses") != [ + "awaiting_authorization", + "awaiting_authorization", + ]: + errors.append("Active Metadata requests must await authorization") + if evidence.get("deployment_expected_replicas") != 0: + errors.append("Active Metadata consumer base must remain inert") + return errors + + +def main(argv: list[str] | None = None) -> int: + parser = argparse.ArgumentParser(description=__doc__) + subparsers = parser.add_subparsers(dest="command", required=True) + validate = subparsers.add_parser("validate") + validate.add_argument("--evidence", type=Path, default=DEFAULT_EVIDENCE_PATH) + rehearse = subparsers.add_parser("rehearse") + rehearse.add_argument("--database-url", required=True) + rehearse.add_argument("--evidence-out", type=Path, required=True) + args = parser.parse_args(argv) + + if args.command == "validate": + report = build_contract_report() + try: + evidence = _load_json_object(args.evidence) + report["errors"].extend(validate_rehearsal_evidence(evidence)) + except (OSError, ValueError) as exc: + report["errors"].append( + f"Active Metadata consumer 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 + + evidence = run_local_rehearsal(args.database_url) + args.evidence_out.write_text( + json.dumps(evidence, ensure_ascii=True, indent=2, sort_keys=True) + "\n", + encoding="utf-8", + ) + print(json.dumps(evidence, ensure_ascii=True, indent=2, sort_keys=True)) + return 0 if not evidence["errors"] else 1 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/data_agent/migrations/100_active_metadata_activation_request.sql b/data_agent/migrations/100_active_metadata_activation_request.sql new file mode 100644 index 00000000..de7cb45f --- /dev/null +++ b/data_agent/migrations/100_active_metadata_activation_request.sql @@ -0,0 +1,322 @@ +-- 100: Durable Active Metadata activation requests. +-- +-- A MetadataChangeEvent may be marked processed only when its deterministic +-- activation request is persisted in the same transaction. The request is +-- deliberately inert: authorization, scheduler submission and provider +-- mutation remain separate, evidence-gated operations. + +CREATE TABLE IF NOT EXISTS gda_control.metadata_activation_request ( + tenant_id TEXT NOT NULL, + request_id UUID PRIMARY KEY, + event_id UUID NOT NULL, + event_sha256 CHAR(64) NOT NULL, + resource_urn TEXT NOT NULL, + resource_version_id UUID NOT NULL, + content_sha256 CHAR(64) NOT NULL, + activation_intent_sha256 CHAR(64) NOT NULL, + route TEXT NOT NULL, + requested_by TEXT NOT NULL, + request JSONB NOT NULL, + request_sha256 CHAR(64) NOT NULL, + status TEXT NOT NULL DEFAULT 'awaiting_authorization', + created_at TIMESTAMPTZ NOT NULL DEFAULT clock_timestamp(), + CONSTRAINT uq_gda_metadata_activation_tenant_request + UNIQUE (tenant_id, request_id), + CONSTRAINT uq_gda_metadata_activation_event + UNIQUE (tenant_id, event_id), + CONSTRAINT uq_gda_metadata_activation_request_sha + UNIQUE (tenant_id, request_sha256), + CONSTRAINT fk_gda_metadata_activation_event + FOREIGN KEY (tenant_id, event_id) + REFERENCES gda_control.metadata_change_outbox(tenant_id, event_id), + CONSTRAINT fk_gda_metadata_activation_version + FOREIGN KEY ( + tenant_id, resource_urn, resource_version_id, content_sha256 + ) REFERENCES gda_control.resource_version( + tenant_id, resource_urn, resource_version_id, content_sha256 + ), + CONSTRAINT ck_gda_metadata_activation_hashes CHECK ( + event_sha256 ~ '^[0-9a-f]{64}$' + AND content_sha256 ~ '^[0-9a-f]{64}$' + AND activation_intent_sha256 ~ '^[0-9a-f]{64}$' + AND request_sha256 ~ '^[0-9a-f]{64}$' + ), + CONSTRAINT ck_gda_metadata_activation_route CHECK ( + route = 'metadata_fabric.projection_plan' + ), + CONSTRAINT ck_gda_metadata_activation_requester CHECK ( + requested_by ~ '^workload:.+' + ), + CONSTRAINT ck_gda_metadata_activation_status CHECK ( + status = 'awaiting_authorization' + ), + CONSTRAINT ck_gda_metadata_activation_document CHECK ( + jsonb_typeof(request) = 'object' + AND request ?& ARRAY[ + 'schema', 'request_id', 'intent', 'status', + 'provider_apply_authorized', 'provider_mutations_executed', + 'production_scheduler_submission_verified', + 'production_ingestion_verified', 'production_ready', + 'request_sha256' + ] + AND request - ARRAY[ + 'schema', 'request_id', 'intent', 'status', + 'provider_apply_authorized', 'provider_mutations_executed', + 'production_scheduler_submission_verified', + 'production_ingestion_verified', 'production_ready', + 'request_sha256' + ] = '{}'::jsonb + AND request->>'schema' = 'gda.metadata_activation_request.v1' + AND request->>'request_id' = request_id::text + AND request->>'status' = status + AND request->>'request_sha256' = request_sha256 + AND (request->>'provider_apply_authorized')::boolean = false + AND (request->>'provider_mutations_executed')::boolean = false + AND ( + request->>'production_scheduler_submission_verified' + )::boolean = false + AND (request->>'production_ingestion_verified')::boolean = false + AND (request->>'production_ready')::boolean = false + AND jsonb_typeof(request->'intent') = 'object' + AND (request->'intent') ?& ARRAY[ + 'schema', 'event_id', 'event_sha256', 'tenant_id', + 'resource_urn', 'resource_version_id', 'content_sha256', + 'route', 'routed_by', 'provider_apply_authorized', + 'provider_mutations_executed', + 'production_ingestion_verified', 'intent_sha256' + ] + AND (request->'intent') - ARRAY[ + 'schema', 'event_id', 'event_sha256', 'tenant_id', + 'resource_urn', 'resource_version_id', 'content_sha256', + 'route', 'routed_by', 'provider_apply_authorized', + 'provider_mutations_executed', + 'production_ingestion_verified', 'intent_sha256' + ] = '{}'::jsonb + AND request->'intent'->>'schema' = 'gda.metadata_activation_intent.v1' + AND request->'intent'->>'event_id' = event_id::text + AND request->'intent'->>'event_sha256' = event_sha256 + AND request->'intent'->>'tenant_id' = tenant_id + AND request->'intent'->>'resource_urn' = resource_urn + AND request->'intent'->>'resource_version_id' = resource_version_id::text + AND request->'intent'->>'content_sha256' = content_sha256 + AND request->'intent'->>'intent_sha256' = activation_intent_sha256 + AND request->'intent'->>'route' = route + AND request->'intent'->>'routed_by' = requested_by + AND ( + request->'intent'->>'provider_apply_authorized' + )::boolean = false + AND ( + request->'intent'->>'provider_mutations_executed' + )::boolean = false + AND ( + request->'intent'->>'production_ingestion_verified' + )::boolean = false + ) +); + +CREATE INDEX IF NOT EXISTS idx_gda_metadata_activation_pending + ON gda_control.metadata_activation_request( + tenant_id, status, created_at, request_id + ); +CREATE INDEX IF NOT EXISTS idx_gda_metadata_activation_resource + ON gda_control.metadata_activation_request( + tenant_id, resource_urn, created_at DESC + ); + +ALTER TABLE gda_control.metadata_activation_request ENABLE ROW LEVEL SECURITY; +ALTER TABLE gda_control.metadata_activation_request FORCE ROW LEVEL SECURITY; +DROP POLICY IF EXISTS gda_metadata_activation_tenant_isolation + ON gda_control.metadata_activation_request; +CREATE POLICY gda_metadata_activation_tenant_isolation + ON gda_control.metadata_activation_request + USING (tenant_id = gda_control.current_tenant()) + WITH CHECK (tenant_id = gda_control.current_tenant()); + +CREATE OR REPLACE FUNCTION gda_control.stage_metadata_activation_request( + p_tenant_id TEXT, + p_event_id UUID, + p_worker_id TEXT, + p_request JSONB +) +RETURNS TABLE(activation_request JSONB, created BOOLEAN) +LANGUAGE plpgsql +SECURITY DEFINER +SET search_path = pg_catalog, gda_control +SET row_security = on +AS $$ +DECLARE + delivery gda_control.metadata_change_outbox%ROWTYPE; + stored gda_control.metadata_activation_request%ROWTYPE; + inserted_rows INTEGER := 0; +BEGIN + IF gda_control.current_tenant() IS DISTINCT FROM p_tenant_id THEN + RAISE EXCEPTION 'tenant context mismatch' USING ERRCODE = '42501'; + END IF; + IF COALESCE(btrim(p_worker_id), '') !~ '^worker:.+' THEN + RAISE EXCEPTION 'worker identity is invalid' USING ERRCODE = '22023'; + END IF; + IF jsonb_typeof(p_request) IS DISTINCT FROM 'object' THEN + RAISE EXCEPTION 'activation request must be an object' + USING ERRCODE = '22023'; + END IF; + + SELECT * + INTO delivery + FROM gda_control.metadata_change_outbox + WHERE tenant_id = p_tenant_id + AND event_id = p_event_id + FOR UPDATE; + IF NOT FOUND THEN + RAISE EXCEPTION 'metadata change was not found' USING ERRCODE = 'P0002'; + END IF; + + IF delivery.status = 'processed' THEN + SELECT * + INTO stored + FROM gda_control.metadata_activation_request + WHERE tenant_id = p_tenant_id + AND event_id = p_event_id; + IF NOT FOUND + OR stored.request IS DISTINCT FROM p_request + OR stored.activation_intent_sha256 IS DISTINCT FROM + delivery.activation_intent_sha256 THEN + RAISE EXCEPTION 'processed metadata change has no exact activation request' + USING ERRCODE = '23505'; + END IF; + RETURN QUERY SELECT stored.request, false; + RETURN; + END IF; + + IF delivery.status <> 'in_flight' + OR delivery.claimed_by IS DISTINCT FROM p_worker_id + OR delivery.claimed_until <= clock_timestamp() THEN + RAISE EXCEPTION 'metadata change claim is missing, expired, or owned by another worker' + USING ERRCODE = '40001'; + END IF; + + INSERT INTO gda_control.metadata_activation_request ( + tenant_id, request_id, event_id, event_sha256, resource_urn, + resource_version_id, content_sha256, activation_intent_sha256, + route, requested_by, request, request_sha256, status + ) VALUES ( + p_tenant_id, + (p_request->>'request_id')::uuid, + p_event_id, + delivery.event_sha256, + delivery.resource_urn, + delivery.resource_version_id, + delivery.content_sha256, + p_request->'intent'->>'intent_sha256', + p_request->'intent'->>'route', + p_request->'intent'->>'routed_by', + p_request, + p_request->>'request_sha256', + p_request->>'status' + ) + ON CONFLICT DO NOTHING; + GET DIAGNOSTICS inserted_rows = ROW_COUNT; + + SELECT * + INTO stored + FROM gda_control.metadata_activation_request + WHERE tenant_id = p_tenant_id + AND event_id = p_event_id; + IF NOT FOUND OR stored.request IS DISTINCT FROM p_request THEN + RAISE EXCEPTION 'activation request identity already has different content' + USING ERRCODE = '23505'; + END IF; + + UPDATE gda_control.metadata_change_outbox + SET status = 'processed', + claimed_by = NULL, + claimed_until = NULL, + last_error_code = NULL, + activation_intent_sha256 = stored.activation_intent_sha256, + completed_at = clock_timestamp() + WHERE tenant_id = p_tenant_id + AND event_id = p_event_id + AND status = 'in_flight' + AND claimed_by = p_worker_id + AND claimed_until > clock_timestamp(); + IF NOT FOUND THEN + RAISE EXCEPTION 'metadata change claim changed during activation staging' + USING ERRCODE = '40001'; + END IF; + + RETURN QUERY SELECT stored.request, inserted_rows = 1; +END; +$$; + +-- Migration 099 allowed completion with only an intent fingerprint. Once this +-- migration is present, completion also requires the durable request row. +CREATE OR REPLACE FUNCTION gda_control.complete_metadata_change( + p_tenant_id TEXT, + p_event_id UUID, + p_worker_id TEXT, + p_activation_intent_sha256 TEXT +) +RETURNS SETOF gda_control.metadata_change_outbox +LANGUAGE plpgsql +SECURITY DEFINER +SET search_path = pg_catalog, gda_control +SET row_security = on +AS $$ +BEGIN + IF gda_control.current_tenant() IS DISTINCT FROM p_tenant_id THEN + RAISE EXCEPTION 'tenant context mismatch' USING ERRCODE = '42501'; + END IF; + IF p_activation_intent_sha256 !~ '^[0-9a-f]{64}$' THEN + RAISE EXCEPTION 'activation intent fingerprint is invalid' + USING ERRCODE = '22023'; + END IF; + IF NOT EXISTS ( + SELECT 1 + FROM gda_control.metadata_activation_request + WHERE tenant_id = p_tenant_id + AND event_id = p_event_id + AND activation_intent_sha256 = p_activation_intent_sha256 + AND status = 'awaiting_authorization' + ) THEN + RAISE EXCEPTION 'durable activation request is required before completion' + USING ERRCODE = '22023'; + END IF; + RETURN QUERY + UPDATE gda_control.metadata_change_outbox AS delivery + SET status = 'processed', + claimed_by = NULL, + claimed_until = NULL, + last_error_code = NULL, + activation_intent_sha256 = p_activation_intent_sha256, + completed_at = clock_timestamp() + WHERE delivery.tenant_id = p_tenant_id + AND delivery.event_id = p_event_id + AND delivery.status = 'in_flight' + AND delivery.claimed_by = p_worker_id + AND delivery.claimed_until > clock_timestamp() + RETURNING delivery.*; + IF NOT FOUND THEN + RAISE EXCEPTION 'metadata change claim is missing, expired, or owned by another worker' + USING ERRCODE = '40001'; + END IF; +END; +$$; + +REVOKE ALL ON TABLE gda_control.metadata_activation_request FROM PUBLIC; +REVOKE ALL ON TABLE gda_control.metadata_activation_request + FROM gda_control_gateway; +GRANT SELECT, INSERT ON gda_control.metadata_activation_request + TO gda_control_gateway; + +REVOKE ALL ON FUNCTION gda_control.stage_metadata_activation_request( + text, uuid, text, jsonb +) FROM PUBLIC; +GRANT EXECUTE ON FUNCTION gda_control.stage_metadata_activation_request( + text, uuid, text, jsonb +) TO gda_control_gateway; + +REVOKE ALL ON FUNCTION gda_control.complete_metadata_change( + text, uuid, text, text +) FROM PUBLIC; +GRANT EXECUTE ON FUNCTION gda_control.complete_metadata_change( + text, uuid, text, text +) TO gda_control_gateway; diff --git a/data_agent/platform_gateway.py b/data_agent/platform_gateway.py index 814856b7..886dff67 100644 --- a/data_agent/platform_gateway.py +++ b/data_agent/platform_gateway.py @@ -20,6 +20,7 @@ from .active_metadata_change_contract import ( ActiveMetadataRegistration, MetadataActivationIntent, + MetadataActivationRequest, MetadataChangeDelivery, MetadataChangeDeliveryStatus, MetadataChangeEvent, @@ -104,6 +105,11 @@ / "migrations" / "099_active_metadata_change_outbox.sql" ) +ACTIVE_METADATA_ACTIVATION_MIGRATION = ( + Path(__file__).resolve().parent + / "migrations" + / "100_active_metadata_activation_request.sql" +) USER_TENANT_MIGRATION = ( Path(__file__).resolve().parent / "migrations" @@ -612,6 +618,82 @@ def complete_metadata_change( ).mappings().one() return self._metadata_change_from_row(row) + def get_metadata_activation_request( + self, + tenant_id: str, + 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: + raise GatewayNotFoundError( + "MetadataActivationRequest was not found" + ) + return MetadataActivationRequest.model_validate( + _as_json(row["request"]) + ) + + def stage_metadata_activation_request( + self, + tenant_id: str, + event_id: UUID, + *, + worker_id: str, + request: MetadataActivationRequest, + ) -> GatewayWriteResult: + try: + request = MetadataActivationRequest.model_validate( + request.model_dump(mode="json", by_alias=True) + ) + except ValueError as exc: + raise GatewayValidationError( + "metadata activation request is not content-bound" + ) from exc + if ( + request.intent.tenant_id != tenant_id + or request.intent.event_id != event_id + ): + raise GatewayValidationError( + "metadata activation request does not match the claimed event" + ) + with self._transaction(tenant_id) as connection: + row = connection.execute( + text( + """ + SELECT * + FROM gda_control.stage_metadata_activation_request( + :tenant_id, :event_id, :worker_id, + CAST(:request AS jsonb) + ) + """ + ), + { + "tenant_id": tenant_id, + "event_id": event_id, + "worker_id": worker_id, + "request": _json( + request.model_dump(mode="json", by_alias=True) + ), + }, + ).mappings().one() + stored = MetadataActivationRequest.model_validate( + _as_json(row["activation_request"]) + ) + if stored != request: + raise GatewayConflictError( + "stored metadata activation request differs from input" + ) + return GatewayWriteResult(stored, bool(row["created"])) + def fail_metadata_change( self, tenant_id: str, @@ -2187,6 +2269,7 @@ def build_gateway_report( binding_migration: Path | None = None, lineage_migration: Path | None = None, active_metadata_migration: Path | None = None, + activation_request_migration: Path | None = None, gateway_source: Path | None = None, routes_source: Path | None = None, command_consumer_source: Path | None = None, @@ -2211,6 +2294,10 @@ def build_gateway_report( "active_metadata_migration": ( active_metadata_migration or ACTIVE_METADATA_CHANGE_MIGRATION ).resolve(), + "activation_request_migration": ( + activation_request_migration + or ACTIVE_METADATA_ACTIVATION_MIGRATION + ).resolve(), "gateway_source": (gateway_source or Path(__file__)).resolve(), "routes_source": (routes_source or GATEWAY_ROUTES_SOURCE).resolve(), "command_consumer_source": ( @@ -2301,6 +2388,15 @@ def build_gateway_report( "ALTER TABLE gda_control.metadata_change_outbox FORCE ROW LEVEL SECURITY", "GRANT SELECT, INSERT ON gda_control.metadata_change_outbox", ), + "activation_request_migration": ( + "CREATE TABLE IF NOT EXISTS gda_control.metadata_activation_request", + "status = 'awaiting_authorization'", + "stage_metadata_activation_request", + "processed metadata change has no exact activation request", + "durable activation request is required before completion", + "ALTER TABLE gda_control.metadata_activation_request FORCE ROW LEVEL SECURITY", + "GRANT SELECT, INSERT ON gda_control.metadata_activation_request", + ), "gateway_source": ( 'SET LOCAL ROLE "{GATEWAY_DATABASE_ROLE}"', "SELECT set_config('app.current_tenant', :tenant, true)", @@ -2321,6 +2417,8 @@ def build_gateway_report( "def claim_metadata_changes(", "def complete_metadata_change(", "def fail_metadata_change(", + "def get_metadata_activation_request(", + "def stage_metadata_activation_request(", ), "routes_source": ( 'base = "/api/platform/v1"', diff --git a/data_agent/platform_truth.py b/data_agent/platform_truth.py index b55c3b70..dbd36176 100644 --- a/data_agent/platform_truth.py +++ b/data_agent/platform_truth.py @@ -283,6 +283,69 @@ def _config( maximum=7200, owner="sre", ), + _config( + "ACTIVE_METADATA_CONSUMER_ENABLED", + "bool", + False, + owner="metadata-platform", + description="Enable the managed Active Metadata request staging worker.", + ), + _config( + "ACTIVE_METADATA_CONSUMER_TENANT_ID", + "str", + None, + owner="metadata-platform", + ), + _config( + "ACTIVE_METADATA_CONSUMER_WORKER_ID", + "str", + None, + owner="sre", + ), + _config( + "ACTIVE_METADATA_CONSUMER_SUBJECT", + "str", + None, + owner="security", + ), + _config( + "ACTIVE_METADATA_CONSUMER_BATCH_SIZE", + "int", + 10, + minimum=1, + maximum=100, + owner="metadata-platform", + ), + _config( + "ACTIVE_METADATA_CONSUMER_LEASE_SECONDS", + "int", + 60, + minimum=5, + maximum=3600, + owner="metadata-platform", + ), + _config( + "ACTIVE_METADATA_CONSUMER_POLL_INTERVAL_SECONDS", + "float", + 5, + minimum=0.1, + maximum=3600, + owner="metadata-platform", + ), + _config( + "ACTIVE_METADATA_CONSUMER_STATUS_FILE", + "str", + "/tmp/gda-active-metadata-consumer.json", + owner="sre", + ), + _config( + "ACTIVE_METADATA_CONSUMER_HEALTH_MAX_AGE_SECONDS", + "float", + 30, + minimum=1, + maximum=7200, + owner="sre", + ), _config("ARCPY_MCP_ENABLED", "bool", False, owner="gis-runtime"), _config("ARCPY_MCP_URL", "url", None, owner="gis-runtime"), _config("ARCPY_MCP_TOKEN", "str", None, secret=True, owner="gis-runtime"), @@ -396,6 +459,33 @@ def _config( ), "Retain as tenant-scoped managed provider command delivery", ), + RuntimeSpec( + "active_metadata_consumer_worker", + "outbox_worker", + "governed", + "database_durable", + ( + "gda_control.metadata_change_outbox + " + "gda_control.metadata_activation_request" + ), + "metadata-platform", + "activation_request_staging_only", + ( + "data_agent/active_metadata_consumer_worker.py", + "data_agent/active_metadata_consumer.py", + ), + ( + ( + "data_agent/active_metadata_consumer_worker.py", + "class ActiveMetadataConsumerWorker", + ), + ( + "data_agent/active_metadata_consumer.py", + "self.gateway.stage_metadata_activation_request(", + ), + ), + "Retain as tenant-scoped inert activation request staging", + ), RuntimeSpec( "api_workflow_background", "api_background_task", @@ -766,7 +856,27 @@ def _config( "def run_local_rehearsal", ), ), - "Managed consumer submission to DolphinScheduler with protected identity", + "Durable inert activation request staging before authorization", + ), + RuntimeSpec( + "metadata_active_metadata_consumer_rehearsal", + "active_metadata_consumer_rehearsal", + "governed", + "evidence_durable", + "temporary PostgreSQL activation requests + committed local evidence", + "metadata-platform", + "local_verification_only", + ( + "data_agent/metadata_fabric_active_metadata_consumer.py", + "scripts/metadata-fabric-active-metadata-consumer.sh", + ), + ( + ( + "data_agent/metadata_fabric_active_metadata_consumer.py", + "def run_local_rehearsal", + ), + ), + "Protected authorization and scheduler promotion of durable requests", ), RuntimeSpec( "datalake_monitor", @@ -829,7 +939,7 @@ def _config( # Fingerprint of literal environment reads in production Python modules. It is # intentionally updated only with an explicit config-contract review. ENV_ACCESS_BASELINE_FINGERPRINT = ( - "41949811ca1d12a9d8bdbd5e7ecb1ba528be7049af96742306d6f317ab0791b8" + "5ee717911c109b480328a050893296e37591bfca748e3ed1743b7e3def3d9048" ) RUNTIME_PRIMITIVE_BASELINE_FINGERPRINT = ( "d6402d91e40ddb61591a7d258925d79e5eee964c3a9c0ace7de34acd10facbfd" diff --git a/data_agent/test_active_metadata_change_contract.py b/data_agent/test_active_metadata_change_contract.py index e6f2b335..905e5bb2 100644 --- a/data_agent/test_active_metadata_change_contract.py +++ b/data_agent/test_active_metadata_change_contract.py @@ -7,10 +7,12 @@ from data_agent.active_metadata_change_contract import ( ActiveMetadataContractError, MetadataActivationIntent, + MetadataActivationRequest, MetadataChangeDelivery, MetadataChangeEvent, build_active_metadata_registration, build_metadata_activation_intent, + build_metadata_activation_request, build_metadata_change_delivery, ) from data_agent.platform_contracts import ResourceVersion @@ -50,6 +52,7 @@ def test_registration_event_and_activation_intent_are_deterministic(): first.event, routed_by=CONSUMER, ) + request = build_metadata_activation_request(intent) assert first == second assert str(first.event.event_id) == "23bce695-edf5-53ef-b266-f053628f3446" @@ -63,6 +66,14 @@ def test_registration_event_and_activation_intent_are_deterministic(): assert intent.intent_sha256 == ( "169ac6b822d9af2ff75071eccc69a9468bffcacbe576bdad8430b95215d1a88c" ) + assert str(request.request_id) == "dc9257ee-7103-5eac-9c56-74830839a678" + assert request.status == "awaiting_authorization" + assert request.provider_apply_authorized is False + assert request.production_scheduler_submission_verified is False + assert request.production_ready is False + assert request.request_sha256 == ( + "ce2cc874a5306d8370bbbcdb1f5d17df47051ab523deffe4eb9a800578b05771" + ) def test_event_and_activation_intent_reject_content_tampering(): @@ -84,6 +95,12 @@ def test_event_and_activation_intent_reject_content_tampering(): with pytest.raises(ValidationError, match="SHA-256"): MetadataActivationIntent.model_validate(intent_payload) + request = build_metadata_activation_request(intent) + request_payload = request.model_dump(mode="json", by_alias=True) + request_payload["production_ready"] = True + with pytest.raises(ValidationError): + MetadataActivationRequest.model_validate(request_payload) + def test_authenticated_producer_and_exact_consumer_are_required(): with pytest.raises(ActiveMetadataContractError, match="authenticated subject"): diff --git a/data_agent/test_active_metadata_consumer.py b/data_agent/test_active_metadata_consumer.py new file mode 100644 index 00000000..238d44b4 --- /dev/null +++ b/data_agent/test_active_metadata_consumer.py @@ -0,0 +1,146 @@ +from datetime import timedelta + +import pytest + +from data_agent.active_metadata_change_contract import ( + MetadataChangeDelivery, + MetadataChangeDeliveryStatus, + build_metadata_activation_intent, + build_metadata_activation_request, +) +from data_agent.active_metadata_consumer import ActiveMetadataConsumer +from data_agent.metadata_fabric_active_metadata_outbox import ( + CONSUMER_SUBJECT, + TENANT, + WORKER_1, + build_active_metadata_bundle, +) +from data_agent.platform_gateway import ( + GatewayConflictError, + GatewayUnavailableError, + GatewayValidationError, + GatewayWriteResult, +) + + +def _claimed_delivery() -> MetadataChangeDelivery: + event = build_active_metadata_bundle().registration.event + return MetadataChangeDelivery( + event=event, + status=MetadataChangeDeliveryStatus.IN_FLIGHT, + attempt_count=1, + max_attempts=3, + available_at=event.occurred_at, + claimed_by=WORKER_1, + claimed_until=event.occurred_at + timedelta(minutes=1), + ) + + +class _Gateway: + def __init__(self, outcome=None): + self.delivery = _claimed_delivery() + self.outcome = outcome + self.claim_calls = [] + self.stage_calls = [] + self.fail_calls = [] + + def claim_metadata_changes(self, tenant_id, worker_id, **kwargs): + self.claim_calls.append((tenant_id, worker_id, kwargs)) + return [self.delivery] + + def stage_metadata_activation_request( + self, tenant_id, event_id, *, worker_id, request + ): + self.stage_calls.append((tenant_id, event_id, worker_id, request)) + if isinstance(self.outcome, Exception): + raise self.outcome + created = True if self.outcome is None else bool(self.outcome) + return GatewayWriteResult(request, created) + + def fail_metadata_change(self, *args, **kwargs): + self.fail_calls.append((args, kwargs)) + return self.delivery + + +def test_consumer_stages_deterministic_inert_request(): + gateway = _Gateway() + consumer = ActiveMetadataConsumer( + gateway, + consumer_subject=CONSUMER_SUBJECT, + ) + + result = consumer.run_once( + TENANT, + worker_id=WORKER_1, + limit=3, + lease_seconds=45, + ) + + expected = build_metadata_activation_request( + build_metadata_activation_intent( + gateway.delivery.event, + routed_by=CONSUMER_SUBJECT, + ) + ) + assert result.claimed == result.staged == 1 + assert result.replayed == result.retry_pending == result.failed == 0 + assert result.request_ids == (expected.request_id,) + assert gateway.stage_calls[0][3] == expected + assert expected.status == "awaiting_authorization" + assert expected.provider_apply_authorized is False + assert expected.production_scheduler_submission_verified is False + assert expected.production_ready is False + + +def test_consumer_reports_exact_stage_replay_without_duplicate_work(): + gateway = _Gateway(outcome=False) + result = ActiveMetadataConsumer( + gateway, + consumer_subject=CONSUMER_SUBJECT, + ).run_once(TENANT, worker_id=WORKER_1) + + assert result.claimed == result.replayed == 1 + assert result.staged == result.retry_pending == result.failed == 0 + + +def test_consumer_leaves_conflict_for_lease_reclaim(): + gateway = _Gateway(GatewayConflictError("uncertain stage")) + result = ActiveMetadataConsumer( + gateway, + consumer_subject=CONSUMER_SUBJECT, + ).run_once(TENANT, worker_id=WORKER_1) + + assert result.claimed == result.retry_pending == 1 + assert result.staged == result.replayed == result.failed == 0 + assert gateway.fail_calls == [] + + +def test_consumer_terminally_rejects_invalid_request_contract(): + gateway = _Gateway(GatewayValidationError("invalid request")) + result = ActiveMetadataConsumer( + gateway, + consumer_subject=CONSUMER_SUBJECT, + ).run_once(TENANT, worker_id=WORKER_1) + + assert result.claimed == result.failed == 1 + assert gateway.fail_calls[0][1] == { + "worker_id": WORKER_1, + "error_code": "activation_contract_rejected", + "retryable": False, + } + + +def test_consumer_propagates_database_unavailability_to_managed_worker(): + gateway = _Gateway(GatewayUnavailableError("database unavailable")) + consumer = ActiveMetadataConsumer( + gateway, + consumer_subject=CONSUMER_SUBJECT, + ) + + with pytest.raises(GatewayUnavailableError): + consumer.run_once(TENANT, worker_id=WORKER_1) + + +def test_consumer_requires_workload_subject(): + with pytest.raises(ValueError, match="workload identity"): + ActiveMetadataConsumer(_Gateway(), consumer_subject="human:operator") diff --git a/data_agent/test_active_metadata_consumer_deployment.py b/data_agent/test_active_metadata_consumer_deployment.py new file mode 100644 index 00000000..d99a6938 --- /dev/null +++ b/data_agent/test_active_metadata_consumer_deployment.py @@ -0,0 +1,98 @@ +from copy import deepcopy + +import yaml + +from data_agent import active_metadata_consumer_deployment as deployment + + +def _documents(): + return list( + yaml.safe_load_all( + deployment.DEFAULT_MANIFEST.read_text(encoding="utf-8") + ) + ) + + +def _write_documents(path, documents): + path.write_text( + yaml.safe_dump_all(documents, sort_keys=False), + encoding="utf-8", + ) + + +def test_inert_consumer_deployment_is_database_only_and_fail_closed(): + report = deployment.build_deployment_report() + + assert report["status"] == "valid" + assert report["errors"] == [] + assert report["expected_replicas"] == 0 + assert report["provider_credentials_present"] is False + assert report["scheduler_credentials_present"] is False + assert report["deployment_applied"] is False + assert report["production_scheduler_submission_verified"] is False + assert report["production_ready"] is False + + +def test_validator_rejects_enabled_base_and_scheduler_secret(tmp_path): + documents = deepcopy(_documents()) + workload = documents[0] + workload["spec"]["replicas"] = 1 + container = workload["spec"]["template"]["spec"]["containers"][0] + container["env"].append( + {"name": "DOLPHINSCHEDULER_TOKEN", "value": "must-not-be-here"} + ) + unsafe = tmp_path / "unsafe-consumer.yaml" + _write_documents(unsafe, documents) + + report = deployment.build_deployment_report(unsafe) + + assert report["status"] == "invalid" + assert "consumer replicas do not match the inert deployment gate" in report[ + "errors" + ] + assert "consumer must not receive provider or scheduler secrets" in report[ + "errors" + ] + + +def test_validator_rejects_kubernetes_token_and_missing_network_access(tmp_path): + documents = deepcopy(_documents()) + documents[0]["spec"]["template"]["spec"][ + "automountServiceAccountToken" + ] = True + unsafe = tmp_path / "token-enabled.yaml" + _write_documents(unsafe, documents) + + policies = list( + yaml.safe_load_all( + deployment.DEFAULT_NETWORK_POLICY.read_text(encoding="utf-8") + ) + ) + postgres = next( + item + for item in policies + if item["kind"] == "NetworkPolicy" + and item["metadata"]["name"] == "postgres-access" + ) + sources = postgres["spec"]["ingress"][0]["from"] + postgres["spec"]["ingress"][0]["from"] = [ + source + for source in sources + if source.get("podSelector", {}).get("matchLabels", {}).get( + "app.kubernetes.io/name" + ) + != deployment.DEPLOYMENT_NAME + ] + unsafe_network = tmp_path / "networkpolicy.yaml" + _write_documents(unsafe_network, policies) + + report = deployment.build_deployment_report( + unsafe, + network_policy_path=unsafe_network, + ) + + assert report["status"] == "invalid" + assert "consumer must disable Kubernetes API token mounting" in report["errors"] + assert "PostgreSQL NetworkPolicy must admit the consumer selector" in report[ + "errors" + ] diff --git a/data_agent/test_active_metadata_consumer_postgres.py b/data_agent/test_active_metadata_consumer_postgres.py new file mode 100644 index 00000000..1376f9a8 --- /dev/null +++ b/data_agent/test_active_metadata_consumer_postgres.py @@ -0,0 +1,254 @@ +import os +from pathlib import Path +from uuid import UUID, uuid4 + +import pytest +from sqlalchemy import create_engine, text +from sqlalchemy.engine import make_url +from sqlalchemy.exc import DBAPIError + +from data_agent.active_metadata_change_contract import ( + build_active_metadata_registration, + build_metadata_activation_intent, + build_metadata_activation_request, +) +from data_agent.active_metadata_consumer import ActiveMetadataConsumer +from data_agent.metadata_fabric_active_metadata_outbox import ( + CONSUMER_SUBJECT, + TENANT, + WORKER_1, + WORKER_2, + build_active_metadata_bundle, +) +from data_agent.platform_gateway import ( + GatewayNotFoundError, + GatewayValidationError, + PlatformGateway, +) + +DATABASE_URL = os.environ.get("DATABASE_URL") +MIGRATIONS = tuple( + Path(__file__).resolve().parent / "migrations" / filename + for filename in ( + "092_platform_control_ledger.sql", + "093_app_user_tenant_context.sql", + "094_platform_control_gateway.sql", + "099_active_metadata_change_outbox.sql", + "100_active_metadata_activation_request.sql", + ) +) + + +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 consumer test requires a superuser") + database_name = f"gda_active_consumer_{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 _apply_migrations(engine) -> None: + with engine.begin() as connection: + 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"))) + + +@pytest.mark.skipif(not DATABASE_URL, reason="DATABASE_URL is not configured") +def test_postgres_consumer_stages_request_atomically_and_fail_closed(): + admin_engine, database_name, database_url = _temporary_database_url() + engine = create_engine(database_url) + try: + _apply_migrations(engine) + gateway = PlatformGateway(engine) + bundle = build_active_metadata_bundle() + gateway.register_resource(bundle.resource) + gateway.register_resource_version_with_metadata_event( + bundle.registration, + max_attempts=3, + ) + claimed = gateway.claim_metadata_changes( + TENANT, + WORKER_1, + consumer_subject=CONSUMER_SUBJECT, + lease_seconds=60, + ) + assert len(claimed) == 1 + intent = build_metadata_activation_intent( + claimed[0].event, + routed_by=CONSUMER_SUBJECT, + ) + request = build_metadata_activation_request(intent) + + with pytest.raises( + GatewayValidationError, + match="platform contract was rejected", + ): + gateway.complete_metadata_change( + TENANT, + claimed[0].event.event_id, + worker_id=WORKER_1, + activation_intent=intent, + ) + still_claimed = gateway.get_metadata_change_delivery( + TENANT, + claimed[0].event.event_id, + ) + assert still_claimed.status.value == "in_flight" + + first = gateway.stage_metadata_activation_request( + TENANT, + claimed[0].event.event_id, + worker_id=WORKER_1, + request=request, + ) + replay = gateway.stage_metadata_activation_request( + TENANT, + claimed[0].event.event_id, + worker_id=WORKER_1, + request=request, + ) + assert first.created is True + assert replay.created is False + assert replay.value == request + assert gateway.get_metadata_activation_request( + TENANT, + request.request_id, + ) == request + completed = gateway.get_metadata_change_delivery( + TENANT, + claimed[0].event.event_id, + ) + assert completed.status.value == "processed" + assert completed.activation_intent_sha256 == intent.intent_sha256 + + next_version = bundle.registration.resource_version.model_copy( + update={ + "resource_version_id": UUID( + "a4000000-0000-4000-8000-000000000003" + ), + "version_key": "snapshot-2", + "predecessor_version_id": ( + bundle.registration.resource_version.resource_version_id + ), + "content_sha256": "c" * 64, + "authority_version_ref": {"snapshot_id": 2}, + } + ) + next_registration = build_active_metadata_registration( + next_version, + consumer_subject=CONSUMER_SUBJECT, + ) + gateway.register_resource_version_with_metadata_event( + next_registration, + max_attempts=3, + ) + result = ActiveMetadataConsumer( + gateway, + consumer_subject=CONSUMER_SUBJECT, + ).run_once( + TENANT, + worker_id=WORKER_2, + limit=1, + lease_seconds=60, + ) + assert result.claimed == result.staged == 1 + assert result.replayed == result.retry_pending == result.failed == 0 + + with pytest.raises(GatewayNotFoundError): + gateway.get_metadata_activation_request( + "active-metadata-isolated", + request.request_id, + ) + + with engine.connect() as connection: + privileges = connection.exec_driver_sql( + """ + SELECT + has_table_privilege( + 'gda_control_gateway', + 'gda_control.metadata_activation_request', + 'SELECT,INSERT' + ), + has_table_privilege( + 'gda_control_gateway', + 'gda_control.metadata_activation_request', 'UPDATE' + ), + has_table_privilege( + 'gda_control_gateway', + 'gda_control.metadata_activation_request', 'DELETE' + ), + has_function_privilege( + 'gda_control_gateway', + 'gda_control.stage_metadata_activation_request(text,uuid,text,jsonb)', + 'EXECUTE' + ) + """ + ).one() + assert privileges == (True, False, False, True) + + with gateway._transaction(TENANT) as connection: + for statement in ( + """ + UPDATE gda_control.metadata_activation_request + SET status = 'awaiting_authorization' + WHERE request_id = :request_id + """, + """ + DELETE FROM gda_control.metadata_activation_request + WHERE request_id = :request_id + """, + ): + with pytest.raises(DBAPIError): + with connection.begin_nested(): + connection.execute(text(statement), {"request_id": request.request_id}) + + with engine.connect() as connection: + request_count = connection.exec_driver_sql( + "SELECT COUNT(*) FROM gda_control.metadata_activation_request" + ).scalar_one() + processed_count = connection.exec_driver_sql( + """ + SELECT COUNT(*) + FROM gda_control.metadata_change_outbox + WHERE status = 'processed' + """ + ).scalar_one() + assert request_count == processed_count == 2 + finally: + engine.dispose() + _drop_temporary_database(admin_engine, database_name) diff --git a/data_agent/test_active_metadata_consumer_worker.py b/data_agent/test_active_metadata_consumer_worker.py new file mode 100644 index 00000000..60284d38 --- /dev/null +++ b/data_agent/test_active_metadata_consumer_worker.py @@ -0,0 +1,200 @@ +import stat +from datetime import UTC, datetime, timedelta + +import pytest +from pydantic import ValidationError + +from data_agent.active_metadata_consumer import ActiveMetadataBatchResult +from data_agent.active_metadata_consumer_worker import ( + ActiveMetadataConsumerStatusStore, + ActiveMetadataConsumerWorker, + ActiveMetadataConsumerWorkerConfig, + ActiveMetadataWorkerConfigurationError, + evaluate_worker_health, + evaluate_worker_liveness, +) +from data_agent.platform_gateway import GatewayUnavailableError + +NOW = datetime(2026, 7, 30, 12, 0, tzinfo=UTC) + + +def _config(tmp_path, **overrides): + values = { + "enabled": True, + "tenant_id": "tenant-a", + "worker_id": "worker:active-metadata:pod-a", + "consumer_subject": "workload:metadata-router", + "batch_size": 10, + "lease_seconds": 60, + "poll_interval_seconds": 5, + "status_file": tmp_path / "active-metadata-status.json", + "health_max_age_seconds": 30, + } + values.update(overrides) + return ActiveMetadataConsumerWorkerConfig(**values) + + +def _batch(**overrides): + values = { + "claimed": 1, + "staged": 1, + "replayed": 0, + "retry_pending": 0, + "failed": 0, + "request_ids": (), + } + values.update(overrides) + return ActiveMetadataBatchResult(**values) + + +class _Consumer: + def __init__(self, outcomes): + self.outcomes = list(outcomes) + self.calls = [] + + def run_once(self, tenant_id, *, worker_id, limit, lease_seconds): + self.calls.append((tenant_id, worker_id, limit, lease_seconds)) + outcome = self.outcomes.pop(0) + if isinstance(outcome, Exception): + raise outcome + return outcome + + +def test_worker_config_is_strict_and_has_no_provider_credentials(tmp_path): + with pytest.raises(ValidationError, match="worker_id"): + _config(tmp_path, worker_id="shared-process") + with pytest.raises(ValidationError, match="consumer_subject"): + _config(tmp_path, consumer_subject="human:operator") + with pytest.raises(ValidationError, match="absolute"): + _config(tmp_path, status_file="relative/status.json") + + summary = _config(tmp_path).safe_summary() + assert summary["provider_credentials_configured"] is False + assert summary["scheduler_credentials_configured"] is False + + +def test_config_from_env_requires_identity_and_safe_health_window( + tmp_path, monkeypatch +): + monkeypatch.setenv("ACTIVE_METADATA_CONSUMER_ENABLED", "true") + monkeypatch.setenv("ACTIVE_METADATA_CONSUMER_TENANT_ID", "tenant-a") + monkeypatch.setenv( + "ACTIVE_METADATA_CONSUMER_WORKER_ID", "worker:active-metadata:pod-a" + ) + monkeypatch.setenv( + "ACTIVE_METADATA_CONSUMER_SUBJECT", "workload:metadata-router" + ) + monkeypatch.setenv( + "ACTIVE_METADATA_CONSUMER_STATUS_FILE", str(tmp_path / "status.json") + ) + config = ActiveMetadataConsumerWorkerConfig.from_env() + assert config.tenant_id == "tenant-a" + + monkeypatch.setenv("ACTIVE_METADATA_CONSUMER_ENABLED", "false") + with pytest.raises(ActiveMetadataWorkerConfigurationError): + ActiveMetadataConsumerWorkerConfig.from_env() + + monkeypatch.setenv("ACTIVE_METADATA_CONSUMER_ENABLED", "true") + monkeypatch.setenv("ACTIVE_METADATA_CONSUMER_HEALTH_MAX_AGE_SECONDS", "5") + with pytest.raises(ActiveMetadataWorkerConfigurationError, match="two polling"): + ActiveMetadataConsumerWorkerConfig.from_env() + + +def test_worker_writes_sanitized_atomic_status_and_accumulates_counts(tmp_path): + config = _config(tmp_path) + store = ActiveMetadataConsumerStatusStore(config.status_file) + consumer = _Consumer([_batch(replayed=1, retry_pending=1)]) + worker = ActiveMetadataConsumerWorker( + consumer, + config, + status_store=store, + clock=lambda: NOW, + ) + + result = worker.run_cycle() + + assert result is not None + status = store.read() + assert status.state == "ready" + assert status.claimed == status.staged == status.replayed == 1 + assert status.retry_pending == 1 + assert stat.S_IMODE(config.status_file.stat().st_mode) == 0o600 + rendered = config.status_file.read_text(encoding="utf-8") + assert "DATABASE_URL" not in rendered + assert "provider" not in rendered + + +def test_worker_degrades_on_gateway_error_and_recovers_next_cycle(tmp_path): + config = _config(tmp_path) + store = ActiveMetadataConsumerStatusStore(config.status_file) + consumer = _Consumer( + [GatewayUnavailableError("database unavailable"), _batch(claimed=0, staged=0)] + ) + worker = ActiveMetadataConsumerWorker( + consumer, + config, + status_store=store, + clock=lambda: NOW, + ) + + assert worker.run_cycle() is None + assert store.read().state == "degraded" + assert store.read().last_error_code == "platform_unavailable" + + assert worker.run_cycle() is not None + assert store.read().state == "ready" + assert store.read().consecutive_gateway_failures == 0 + + +def test_health_and_liveness_separate_database_readiness_from_process(tmp_path): + config = _config(tmp_path) + store = ActiveMetadataConsumerStatusStore(config.status_file) + worker = ActiveMetadataConsumerWorker( + _Consumer([_batch(claimed=0, staged=0)]), + config, + status_store=store, + clock=lambda: NOW, + ) + worker.run_cycle() + + health, ready = evaluate_worker_health( + store, + max_age_seconds=30, + now=NOW + timedelta(seconds=10), + ) + liveness, alive = evaluate_worker_liveness( + store, + max_age_seconds=30, + now=NOW + timedelta(seconds=10), + ) + assert ready is alive is True + assert health["status"] == liveness["status"] == "healthy" + + stale, ready = evaluate_worker_health( + store, + max_age_seconds=5, + now=NOW + timedelta(seconds=10), + ) + assert ready is False + assert stale["reason"] == "status_stale" + + +def test_once_mode_stops_after_one_cycle_and_returns_failure_on_gateway_error( + tmp_path, +): + config = _config(tmp_path) + success = ActiveMetadataConsumerWorker( + _Consumer([_batch(claimed=0, staged=0)]), + config, + clock=lambda: NOW, + ) + assert success.run(once=True) == 0 + assert success.status is not None and success.status.state == "stopped" + + failed = ActiveMetadataConsumerWorker( + _Consumer([GatewayUnavailableError("database unavailable")]), + config, + clock=lambda: NOW, + ) + assert failed.run(once=True) == 1 + assert failed.status is not None and failed.status.state == "stopped" diff --git a/data_agent/test_metadata_fabric_active_metadata_consumer.py b/data_agent/test_metadata_fabric_active_metadata_consumer.py new file mode 100644 index 00000000..ac1f545e --- /dev/null +++ b/data_agent/test_metadata_fabric_active_metadata_consumer.py @@ -0,0 +1,57 @@ +import json +from copy import deepcopy + +from data_agent import metadata_fabric_active_metadata_consumer as consumer + + +def test_static_contract_binds_inert_request_staging_and_deployment(): + report = consumer.build_contract_report() + + assert report["status"] == "valid" + assert report["errors"] == [] + assert report["activation_boundary"] == ( + "durable_request_awaiting_authorization" + ) + assert report["deployment_expected_replicas"] == 0 + 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_evidence_is_current_content_bound_and_locally_scoped(): + evidence = json.loads( + consumer.DEFAULT_EVIDENCE_PATH.read_text(encoding="utf-8") + ) + + assert consumer.validate_rehearsal_evidence(evidence) == [] + assert evidence["processed_event_count"] == 2 + assert evidence["activation_request_count"] == 2 + assert evidence["platform_command_count"] == 0 + assert evidence["atomic_completion_guard_verified"] is True + assert evidence["requests_inert"] is True + assert evidence["deployment_applied"] is False + 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( + consumer.DEFAULT_EVIDENCE_PATH.read_text(encoding="utf-8") + ) + tampered = deepcopy(evidence) + tampered["platform_command_count"] = 1 + tampered["production_ready"] = True + + errors = consumer.validate_rehearsal_evidence(tampered) + + assert "Active Metadata consumer evidence SHA-256 does not match" in errors + assert ( + "local Active Metadata consumer evidence may not claim production_ready" + in errors + ) + assert "Active Metadata consumer must not create platform commands" in errors diff --git a/data_agent/test_platform_contracts.py b/data_agent/test_platform_contracts.py index c60f564f..fb6eae6c 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"] == ( - "099_active_metadata_change_outbox" + "100_active_metadata_activation_request" ) diff --git a/data_agent/test_platform_truth.py b/data_agent/test_platform_truth.py index 3c3f7526..5a130abd 100644 --- a/data_agent/test_platform_truth.py +++ b/data_agent/test_platform_truth.py @@ -247,6 +247,16 @@ 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"] == "active_metadata_consumer_worker" + and item["production_role"] == "activation_request_staging_only" + for item in static_report["runtime"]["inventory"] + ) + assert any( + item["runtime_id"] == "metadata_active_metadata_consumer_rehearsal" + and item["production_role"] == "local_verification_only" + for item in static_report["runtime"]["inventory"] + ) def test_runtime_report_detects_unregistered_background_mechanism(tmp_path): diff --git a/docs/architecture-decisions/adr-060-transactional-active-metadata-change-outbox.md b/docs/architecture-decisions/adr-060-transactional-active-metadata-change-outbox.md index 92bc40d1..2f5a691d 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 -`b4b9d77d716f0e5d62389378aec60bdc366b2d847ace07a6b476310eb3ff6732`, +`b6162ed687c6554859b0056e34dd5cf9f051818deb533daf708eeea67b18eca5`, and evidence fingerprint -`fcd44cf9043e144016644b0d04e72f348a777668d853081d9d21ded341bb2b59`. +`807444bb1adb2ea633bf98d845ba8a62fee608a4c6eff18c17c0c94f406e8a2`. 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 new file mode 100644 index 00000000..68d7e68c --- /dev/null +++ b/docs/architecture-decisions/adr-061-durable-inert-active-metadata-activation-request.md @@ -0,0 +1,103 @@ +# ADR-061: Durable inert Active Metadata activation request + +- 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 + +## Context + +ADR-060 establishes an authoritative `MetadataChangeEvent` when a new +ResourceVersion is registered. The event identifies a changed resource version +and a deterministic projection route, but it is not executable work. It does +not contain a PlatformDefinitionVersion, PlatformRun, execution-plan Artifact, +PolicyDecision, ApprovalRecord, or provider read-back contract. + +Turning that event directly into a DolphinScheduler `PlatformCommand` would +silently invent the missing authorization and execution context. Letting the +consumer call OpenMetadata or Gravitino would also combine observation, +authorization, scheduling, and provider mutation in one credential-bearing +process. + +## Options considered + +| Option | Benefit | Cost / risk | Decision | +|---|---|---|---| +| Convert each event directly to a `PlatformCommand` | Shortest path to a scheduler | Fabricates Definition/Run/policy/approval bindings and bypasses command authorization | Rejected | +| Let the consumer mutate metadata providers | Fewer durable hops | Gives the event consumer provider credentials and collapses policy, execution, and evidence boundaries | Rejected | +| Persist an inert activation request, then authorize it separately | Durable intent, replay safety, explicit authorization boundary | Adds one ledger object and a later promotion step | Accepted | + +## Decision + +Migration 100 adds the tenant-scoped +`gda_control.metadata_activation_request` ledger. A managed consumer claims an +event, derives a content-bound `MetadataActivationRequest`, and calls +`stage_metadata_activation_request`. The database persists the request and +marks the event processed in the same PostgreSQL transaction. Exact replay +returns the existing request; conflicting content or an expired/wrong worker +claim fails closed. + +Every new request has immutable status `awaiting_authorization` and fixes these +claims to `false`: + +- `provider_apply_authorized` +- `provider_mutations_executed` +- `production_scheduler_submission_verified` +- `production_ingestion_verified` +- `production_ready` + +The previous `complete_metadata_change` path may not process an event unless an +exact durable request already exists. The managed consumer owns only its +tenant-scoped PostgreSQL credential. It has no provider or scheduler secret, +no Kubernetes API token, and no provider mutation or command-submission client. +The base Deployment remains at zero replicas until an environment explicitly +supplies its database identity and scales it. + +A later authorization component must bind the request to a real Definition, +Run, execution plan, PolicyDecision, and Approval before the existing +DolphinScheduler command boundary can be used. That promotion is outside this +decision. + +## Rationale + +The request ledger preserves the fact that an event was routed without +pretending that routing authorized an action. Atomic stage-and-complete avoids +both processed events with no durable request and requests detached from event +delivery. Separating scheduler/provider credentials limits the consumer's +blast radius and keeps production claims evidence-driven. + +## Trade-offs + +- Provider ingestion is not immediate; an authorization/promotion controller + is still required. +- `awaiting_authorization` requests can accumulate, so production needs queue + age, retry/dead-letter, alert, and SLO ownership before scaling the worker. +- Database durability does not prove protected workload identity, scheduler + submission, provider mutation, or production ingestion. + +## Local verification + +Checked PostgreSQL evidence contains two processed events and exactly two +inert activation requests. It proves the no-request completion guard, exact +request replay, managed consumer staging, tenant isolation, forced RLS, +gateway direct UPDATE/DELETE denial, and zero `platform_command_outbox` rows. +The static deployment contract proves zero base replicas, no provider or +scheduler secret, disabled Kubernetes token mounting, and admission of the +consumer selector to PostgreSQL ingress. + +- Contract fingerprint: + `7a0bc0e3aaa53509443f031e7c40d7c5ff0f304510ddfeceedbad00147b505f5` +- Evidence fingerprint: + `480a2d651c1389feed471cb6ff082537ddfcd2f108e2652b9168b16ea4d02977` + +`deployment_applied`, `production_workload_identity_verified`, +`production_scheduler_submission_verified`, `production_ingestion_verified`, +and `production_ready` remain `false`. + +## Revisit triggers + +Revisit this decision only if an event schema itself becomes an authorized, +content-bound command carrying the complete Definition/Run/execution-plan/ +PolicyDecision/Approval chain, or if a different durable authorization service +can prove equivalent atomicity, tenant isolation, idempotency, and credential +separation. 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 new file mode 100644 index 00000000..5b2c63ac --- /dev/null +++ b/docs/evidence/metadata-fabric-active-metadata-consumer-2026-07-30.json @@ -0,0 +1,47 @@ +{ + "activation_request_count": 2, + "activation_route": "metadata_fabric.projection_plan", + "atomic_completion_guard_verified": true, + "contract_sha256": "7a0bc0e3aaa53509443f031e7c40d7c5ff0f304510ddfeceedbad00147b505f5", + "cross_tenant_read_blocked": true, + "deployment_applied": false, + "deployment_contract_verified": true, + "deployment_expected_replicas": 0, + "direct_request_mutation_blocked": true, + "errors": [], + "event_ids": [ + "3a7303a0-06a8-5f60-aa16-1abe23a050d5", + "52073ce1-db1f-526c-a6a7-51976c2ec8cd" + ], + "evidence_sha256": "480a2d651c1389feed471cb6ff082537ddfcd2f108e2652b9168b16ea4d02977", + "exact_request_replay_created": false, + "force_rls_verified": true, + "gateway_select_insert_only_verified": true, + "legacy_completion_without_request_blocked": true, + "local_postgresql_activation_request_staging_verified": true, + "managed_consumer_staged": true, + "platform_command_count": 0, + "processed_event_count": 2, + "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, + "request_absent_before_atomic_stage": true, + "request_ids": [ + "478851c6-d997-59fd-b19e-ebaefdf4225d", + "66970966-d9d4-5d94-aa4b-665331882fe7" + ], + "request_sha256": [ + "5b2ee475d663a0659b752dd4403952901e69385a36f5ad9a8bd2a9e4330dc394", + "bbbf6e5a5647eacd67b02bacca14b67cb68e8e6b0f11ebc714781cf7eaeb5520" + ], + "request_statuses": [ + "awaiting_authorization", + "awaiting_authorization" + ], + "requests_inert": true, + "schema": "gda.active_metadata_consumer_evidence.v1", + "status": "local_postgresql_activation_request_staging_verified" +} 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 544faaf8..22f57fb7 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": "b4b9d77d716f0e5d62389378aec60bdc366b2d847ace07a6b476310eb3ff6732", + "contract_sha256": "b6162ed687c6554859b0056e34dd5cf9f051818deb533daf708eeea67b18eca5", "cross_tenant_read_blocked": true, "errors": [], "event_id": "52073ce1-db1f-526c-a6a7-51976c2ec8cd", "event_sha256": "08087307bdd6694a3b2de2176fdb02ab7db720728b65b6da641e0b25dfa5edcf", - "evidence_sha256": "fcd44cf9043e144016644b0d04e72f348a777668d853081d9d21ded341bb2b59", + "evidence_sha256": "807444bb1adb2ea633bf98d845ba8a62fee608a4c6eff18c17c0c94f406e8a2c", "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 51aebd6a..654e612b 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-14(本地 Active Metadata 事务 outbox 已验证,生产验证待执行) +### 4.8 Metadata Fabric Bridge M1 + M2 + M3-15(本地 Active Metadata durable request consumer 已验证,生产验证待执行) 第八块回到 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,9 +233,10 @@ Temporal 继续保持目标组件状态,不在这一包并行接入。OpenMeta 28. [ADR-057](architecture-decisions/adr-057-production-object-store-readiness-gate.md) 已将 M3-10 evidence、S3-compatible provider/account/region/bucket、独立 failure domain、OIDC workload federation、精确八项 S3 permission、TLS/private path、KMS、versioning、cross-region replication、strong read/list consistency、tenant isolation、owner/SLO/runbook 和 26 项 protected attestation check 冻结为 fail-closed profile。当前 profile fingerprint 为 `668e194b3c688307014148391e7f389c9d6e9ca69c95d7b4cc92b4acae93181a`,report fingerprint 为 `85362dd10b7dc565f9fa567673d90b774cdec714bd1e70fb2c3c83c1af48b5ea`,合同有效但 43 项生产输入仍 blocked,全部 production claims 为 `false`。这只是 provider-neutral 决策和验收合同:没有选择、部署或验证 AWS S3、华为云 OBS 或其他生产对象存储;原生非 S3 provider 必须进入新的 conformance slice。 29. [ADR-058](architecture-decisions/adr-058-local-spark-commit-failure-recovery.md) 已在 Spark driver 的 loopback Iceberg REST proxy 中于 provider 转发前注入 HTTP 503。baseline 为 1 个 append snapshot、2 行和 1 个 referenced Parquet;失败调用经过精确 2 次 503 后,snapshot/row/file 均零漂移;对同一 `spark-recovery` 行做一次显式重试后为父子相连的 2 个 append snapshots、3 行和 2 个 referenced Parquet。直接 MinIO inventory 精确为 2 data + 3 metadata + 4 manifest = 9 objects,没有孤儿 data file;namespace、两块 PV 和 port-forward 均清理。contract fingerprint 为 `6d8944ab80246dc65891aa81118cb8b73f7ecad699be9a2af5e62d8260c41002`,evidence fingerprint 为 `39571cdac1e4043bcfc2d03a73b2b12ff925210daf8ae36bc640b8cb14d89401`。该结果只证明已知 pre-forward 失败下的本地原子性和一次显式重试,不证明 uncertain commit reconciliation、网络 exactly-once、生产对象存储或完整 engine conformance。 30. [ADR-059](architecture-decisions/adr-059-local-spark-uncertain-commit-reconciliation.md) 已将一个 armed commit 转发给 Gravitino,并在 provider 返回 200 后丢弃成功响应、向 Spark 返回 Iceberg `CommitStateUnknownException` 所需的 HTTP 504;一次传输重试被抑制。Spark 不重提逻辑写,而是 readback 得到父子相连的 2 个 append snapshots、3 行和 2 个 referenced Parquet,决策为 `committed_do_not_resubmit`、`write_resubmitted=false`。MinIO inventory 为 2 data + 3 metadata + 4 manifest = 9 objects;Job `Complete 1/1`,namespace、两块 PV 和 port-forward 均清理。contract fingerprint 为 `7a8d75a1d6b4558b982c6c3242d8d356c5046955f8aae7a45e5c297b6f4d4132`,evidence fingerprint 为 `d6462fff78d07047311b1f715d5f2c7f08c0ce8fbdd5c8b26a3d95ddc3474786`。该结果只证明一个本地 append 的确定性 readback/no-resubmit,不证明持久 reconcile controller、并发写、进程崩溃恢复、网络 exactly-once 或生产能力。 -31. [ADR-060](architecture-decisions/adr-060-transactional-active-metadata-change-outbox.md) 已新增 migration 099、内容绑定 `MetadataChangeEvent`、deterministic activation intent 与 PlatformGateway 原子注册/claim/fail/complete API。真实 PostgreSQL 16 演练中,ResourceVersion 与事件同事务创建,精确 replay 与 processed replay 均不新增事件;错误 consumer/worker 被拒绝,一次 retry 和一次强制租约过期后由第三个 worker 完成,旧 ResourceVersion 的补事件尝试整笔回滚。最终只有 1 条权威事件、3 次 attempt,FORCE RLS、跨租户拒绝与 gateway 无直接 UPDATE/DELETE 均通过。contract fingerprint 为 `b4b9d77d716f0e5d62389378aec60bdc366b2d847ace07a6b476310eb3ff6732`,evidence fingerprint 为 `fcd44cf9043e144016644b0d04e72f348a777668d853081d9d21ded341bb2b59`。该切片只生成 `metadata_fabric.projection_plan` 意图,不新增常驻 consumer、不提交 DolphinScheduler、不授权或执行 provider mutation,也不证明 production ingestion。 +31. [ADR-060](architecture-decisions/adr-060-transactional-active-metadata-change-outbox.md) 已新增 migration 099、内容绑定 `MetadataChangeEvent`、deterministic activation intent 与 PlatformGateway 原子注册/claim/fail/complete API。真实 PostgreSQL 16 演练中,ResourceVersion 与事件同事务创建,精确 replay 与 processed replay 均不新增事件;错误 consumer/worker 被拒绝,一次 retry 和一次强制租约过期后由第三个 worker 完成,旧 ResourceVersion 的补事件尝试整笔回滚。最终只有 1 条权威事件、3 次 attempt,FORCE RLS、跨租户拒绝与 gateway 无直接 UPDATE/DELETE 均通过。因共享 contract/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`。 -此处 M1 只证明静态合同和只读 HTTP 边界;M2a 只证明本地 live foundation 与 PVC 重挂载连续性;M2b-1/M2b-2 分别限定在同集群新 PVC 和同集群隔离 repository;M2b-3 的 `local_cross_cluster_recovery_verified=true` 只限定在 `local_same_host_distinct_kubernetes_clusters_external_s3_repository`;M2c-1/M2c-2/M2c-3 分别限定本地 provider metrics、临时双周期 OTel 和单 job scrape recovery;M2c-4/M2d-2 只证明 production observability/NetworkPolicy profile 与 attestation 合同可校验;M2d-1 只证明本地两节点 kindnet 的隔离合成流量;M3-1 的 terminal evidence 与 M3-2 的 PolicyDecision/Approval 仍是 deterministic local fixtures。M3-2 只把 projection 写入本地 provider 并证明 retained target 的单次零写入 replay;M3-3 只把该本地 evidence 对应的 binding 写入临时 GDA Control 账本;M3-4 只向无认证 loopback receiver 发送精确 candidate 并验证 503 后幂等恢复;M3-5 只证明 OpenMetadata 在 provider 强制默认 role 之上的项目新增 grant 限定为 `table/Create`,以及本地 JWT 轮换/吊销和越权拒绝;M3-6 只证明隔离 Gravitino Basic IdP 的 bounded table-create、catalog-create 拒绝、登录轮换/吊销和完整清理;M3-7 只证明 pending production identity profile、profile-bound attestation 和派生 claim 的 fail-closed 合同可校验,没有部署或证明真实身份路径;M3-8 只证明同一 Docker Desktop 集群内 Basic 用户、JDBC metadata 与 file warehouse PVC 在受控 Pod restart 后连续;M3-9 只证明同节点共享 RWO PVC 的 Spark interoperability;M3-10 移除了该共享 PVC,并证明同一 Docker Desktop 主机/集群内 Spark 与 MinIO 的跨节点 S3-compatible 互操作,但不证明生产云对象存储、独立 failure domain、持久 identity binding、Flink 或完整 engine conformance;M3-11 只冻结 provider-neutral production object-store profile、精确 attestation binding 与 fail-closed claims,没有选择 provider、部署 bucket/KMS/policy 或提交真实 attestation;M3-12 只证明同一本地路径的 pre-forward commit failure 不改变可见 table state,随后一次显式重试产生一个新 snapshot/row,且无孤儿 data file;M3-13 只证明单次本地 append 在 provider 200 响应丢失并映射为 commit-state-unknown 后,可以由即时 table readback 判定 committed 且不重提,不覆盖持久 controller、进程崩溃、并发写或任意 mutation;M3-14 只证明 ResourceVersion 注册与 Active Metadata 事件在本地 PostgreSQL 同事务创建,并验证租户/workload scoped claim/retry/complete,不包含常驻 consumer、DolphinScheduler submission、provider policy 或 provider mutation。M3-2 ingestion 仍使用 bootstrap admin,生产持久 binding、ResourceVersion 和 legacy authority 都未写入;生产对象存储、双 provider/生产最小权限、protected workload identity、OIDC、TLS、生产持久 catalog、tenant isolation、真实 receiver/alert/SLO、受保护 provider policy、生产故障注入、source-loss recovery、cancel/reconcile/lineage、完整 Spark/Flink conformance、生产 ingest、四项 production gate 和 `production_ready` 仍为 `false`。 +此处 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`。 ## 5. 重新评估条件 diff --git a/docs/system-of-record-matrix-2026-07-24.md b/docs/system-of-record-matrix-2026-07-24.md index 42b0c4aa..055a8657 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 已验证,生产 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 已验证,生产 provider ingestion、生产观测、生产 policy/tenant isolation、生产 identity/object-store attestation、生产 consumer/scheduler 和生产切换仍 `in_progress` -适用分支:`feat/ar1-metadata-fabric-active-metadata-outbox` +适用分支:`feat/ar1-metadata-fabric-active-metadata-consumer` ## 判定规则 @@ -18,9 +18,9 @@ | 事实域 | 当前权威/状态 | 非权威副本或投影 | 目标权威与迁移规则 | Owner | 阶段 | |---|---|---|---|---|---| | SQL schema 历史 | PostgreSQL `schema_migrations`,以完整 migration ID + checksum 为权威 | migration CLI 的 JSON 报告 | 保持现有 ledger;任何 drift fail closed | Data Platform | AR-0,已验证 | -| 部署配置策略 | Compose/K8s/进程环境;`platform_truth.CONFIG_SPECS` 定义关键类型与策略;DolphinScheduler worker 有默认零副本、外部 ConfigMap/Secret 驱动的 Kustomize 模板、静态 validator 和 staging activation preflight | `.env` 仅补默认;脱敏 snapshot、Secret key attestation、未扩容 Deployment 和 `ready_for_activation` 都是观测/模板 | 版本化 DeploymentProfile + secret reference;部署环境始终优先;模板或 preflight 通过都不等于环境已启用 | Platform/SRE/Security | AR-0,部分实现;worker 模板/preflight 本地已验证 | +| 部署配置策略 | 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 已登记但尚无生产调用方;Metadata Fabric recovery/metrics/policy/catalog/interoperability/failure/uncertain-commit 与 Active Metadata outbox rehearsal 均登记为 `local_verification_only`,不是 scheduler、worker、持续监控、生产 policy/catalog controller 或状态权威 | AST primitive report、worker status JSON、FrameworkAttemptObservation、DolphinScheduler instance state、本地 recovery/metrics/network-policy/catalog/interoperability/failure/outbox evidence | PlatformRun ledger 唯一登记最终状态;framework/provider attempt 只能回报观测;本地演练进程与 evidence 不得变成生产控制器、监控后端、catalog authority 或 tenant-isolation 权威 | Platform Architecture | AR-1 adapter/worker 本地已验证;metadata runner 仅本地验证 -> staging 控制链待接入 | +| 后台运行时清单 | `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 控制链待接入 | | 原始文件/对象 | 当前 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 与内容绑定事件同事务写入,并提供 tenant/workload scoped claim/retry/complete、租约和终态,当前没有常驻 consumer | command/metadata delivery status、消费者 claim、activation intent、worker status JSON、WebSocket 消息 | command/event 与源事实同事务入 outbox,幂等 consumer 交付;Active Metadata consumer 只生成 projection-plan intent 并通过 DolphinScheduler 提交授权工作,不能自行执行 provider mutation | Platform/Integrations/Metadata Platform | AR-1 command worker 代码与 M3-14 本地 outbox 已验证 -> staging managed consumer/scheduler 待部署 | +| 事件交付 | 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 待执行 | | 质量结果 | `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 待执行 | @@ -58,6 +58,7 @@ 13. `candidate_validated`、`registry_subject_bound`、GitHub provenance action 成功、CI artifact、离线 preflight、未独立 attested 的 live observation JSON 或人工批准都不能单独授权 production;缺少同一 source revision 的 OCI subject 独立验证、registry/live revision/identity/health/golden-slice 绑定及受保护 provenance 时,promotion 必须失败。 14. Metadata Fabric M1 只允许 OpenMetadata/Gravitino GET;M2 只执行本地 foundation/recovery/metrics/policy 演练或验证 production readiness profile;M3-1 只从 synthetic terminal evidence 生成 plan/candidate;M3-2 只允许 exact local PolicyDecision/Approval 后向本地 provider 写 projection;M3-3 只将同一 source evidence 经 PlatformGateway 写入临时 append-only binding ledger;M3-4 只经 tenant-scoped outbox 向无认证 loopback receiver 投递精确 candidate,并验证 at-least-once + receiver idempotency;M3-5 只证明临时 OpenMetadata bot 在 provider 强制 `DefaultBotRole` 之上的项目新增 grant 是 `table/Create`,并验证 policy-create 拒绝与本地 JWT 轮换/吊销;M3-6 只证明隔离 Gravitino Basic user 的 bounded table-create、catalog-create 拒绝、密码轮换/用户吊销和完整清理;M3-7 只冻结 production identity profile、精确 attestation binding 与 fail-closed 派生 claims,既不部署 identity path,也不提交真实 production attestation;M3-8 只证明 Docker Desktop 单集群中 Basic role、PostgreSQL JDBC metadata 与 file warehouse PVC 在受控 Pod restart 后连续;M3-9 只证明同节点共享 RWO PVC 的 Spark/Iceberg REST 互操作;M3-10 移除 Spark/Gravitino 共享 warehouse PVC,且只证明同一 Docker Desktop 主机/集群内跨节点 MinIO 的 read/write/schema evolution/snapshot/time travel 与对象级 metadata 一致;M3-11 只冻结 S3-compatible production profile、精确 attestation binding 与 fail-closed claims,既不选择/部署 provider,也不创建 bucket/KMS/policy 或提交真实 attestation;M3-12 只在 Spark driver loopback proxy 中于转发前注入 503,证明失败尝试零可见漂移、随后一次显式重试和直接对象清单无孤儿 data file;M3-13 只在 provider 200 响应丢失后以 HTTP 504 触发 commit-state-unknown,并对一个本地 append 做即时 readback/no-resubmit,不构成持久生产 reconcile controller、crash/concurrency proof 或网络 exactly-once。生产对象存储、identity/TLS、受保护环境故障注入、source-loss recovery、cancel/reconcile/lineage 与 Flink 仍未证明。Gravitino `1.3.0` Basic IdP 不算 OIDC,生产必须明确选择并证明 custom OIDC authenticator 或 identity-aware proxy。本地 bootstrap provisioner、Basic IdP、loopback/cluster HTTP、memory/file-backed JDBC catalog、同节点 PVC/MinIO、临时 identity/ledger/outbox、loopback receiver、pending profile、合成 attestation 和 local evidence 都不等于双 provider/生产最小权限、protected workload identity/OIDC、生产持久 catalog/binding、TLS、受保护 OpenLineage receiver、tenant isolation、alert/SLO、生产 ingestion/conformance 或生产写权威。 15. M3-14 只证明本地 PostgreSQL 16 中 ResourceVersion 与 `resource_version.registered` 事件同事务创建,以及 tenant/workload scoped claim、retry、lease reclaim 和 complete;它没有部署常驻 consumer、没有提交 DolphinScheduler、没有取得 provider apply 授权、没有执行 provider mutation,也不构成 production Active Metadata readiness。 +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 仍未验证。 ## 已建立的 AR-0/AR-1 entry 证据 @@ -94,7 +95,8 @@ - Metadata Fabric M3-11 已建立 production object-store profile/attestation gate;checked-in profile fingerprint 为 `668e194b3c688307014148391e7f389c9d6e9ca69c95d7b4cc92b4acae93181a`,report fingerprint 为 `85362dd10b7dc565f9fa567673d90b774cdec714bd1e70fb2c3c83c1af48b5ea`,`profile_valid=true`,43 项 provider/identity/transport/encryption/durability/consistency/tenancy/operations 外部输入以 blockers 暴露,`ready_for_protected_verification=false`、`attestation_valid=false`、`production_object_store_gate_passed=false`、`production_ready=false`。该合同绑定 M3-10 evidence,但没有选择或部署 provider;合成完整 attestation 只验证门禁逻辑,不计入生产证据。下一项真实证据是经 owner 批准并物化的 provider profile,以及来自 `production-object-store` 受保护环境、绑定当前 source/profile 并通过全部 26 项检查的 attestation。 - Metadata Fabric M3-12 已在本地 Spark driver loopback Iceberg REST proxy 中于 provider 转发前注入 HTTP 503。baseline 为 1 个 append snapshot、2 行和 1 个 referenced Parquet;失败写经 2 次 503 后 snapshot/row/file 零漂移;对同一逻辑行一次显式重试后为父子相连的 2 个 append snapshots、3 行和 2 个 referenced Parquet。直接 MinIO inventory 精确为 2 data + 3 metadata + 4 manifest = 9 objects,没有孤儿 data file;namespace、两块 PV 和 port-forward 均清理。contract fingerprint 为 `6d8944ab80246dc65891aa81118cb8b73f7ecad699be9a2af5e62d8260c41002`,evidence fingerprint 为 `39571cdac1e4043bcfc2d03a73b2b12ff925210daf8ae36bc640b8cb14d89401`。这只证明已知 pre-forward failure 的本地原子性与一次显式重试,不证明 provider uncertain outcome reconcile、网络 exactly-once、生产对象存储、cancel/lineage、Flink 或完整 Spark conformance。 - Metadata Fabric M3-13 已在本地 Spark driver loopback proxy 将 armed commit 转发给 Gravitino,并在 provider 200 后丢弃成功响应、返回 HTTP 504;Iceberg 将其映射为 commit-state-unknown,一次传输重试被抑制。Spark 只读 readback 后输出 `committed_do_not_resubmit` 和 `write_resubmitted=false`,最终为父子相连的 2 个 append snapshots、3 行、2 个 referenced Parquet;MinIO 为 2 data + 3 metadata + 4 manifest = 9 objects。Job `Complete 1/1`,namespace、两块 PV 和 port-forward 均清理。contract fingerprint 为 `7a8d75a1d6b4558b982c6c3242d8d356c5046955f8aae7a45e5c297b6f4d4132`,evidence fingerprint 为 `d6462fff78d07047311b1f715d5f2c7f08c0ce8fbdd5c8b26a3d95ddc3474786`。这只证明一个本地 append 的确定性 readback/no-resubmit;持久 controller、crash/concurrency、网络 exactly-once、生产对象存储、cancel/lineage、Flink 和完整 Spark conformance 仍未证明。 -- Metadata Fabric M3-14 已新增内容绑定 `MetadataChangeEvent`、migration 099 transactional outbox 与 PlatformGateway 原子注册/claim/fail/complete。真实 PostgreSQL 16 演练验证首次注册 `created=true`、pending/processed 精确 replay 均 `created=false`、错误 consumer/worker 拒绝、retry、lease-expiry reclaim、第三次 attempt 完成、processed 不再认领、旧 ResourceVersion 补事件整笔回滚、跨租户不可见、FORCE RLS 和 gateway 无直接 UPDATE/DELETE;最终只有 1 条权威事件。contract fingerprint 为 `b4b9d77d716f0e5d62389378aec60bdc366b2d847ace07a6b476310eb3ff6732`,evidence fingerprint 为 `fcd44cf9043e144016644b0d04e72f348a777668d853081d9d21ded341bb2b59`。激活意图只路由到 `metadata_fabric.projection_plan`,`provider_apply_authorized=false`、`provider_mutations_executed=false`、`production_ingestion_verified=false`、`production_scheduler_submission_verified=false`、`production_ready=false`。 +- 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`。 ## 下一验收证据 @@ -102,7 +104,7 @@ - 真实 provenance artifact verify、受保护 overlay 的 `verified_for_staging_apply` release report,以及 staging/production 的 schema、config/runtime snapshot、registry/live DeploymentRevision 绑定、release/live artifact attestation 和环境 compare 报告; - staging 的 migration role、应用 login membership、连接池 role/tenant 复位、双租户 API 和 success finalization 运行产物; - DolphinScheduler adapter 的真实 IAM/OIDC、service token provisioning/轮换、provider 最小权限、binding artifact staging 接入、managed outbox worker/provider callback 实际扩容部署、唯一 worker ID、status/lease 故障恢复和无双写证据; -- Active Metadata managed consumer 的受保护 workload identity、DolphinScheduler submission、幂等 projection execution、policy/approval gate、provider read-back、重试/死信告警与生产 SLO 证据; +- Active Metadata consumer 的 production scale-up、受保护 workload identity,以及 request 绑定 Definition/Run/execution plan/PolicyDecision/Approval 后的 DolphinScheduler submission、幂等 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/k8s/base/active-metadata-consumer.yaml b/k8s/base/active-metadata-consumer.yaml new file mode 100644 index 00000000..e00461eb --- /dev/null +++ b/k8s/base/active-metadata-consumer.yaml @@ -0,0 +1,142 @@ +apiVersion: apps/v1 +kind: Deployment +metadata: + name: gis-agent-active-metadata-consumer + namespace: gis-agent + labels: + app.kubernetes.io/name: gis-agent-active-metadata-consumer + app.kubernetes.io/component: worker +spec: + # The base stays inert until an environment supplies the tenant-scoped + # database login and explicitly scales the deployment. + replicas: 0 + selector: + matchLabels: + app.kubernetes.io/name: gis-agent-active-metadata-consumer + strategy: + type: RollingUpdate + rollingUpdate: + maxUnavailable: 0 + maxSurge: 1 + template: + metadata: + labels: + app.kubernetes.io/name: gis-agent-active-metadata-consumer + app.kubernetes.io/component: worker + spec: + serviceAccountName: gis-agent-active-metadata-consumer + automountServiceAccountToken: false + terminationGracePeriodSeconds: 60 + containers: + - name: worker + image: gis-data-agent:latest + imagePullPolicy: IfNotPresent + command: + - sh + - -ec + args: + - | + export ACTIVE_METADATA_CONSUMER_WORKER_ID="worker:active-metadata:${POD_UID}" + exec python -m data_agent.active_metadata_consumer_worker run + env: + - name: PYTHONDONTWRITEBYTECODE + value: "1" + - name: POD_UID + valueFrom: + fieldRef: + fieldPath: metadata.uid + - name: DATABASE_URL + valueFrom: + secretKeyRef: + name: gis-agent-active-metadata-consumer + key: database-url + - name: ACTIVE_METADATA_CONSUMER_ENABLED + value: "true" + - name: ACTIVE_METADATA_CONSUMER_TENANT_ID + valueFrom: + configMapKeyRef: + name: gis-agent-active-metadata-consumer + key: tenant-id + - name: ACTIVE_METADATA_CONSUMER_SUBJECT + valueFrom: + configMapKeyRef: + name: gis-agent-active-metadata-consumer + key: consumer-subject + - name: ACTIVE_METADATA_CONSUMER_STATUS_FILE + value: /var/run/gis-agent/active-metadata/status.json + - name: ACTIVE_METADATA_CONSUMER_BATCH_SIZE + value: "1" + - name: ACTIVE_METADATA_CONSUMER_LEASE_SECONDS + value: "60" + - name: ACTIVE_METADATA_CONSUMER_POLL_INTERVAL_SECONDS + value: "5" + - name: ACTIVE_METADATA_CONSUMER_HEALTH_MAX_AGE_SECONDS + value: "30" + volumeMounts: + - name: worker-runtime + mountPath: /var/run/gis-agent/active-metadata + startupProbe: + exec: + command: + - python + - -m + - data_agent.active_metadata_consumer_worker + - liveness + - --status-file + - /var/run/gis-agent/active-metadata/status.json + - --max-age-seconds + - "30" + periodSeconds: 5 + timeoutSeconds: 3 + failureThreshold: 120 + readinessProbe: + exec: + command: + - python + - -m + - data_agent.active_metadata_consumer_worker + - health + - --status-file + - /var/run/gis-agent/active-metadata/status.json + - --max-age-seconds + - "30" + periodSeconds: 10 + timeoutSeconds: 3 + failureThreshold: 3 + livenessProbe: + exec: + command: + - python + - -m + - data_agent.active_metadata_consumer_worker + - liveness + - --status-file + - /var/run/gis-agent/active-metadata/status.json + - --max-age-seconds + - "90" + periodSeconds: 30 + timeoutSeconds: 3 + failureThreshold: 3 + securityContext: + allowPrivilegeEscalation: false + readOnlyRootFilesystem: true + capabilities: + drop: ["ALL"] + resources: + requests: + cpu: 50m + memory: 128Mi + limits: + cpu: 500m + memory: 512Mi + volumes: + - name: worker-runtime + emptyDir: + medium: Memory + sizeLimit: 8Mi +--- +apiVersion: v1 +kind: ServiceAccount +metadata: + name: gis-agent-active-metadata-consumer + namespace: gis-agent diff --git a/k8s/base/kustomization.yaml b/k8s/base/kustomization.yaml index 33bd8dee..d6d6b683 100644 --- a/k8s/base/kustomization.yaml +++ b/k8s/base/kustomization.yaml @@ -23,6 +23,7 @@ resources: - app-service.yaml - outbox-worker.yaml - dolphinscheduler-command-worker.yaml + - active-metadata-consumer.yaml # ----- Local LLM backend (host Ollama) ----- - ollama-service.yaml diff --git a/k8s/base/networkpolicy.yaml b/k8s/base/networkpolicy.yaml index 4067c842..e2265a26 100644 --- a/k8s/base/networkpolicy.yaml +++ b/k8s/base/networkpolicy.yaml @@ -30,6 +30,9 @@ spec: - podSelector: matchLabels: app.kubernetes.io/name: gis-agent-dolphinscheduler-command-worker + - podSelector: + matchLabels: + app.kubernetes.io/name: gis-agent-active-metadata-consumer - podSelector: matchLabels: app.kubernetes.io/name: gis-agent-migrate diff --git a/scripts/metadata-fabric-active-metadata-consumer.sh b/scripts/metadata-fabric-active-metadata-consumer.sh new file mode 100755 index 00000000..94370827 --- /dev/null +++ b/scripts/metadata-fabric-active-metadata-consumer.sh @@ -0,0 +1,4 @@ +#!/usr/bin/env bash +set -euo pipefail + +python -m data_agent.metadata_fabric_active_metadata_consumer "$@"