From 9665847b0043e4ee905133cdc29c026ec84afe28 Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Tue, 18 Aug 2026 03:31:06 -0700 Subject: [PATCH 01/13] test: define data management evidence profile acceptance --- tests/test_data_management_evidence.py | 219 +++++++++++++++++++++++++ 1 file changed, 219 insertions(+) create mode 100644 tests/test_data_management_evidence.py diff --git a/tests/test_data_management_evidence.py b/tests/test_data_management_evidence.py new file mode 100644 index 0000000..1bc4657 --- /dev/null +++ b/tests/test_data_management_evidence.py @@ -0,0 +1,219 @@ +"""Buyer acceptance for framework-neutral data-management evidence profiles.""" + +from __future__ import annotations + +from fastapi.testclient import TestClient + +from sdp.api import app +from sdp.tenant_binding import PURPOSE_HEADER, SUBJECT_HEADER, TENANT_HEADER + + +client = TestClient(app) + + +def _headers( + *, + subject: str = "admin", + tenant: str = "demo", + purpose: str = "glossary_stewardship", +) -> dict[str, str]: + """Return an explicitly opted-in demo identity for one tenant.""" + + return { + SUBJECT_HEADER: subject, + TENANT_HEADER: tenant, + PURPOSE_HEADER: purpose, + } + + +def _create_catalog_dataset() -> str: + """Create one governed dataset parent and return its catalog identifier.""" + + response = client.post( + "/plane/catalog-objects", + headers=_headers(), + json={ + "object_kind": "catalog_dataset", + "object_slug": "billing-settlement-evidence", + "display_title": "Billing settlement evidence", + "definition_text": "Usage, invoice, payment, refund, and settlement evidence for reconciliation.", + "preferred_language": "en", + "steward_display_name": "Billing Data Steward", + }, + ) + assert response.status_code == 200, response.text + return response.json()["catalog_object"]["catalog_object_id"] + + +def test_buyer_builds_evidence_complete_data_management_profile() -> None: + """Owner, CDE, rule, and observation create an explainable complete profile.""" + + catalog_object_id = _create_catalog_dataset() + profile_url = f"/plane/catalog-objects/{catalog_object_id}/data-management-profile" + + initial = client.get( + profile_url, + headers=_headers(subject="analyst", purpose="catalog_browse"), + ) + assert initial.status_code == 200, initial.text + assert initial.json()["evidence_complete"] is False + assert initial.json()["factors"] == { + "data_owner_present": False, + "critical_data_element_present": False, + "quality_rule_present": False, + "quality_observation_present": False, + } + assert "data-owner-assignments" in initial.json()["customer_next_action"] + assert initial.json()["policy_decision_id"] + + owner = client.post( + f"/plane/catalog-objects/{catalog_object_id}/data-owner-assignments", + headers=_headers(), + json={ + "owner_subject": "billing-operations-owner", + "owner_display_name": "Billing Operations Owner", + "valid_from": "2026-08-18T00:00:00Z", + "evidence_reference": "https://evidence.example.test/decisions/billing-owner-2026", + "truth_status": "authoritative", + }, + ) + assert owner.status_code == 200, owner.text + assert owner.json()["data_owner_assignment"]["owner_display_name"] == "Billing Operations Owner" + + duplicate_owner = client.post( + f"/plane/catalog-objects/{catalog_object_id}/data-owner-assignments", + headers=_headers(), + json={ + "owner_subject": "billing-operations-owner", + "owner_display_name": "Billing Operations Owner", + "valid_from": "2026-08-18T00:00:00Z", + "evidence_reference": "https://evidence.example.test/decisions/billing-owner-2026", + "truth_status": "authoritative", + }, + ) + assert duplicate_owner.status_code == 400 + + cde = client.post( + f"/plane/catalog-objects/{catalog_object_id}/critical-data-elements", + headers=_headers(), + json={ + "element_key": "settlement_amount", + "display_name": "Settlement amount", + "definition_text": "Cash amount paid out by the commerce provider for one settlement.", + "data_classification": "restricted_financial", + "evidence_reference": "https://evidence.example.test/dictionary/settlement-amount", + "truth_status": "authoritative", + }, + ) + assert cde.status_code == 200, cde.text + critical_data_element_id = cde.json()["critical_data_element"]["critical_data_element_id"] + + duplicate_cde = client.post( + f"/plane/catalog-objects/{catalog_object_id}/critical-data-elements", + headers=_headers(), + json={ + "element_key": "settlement_amount", + "display_name": "Settlement amount duplicate", + "definition_text": "Duplicate definition that must fail closed.", + "data_classification": "restricted_financial", + "evidence_reference": "https://evidence.example.test/dictionary/duplicate", + "truth_status": "proposed", + }, + ) + assert duplicate_cde.status_code == 400 + + rule = client.post( + f"/plane/critical-data-elements/{critical_data_element_id}/quality-rules", + headers=_headers(), + json={ + "rule_code": "settlement_matches_expected_amount", + "rule_description": "Provider settlement plus provider fee equals the captured invoice amount.", + "metric_code": "reconciliation_difference", + "threshold_operator": "equal_to", + "threshold_value": "0", + "unit_code": "KRW", + "evidence_reference": "https://evidence.example.test/controls/three-way-reconciliation", + "truth_status": "authoritative", + }, + ) + assert rule.status_code == 200, rule.text + data_quality_rule_id = rule.json()["data_quality_rule"]["data_quality_rule_id"] + + duplicate_rule = client.post( + f"/plane/critical-data-elements/{critical_data_element_id}/quality-rules", + headers=_headers(), + json={ + "rule_code": "settlement_matches_expected_amount", + "rule_description": "Duplicate rule that must fail closed.", + "metric_code": "reconciliation_difference", + "threshold_operator": "equal_to", + "threshold_value": "0", + "unit_code": "KRW", + "evidence_reference": "https://evidence.example.test/controls/duplicate", + "truth_status": "proposed", + }, + ) + assert duplicate_rule.status_code == 400 + + observation = client.post( + f"/plane/quality-rules/{data_quality_rule_id}/observations", + headers=_headers(), + json={ + "source_observation_id": "reconciliation_run_2026_08_18", + "observed_value": "0", + "observed_at": "2026-08-18T01:00:00Z", + "quality_status": "passed", + "evidence_reference": "https://evidence.example.test/runs/reconciliation-2026-08-18", + "truth_status": "observed", + }, + ) + assert observation.status_code == 200, observation.text + assert observation.json()["data_quality_observation"]["quality_status"] == "passed" + + duplicate_observation = client.post( + f"/plane/quality-rules/{data_quality_rule_id}/observations", + headers=_headers(), + json={ + "source_observation_id": "reconciliation_run_2026_08_18", + "observed_value": "0", + "observed_at": "2026-08-18T01:00:00Z", + "quality_status": "passed", + "evidence_reference": "https://evidence.example.test/runs/reconciliation-2026-08-18", + "truth_status": "observed", + }, + ) + assert duplicate_observation.status_code == 400 + + complete = client.get( + profile_url, + headers=_headers(subject="analyst", purpose="catalog_browse"), + ) + assert complete.status_code == 200, complete.text + body = complete.json() + assert body["evidence_complete"] is True + assert body["factors"] == { + "data_owner_present": True, + "critical_data_element_present": True, + "quality_rule_present": True, + "quality_observation_present": True, + } + assert body["counts"] == { + "data_owner_assignments": 1, + "critical_data_elements": 1, + "data_quality_rules": 1, + "data_quality_observations": 1, + } + assert body["data_quality_observations"][0]["evidence_reference"].startswith("https://") + assert "evidence" in body["customer_next_action"].lower() + assert body["policy_decision_id"] + + +def test_data_management_profile_is_tenant_isolated() -> None: + """A foreign tenant cannot discover a catalog object's governance profile.""" + + catalog_object_id = _create_catalog_dataset() + response = client.get( + f"/plane/catalog-objects/{catalog_object_id}/data-management-profile", + headers=_headers(subject="external-analyst", tenant="external", purpose="catalog_browse"), + ) + assert response.status_code == 404 From 213d469c85936e343a90b9e46dcbf95b1216a21e Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Tue, 18 Aug 2026 03:31:38 -0700 Subject: [PATCH 02/13] test: define data management evidence migration contract --- tests/test_data_management_migration.py | 42 +++++++++++++++++++++++++ 1 file changed, 42 insertions(+) create mode 100644 tests/test_data_management_migration.py diff --git a/tests/test_data_management_migration.py b/tests/test_data_management_migration.py new file mode 100644 index 0000000..7a65c0f --- /dev/null +++ b/tests/test_data_management_migration.py @@ -0,0 +1,42 @@ +"""Static migration contract for the data-management evidence profile.""" + +from __future__ import annotations + +from pathlib import Path + + +MIGRATION = Path(__file__).resolve().parents[1] / "migrations" / "0003_data_management_evidence.sql" + + +def test_data_management_evidence_migration_declares_normalized_tables() -> None: + """The evidence profile persists normalized owner, CDE, rule, and observation rows.""" + + assert MIGRATION.exists(), "0003 data-management evidence migration is required" + sql = MIGRATION.read_text(encoding="utf-8") + + for table_name in ( + "data_owner_assignments", + "critical_data_elements", + "data_quality_rules", + "data_quality_observations", + ): + assert f"CREATE TABLE IF NOT EXISTS {table_name}" in sql + + assert "UNIQUE (catalog_object_id, owner_subject, valid_from)" in sql + assert "UNIQUE (catalog_object_id, element_key)" in sql + assert "UNIQUE (critical_data_element_id, rule_code)" in sql + assert "UNIQUE (data_quality_rule_id, source_observation_id)" in sql + + +def test_data_management_evidence_migration_preserves_authority_and_provenance() -> None: + """Every governance fact carries tenant, truth, time, and HTTPS evidence fields.""" + + assert MIGRATION.exists(), "0003 data-management evidence migration is required" + sql = MIGRATION.read_text(encoding="utf-8") + + assert sql.count("tenant_reference TEXT NOT NULL") >= 4 + assert sql.count("truth_status TEXT NOT NULL") >= 4 + assert sql.count("evidence_reference TEXT NOT NULL") >= 4 + assert "CHECK (truth_status IN ('authoritative', 'observed', 'inferred', 'proposed'))" in sql + assert "CHECK (evidence_reference LIKE 'https://%')" in sql + assert "FOREIGN KEY (catalog_object_id) REFERENCES catalog_objects" in sql From b718dc4ba81cdaba308fa98ece5bcc68d34662ff Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Tue, 18 Aug 2026 03:33:54 -0700 Subject: [PATCH 03/13] feat: add data management evidence contracts --- src/sdp_core/data_management_evidence.py | 275 +++++++++++++++++++++++ 1 file changed, 275 insertions(+) create mode 100644 src/sdp_core/data_management_evidence.py diff --git a/src/sdp_core/data_management_evidence.py b/src/sdp_core/data_management_evidence.py new file mode 100644 index 0000000..7c8ec9e --- /dev/null +++ b/src/sdp_core/data_management_evidence.py @@ -0,0 +1,275 @@ +"""Framework-neutral contracts for actionable data-management evidence. + +The contracts describe CWL-owned ownership, critical-element, quality-rule, and +observation facts. They do not reproduce DAMA-DMBOK or DCAM licensed prose, +challenge questions, official scoring criteria, or evidence lists. +""" + +from __future__ import annotations + +import re +from datetime import datetime, timezone +from decimal import Decimal +from typing import Literal + +from pydantic import BaseModel, Field, field_validator, model_validator + + +TruthStatus = Literal["authoritative", "observed", "inferred", "proposed"] +DataClassification = Literal[ + "public", + "internal", + "confidential", + "restricted_pii", + "restricted_financial", +] +ThresholdOperator = Literal[ + "equal_to", + "not_equal_to", + "greater_than", + "greater_than_or_equal_to", + "less_than", + "less_than_or_equal_to", +] +QualityStatus = Literal["passed", "failed", "warning", "unknown"] +_OPAQUE_REFERENCE_RE = re.compile(r"^[A-Za-z0-9._:@-]+$") + + +def _require_aware_utc(value: datetime, field_name: str) -> datetime: + """Require an aware timestamp and normalize it to UTC.""" + + if value.tzinfo is None or value.utcoffset() is None: + raise ValueError(f"{field_name} must include a timezone offset") + return value.astimezone(timezone.utc) + + +def _require_https(value: str, field_name: str) -> str: + """Require an HTTPS evidence reference rather than a local path or URI.""" + + if not value.startswith("https://"): + raise ValueError(f"{field_name} must be an https evidence reference") + return value + + +def _require_opaque_reference(value: str, field_name: str) -> str: + """Reject path-bearing or whitespace-bearing external identities.""" + + if _OPAQUE_REFERENCE_RE.fullmatch(value) is None: + raise ValueError(f"{field_name} must be an opaque identifier") + return value + + +class DataOwnerAssignmentDraft(BaseModel): + """Effective-dated data-owner assignment backed by explicit evidence.""" + + owner_subject: str = Field(min_length=1, max_length=256) + owner_display_name: str = Field(min_length=1, max_length=256) + valid_from: datetime + valid_to: datetime | None = None + evidence_reference: str = Field(min_length=9, max_length=1024) + truth_status: TruthStatus + + @field_validator("owner_subject") + @classmethod + def owner_subject_is_opaque(cls, value: str) -> str: + """Keep the external identity path-free and bounded.""" + + return _require_opaque_reference(value, "owner_subject") + + @field_validator("valid_from", "valid_to") + @classmethod + def timestamps_are_aware(cls, value: datetime | None, info) -> datetime | None: + """Normalize assignment business times to aware UTC.""" + + if value is None: + return value + return _require_aware_utc(value, info.field_name) + + @field_validator("evidence_reference") + @classmethod + def evidence_reference_is_https(cls, value: str) -> str: + """Require transport-protected owner-decision evidence.""" + + return _require_https(value, "evidence_reference") + + @model_validator(mode="after") + def valid_interval_is_ordered(self) -> "DataOwnerAssignmentDraft": + """Reject empty or reversed effective intervals.""" + + if self.valid_to is not None and self.valid_to <= self.valid_from: + raise ValueError("valid_to must be later than valid_from") + return self + + +class CriticalDataElementDraft(BaseModel): + """One critical data element defined within a catalog dataset.""" + + element_key: str = Field( + min_length=2, + max_length=128, + pattern=r"^[a-z][a-z0-9_]+$", + ) + display_name: str = Field(min_length=1, max_length=256) + definition_text: str = Field(min_length=10, max_length=4000) + data_classification: DataClassification + evidence_reference: str = Field(min_length=9, max_length=1024) + truth_status: TruthStatus + + @field_validator("evidence_reference") + @classmethod + def evidence_reference_is_https(cls, value: str) -> str: + """Require an HTTPS glossary, policy, or decision citation.""" + + return _require_https(value, "evidence_reference") + + +class DataQualityRuleDraft(BaseModel): + """Versioned quality expectation for one critical data element.""" + + rule_code: str = Field( + min_length=2, + max_length=128, + pattern=r"^[a-z][a-z0-9_]+$", + ) + rule_description: str = Field(min_length=10, max_length=4000) + metric_code: str = Field( + min_length=2, + max_length=128, + pattern=r"^[a-z][a-z0-9_]+$", + ) + threshold_operator: ThresholdOperator + threshold_value: Decimal + unit_code: str = Field(min_length=1, max_length=32, pattern=r"^[A-Za-z0-9._%-]+$") + evidence_reference: str = Field(min_length=9, max_length=1024) + truth_status: TruthStatus + + @field_validator("threshold_value") + @classmethod + def threshold_is_finite(cls, value: Decimal) -> Decimal: + """Reject NaN and infinite thresholds from financial or quality logic.""" + + if not value.is_finite(): + raise ValueError("threshold_value must be finite") + return value + + @field_validator("evidence_reference") + @classmethod + def evidence_reference_is_https(cls, value: str) -> str: + """Require an HTTPS control-definition citation.""" + + return _require_https(value, "evidence_reference") + + +class DataQualityObservationDraft(BaseModel): + """Immutable observed result for one quality rule.""" + + source_observation_id: str = Field(min_length=2, max_length=256) + observed_value: Decimal + observed_at: datetime + quality_status: QualityStatus + evidence_reference: str = Field(min_length=9, max_length=1024) + truth_status: TruthStatus + + @field_validator("source_observation_id") + @classmethod + def observation_identity_is_opaque(cls, value: str) -> str: + """Keep the producer's replay identity path-free.""" + + return _require_opaque_reference(value, "source_observation_id") + + @field_validator("observed_value") + @classmethod + def observed_value_is_finite(cls, value: Decimal) -> Decimal: + """Reject NaN and infinite observed values.""" + + if not value.is_finite(): + raise ValueError("observed_value must be finite") + return value + + @field_validator("observed_at") + @classmethod + def observed_at_is_aware(cls, value: datetime) -> datetime: + """Normalize observation time to aware UTC.""" + + return _require_aware_utc(value, "observed_at") + + @field_validator("evidence_reference") + @classmethod + def evidence_reference_is_https(cls, value: str) -> str: + """Require an HTTPS run, receipt, or measurement citation.""" + + return _require_https(value, "evidence_reference") + + +class DataOwnerAssignmentRecord(DataOwnerAssignmentDraft): + """Persisted owner assignment.""" + + data_owner_assignment_id: str + catalog_object_id: str + tenant_reference: str + recorded_at: datetime + + +class CriticalDataElementRecord(CriticalDataElementDraft): + """Persisted critical data element.""" + + critical_data_element_id: str + catalog_object_id: str + tenant_reference: str + recorded_at: datetime + + +class DataQualityRuleRecord(DataQualityRuleDraft): + """Persisted quality rule.""" + + data_quality_rule_id: str + critical_data_element_id: str + catalog_object_id: str + tenant_reference: str + recorded_at: datetime + + +class DataQualityObservationRecord(DataQualityObservationDraft): + """Persisted immutable quality observation.""" + + data_quality_observation_id: str + data_quality_rule_id: str + critical_data_element_id: str + catalog_object_id: str + tenant_reference: str + recorded_at: datetime + + +class DataManagementMutationEnvelope(BaseModel): + """Buyer-facing mutation result with policy evidence and next action.""" + + status: str + tenant_reference: str + catalog_object_id: str + policy_decision_id: str + pii_handling: str = "usable_purpose_limited_no_masking" + customer_next_action: str + data_owner_assignment: DataOwnerAssignmentRecord | None = None + critical_data_element: CriticalDataElementRecord | None = None + data_quality_rule: DataQualityRuleRecord | None = None + data_quality_observation: DataQualityObservationRecord | None = None + generated_at: datetime = Field(default_factory=lambda: datetime.now(timezone.utc)) + + +class DataManagementProfile(BaseModel): + """Explainable evidence-completeness profile for one catalog dataset.""" + + status: str = "data_management_profile" + tenant_reference: str + catalog_object_id: str + policy_decision_id: str + pii_handling: str = "usable_purpose_limited_no_masking" + evidence_complete: bool + factors: dict[str, bool] + counts: dict[str, int] + customer_next_action: str + data_owner_assignments: list[DataOwnerAssignmentRecord] = Field(default_factory=list) + critical_data_elements: list[CriticalDataElementRecord] = Field(default_factory=list) + data_quality_rules: list[DataQualityRuleRecord] = Field(default_factory=list) + data_quality_observations: list[DataQualityObservationRecord] = Field(default_factory=list) + generated_at: datetime = Field(default_factory=lambda: datetime.now(timezone.utc)) From fca89bfe005b58e9f23c78d5c8700117b5d83459 Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Tue, 18 Aug 2026 03:34:29 -0700 Subject: [PATCH 04/13] feat: add normalized data management evidence schema --- migrations/0003_data_management_evidence.sql | 93 ++++++++++++++++++++ 1 file changed, 93 insertions(+) create mode 100644 migrations/0003_data_management_evidence.sql diff --git a/migrations/0003_data_management_evidence.sql b/migrations/0003_data_management_evidence.sql new file mode 100644 index 0000000..eac18e9 --- /dev/null +++ b/migrations/0003_data_management_evidence.sql @@ -0,0 +1,93 @@ +-- Framework-neutral data-management evidence profile (3NF, 2+ word snake_case). +-- +-- This migration stores CWL-owned ownership, critical-element, quality-rule, +-- and observation facts. It deliberately does not reproduce licensed +-- DAMA-DMBOK or DCAM framework prose, questions, scoring, or evidence lists. + +SET search_path = public; + +CREATE TABLE IF NOT EXISTS data_owner_assignments ( + data_owner_assignment_id TEXT PRIMARY KEY, + catalog_object_id TEXT NOT NULL REFERENCES catalog_objects (catalog_object_id), + tenant_reference TEXT NOT NULL, + owner_subject TEXT NOT NULL, + owner_display_name TEXT NOT NULL, + valid_from TIMESTAMPTZ NOT NULL, + valid_to TIMESTAMPTZ, + evidence_reference TEXT NOT NULL CHECK (evidence_reference LIKE 'https://%'), + truth_status TEXT NOT NULL + CHECK (truth_status IN ('authoritative', 'observed', 'inferred', 'proposed')), + recorded_at TIMESTAMPTZ NOT NULL DEFAULT now(), + UNIQUE (catalog_object_id, owner_subject, valid_from), + CHECK (valid_to IS NULL OR valid_to > valid_from) +); +CREATE INDEX IF NOT EXISTS data_owner_assignments_object_idx + ON data_owner_assignments (tenant_reference, catalog_object_id); +CREATE INDEX IF NOT EXISTS data_owner_assignments_active_idx + ON data_owner_assignments (tenant_reference, catalog_object_id, valid_to); + +CREATE TABLE IF NOT EXISTS critical_data_elements ( + critical_data_element_id TEXT PRIMARY KEY, + catalog_object_id TEXT NOT NULL REFERENCES catalog_objects (catalog_object_id), + tenant_reference TEXT NOT NULL, + element_key TEXT NOT NULL, + display_name TEXT NOT NULL, + definition_text TEXT NOT NULL, + data_classification TEXT NOT NULL, + evidence_reference TEXT NOT NULL CHECK (evidence_reference LIKE 'https://%'), + truth_status TEXT NOT NULL + CHECK (truth_status IN ('authoritative', 'observed', 'inferred', 'proposed')), + recorded_at TIMESTAMPTZ NOT NULL DEFAULT now(), + UNIQUE (catalog_object_id, element_key) +); +CREATE INDEX IF NOT EXISTS critical_data_elements_object_idx + ON critical_data_elements (tenant_reference, catalog_object_id); + +CREATE TABLE IF NOT EXISTS data_quality_rules ( + data_quality_rule_id TEXT PRIMARY KEY, + critical_data_element_id TEXT NOT NULL + REFERENCES critical_data_elements (critical_data_element_id), + catalog_object_id TEXT NOT NULL REFERENCES catalog_objects (catalog_object_id), + tenant_reference TEXT NOT NULL, + rule_code TEXT NOT NULL, + rule_description TEXT NOT NULL, + metric_code TEXT NOT NULL, + threshold_operator TEXT NOT NULL, + threshold_value NUMERIC NOT NULL, + unit_code TEXT NOT NULL, + evidence_reference TEXT NOT NULL CHECK (evidence_reference LIKE 'https://%'), + truth_status TEXT NOT NULL + CHECK (truth_status IN ('authoritative', 'observed', 'inferred', 'proposed')), + recorded_at TIMESTAMPTZ NOT NULL DEFAULT now(), + UNIQUE (critical_data_element_id, rule_code) +); +CREATE INDEX IF NOT EXISTS data_quality_rules_element_idx + ON data_quality_rules (tenant_reference, critical_data_element_id); +CREATE INDEX IF NOT EXISTS data_quality_rules_object_idx + ON data_quality_rules (tenant_reference, catalog_object_id); + +CREATE TABLE IF NOT EXISTS data_quality_observations ( + data_quality_observation_id TEXT PRIMARY KEY, + data_quality_rule_id TEXT NOT NULL REFERENCES data_quality_rules (data_quality_rule_id), + critical_data_element_id TEXT NOT NULL + REFERENCES critical_data_elements (critical_data_element_id), + catalog_object_id TEXT NOT NULL REFERENCES catalog_objects (catalog_object_id), + tenant_reference TEXT NOT NULL, + source_observation_id TEXT NOT NULL, + observed_value NUMERIC NOT NULL, + observed_at TIMESTAMPTZ NOT NULL, + quality_status TEXT NOT NULL, + evidence_reference TEXT NOT NULL CHECK (evidence_reference LIKE 'https://%'), + truth_status TEXT NOT NULL + CHECK (truth_status IN ('authoritative', 'observed', 'inferred', 'proposed')), + recorded_at TIMESTAMPTZ NOT NULL DEFAULT now(), + UNIQUE (data_quality_rule_id, source_observation_id) +); +CREATE INDEX IF NOT EXISTS data_quality_observations_rule_idx + ON data_quality_observations (tenant_reference, data_quality_rule_id, observed_at); +CREATE INDEX IF NOT EXISTS data_quality_observations_object_idx + ON data_quality_observations (tenant_reference, catalog_object_id, observed_at); + +INSERT INTO schema_migrations (migration_id) +VALUES ('0003_data_management_evidence') +ON CONFLICT (migration_id) DO NOTHING; From 6cdd29e7c4752c3580978d8b22136c395ff876ef Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Tue, 18 Aug 2026 03:36:39 -0700 Subject: [PATCH 05/13] feat: add data management evidence stores --- src/sdp/data_management_store.py | 794 +++++++++++++++++++++++++++++++ 1 file changed, 794 insertions(+) create mode 100644 src/sdp/data_management_store.py diff --git a/src/sdp/data_management_store.py b/src/sdp/data_management_store.py new file mode 100644 index 0000000..53fb647 --- /dev/null +++ b/src/sdp/data_management_store.py @@ -0,0 +1,794 @@ +"""Persistence backends for data-management evidence profiles. + +The in-memory backend is the CI/pytest default. When ``SDP_DATABASE_DSN`` is +set, the relational backend uses migration 0003. Both backends enforce the +same tenant and uniqueness contracts and preserve observations append-only. +""" + +from __future__ import annotations + +from abc import ABC, abstractmethod +from copy import deepcopy +from datetime import datetime, timezone +from decimal import Decimal +from pathlib import Path +from threading import RLock +from typing import Any, TypeAlias + +from sdp.config import load_bootstrap +from sdp_core.data_management_evidence import ( + CriticalDataElementRecord, + DataOwnerAssignmentRecord, + DataQualityObservationRecord, + DataQualityRuleRecord, +) + +try: + from sqlalchemy import create_engine, text + from sqlalchemy.engine import Engine + from sqlalchemy.exc import IntegrityError +except ImportError: # pragma: no cover - optional graph extra + create_engine = None # type: ignore[assignment] + text = None # type: ignore[assignment] + Engine = Any # type: ignore[misc,assignment] + + class IntegrityError(Exception): # type: ignore[no-redef] + """Placeholder used only when SQLAlchemy is not installed.""" + + +ProfileRows: TypeAlias = tuple[ + list[DataOwnerAssignmentRecord], + list[CriticalDataElementRecord], + list[DataQualityRuleRecord], + list[DataQualityObservationRecord], +] + +_MIGRATION_0003 = ( + Path(__file__).resolve().parents[2] / "migrations" / "0003_data_management_evidence.sql" +) +_GRAPH_EXTRA_HINT = ( + "Install the optional graph extra so SQLAlchemy can open SDP_DATABASE_DSN." +) +_MEMORY_OWNER_ROWS: dict[str, DataOwnerAssignmentRecord] = {} +_MEMORY_ELEMENT_ROWS: dict[str, CriticalDataElementRecord] = {} +_MEMORY_RULE_ROWS: dict[str, DataQualityRuleRecord] = {} +_MEMORY_OBSERVATION_ROWS: dict[str, DataQualityObservationRecord] = {} +_MEMORY_LOCK = RLock() + + +def _as_datetime(value: Any) -> datetime: + """Coerce a driver timestamp or ISO string into an aware UTC datetime.""" + + if isinstance(value, datetime): + if value.tzinfo is None: + return value.replace(tzinfo=timezone.utc) + return value.astimezone(timezone.utc) + return datetime.fromisoformat(str(value).replace("Z", "+00:00")).astimezone(timezone.utc) + + +def _sql_timestamp(value: datetime | None) -> str | None: + """Render an optional timestamp for Postgres TIMESTAMPTZ and SQLite TEXT.""" + + if value is None: + return None + return _as_datetime(value).isoformat() + + +def _open_engine(database_dsn: str) -> Engine: + """Open a SQLAlchemy engine or fail loud when the graph extra is absent.""" + + if create_engine is None: + raise RuntimeError(_GRAPH_EXTRA_HINT) + return create_engine(database_dsn, future=True) + + +def _decimal(value: Any) -> Decimal: + """Coerce database numeric output into an exact Decimal.""" + + return Decimal(str(value)) + + +def _raise_unique_violation(exc: Exception, fallback: str) -> None: + """Map relational uniqueness failures to stable buyer-facing errors.""" + + detail = str(exc).lower() + if "data_owner_assignments" in detail or "owner_subject" in detail: + raise ValueError("duplicate data-owner assignment for this catalog object") from exc + if "critical_data_elements" in detail or "element_key" in detail: + raise ValueError("duplicate critical data element in this catalog object") from exc + if "data_quality_rules" in detail or "rule_code" in detail: + raise ValueError("duplicate data-quality rule for this critical data element") from exc + if "data_quality_observations" in detail or "source_observation_id" in detail: + raise ValueError("duplicate data-quality observation for this rule") from exc + raise ValueError(fallback) from exc + + +def data_management_sqlite_ddl() -> str: + """Return migration 0003 rewritten for SQLite unit-test execution.""" + + sql_text = _MIGRATION_0003.read_text(encoding="utf-8") + rewritten = ( + sql_text.replace("SET search_path = public;", "") + .replace("TIMESTAMPTZ", "TEXT") + .replace("DEFAULT now()", "DEFAULT CURRENT_TIMESTAMP") + ) + statements: list[str] = [] + buffer: list[str] = [] + for line in rewritten.splitlines(): + buffer.append(line) + if line.strip().endswith(";"): + statement = "\n".join(buffer).strip() + buffer = [] + if statement and "schema_migrations" not in statement: + statements.append(statement) + tail = "\n".join(buffer).strip() + if tail and "schema_migrations" not in tail: + statements.append(tail) + return "\n\n".join(statements) + "\n" + + +def apply_data_management_sqlite_schema(engine: Engine) -> None: + """Create migration 0003 table names on a SQLite test engine.""" + + if text is None: # pragma: no cover - SQLAlchemy is a dev extra + raise RuntimeError(_GRAPH_EXTRA_HINT) + statements: list[str] = [] + buffer: list[str] = [] + for line in data_management_sqlite_ddl().splitlines(): + stripped = line.strip() + if stripped.startswith("--"): + continue + buffer.append(line) + if stripped.endswith(";"): + statements.append("\n".join(buffer).strip()) + buffer = [] + tail = "\n".join(buffer).strip() + if tail: + statements.append(tail) + with engine.begin() as connection: + for statement in statements: + if statement: + connection.execute(text(statement)) + + +class DataManagementStore(ABC): + """Tenant-scoped persistence surface for evidence-profile records.""" + + @abstractmethod + def insert_owner_assignment(self, record: DataOwnerAssignmentRecord) -> None: + """Persist one effective-dated owner assignment.""" + + @abstractmethod + def insert_critical_data_element(self, record: CriticalDataElementRecord) -> None: + """Persist one critical data element.""" + + @abstractmethod + def insert_data_quality_rule(self, record: DataQualityRuleRecord) -> None: + """Persist one quality rule whose parent belongs to the same tenant.""" + + @abstractmethod + def insert_data_quality_observation(self, record: DataQualityObservationRecord) -> None: + """Persist one immutable quality observation.""" + + @abstractmethod + def get_critical_data_element( + self, + *, + tenant_reference: str, + critical_data_element_id: str, + ) -> CriticalDataElementRecord | None: + """Return one tenant-owned critical data element.""" + + @abstractmethod + def get_data_quality_rule( + self, + *, + tenant_reference: str, + data_quality_rule_id: str, + ) -> DataQualityRuleRecord | None: + """Return one tenant-owned quality rule.""" + + @abstractmethod + def profile_rows( + self, + *, + tenant_reference: str, + catalog_object_id: str, + ) -> ProfileRows: + """Return every evidence row for one tenant-owned catalog object.""" + + +class InMemoryDataManagementStore(DataManagementStore): + """Process-local evidence store used when no database DSN is configured.""" + + def insert_owner_assignment(self, record: DataOwnerAssignmentRecord) -> None: + """Insert an owner assignment after enforcing its natural key.""" + + with _MEMORY_LOCK: + if any( + row.catalog_object_id == record.catalog_object_id + and row.owner_subject == record.owner_subject + and row.valid_from == record.valid_from + for row in _MEMORY_OWNER_ROWS.values() + ): + raise ValueError("duplicate data-owner assignment for this catalog object") + _MEMORY_OWNER_ROWS[record.data_owner_assignment_id] = record.model_copy(deep=True) + + def insert_critical_data_element(self, record: CriticalDataElementRecord) -> None: + """Insert a CDE after enforcing the tenant-local element key.""" + + with _MEMORY_LOCK: + if any( + row.catalog_object_id == record.catalog_object_id + and row.element_key == record.element_key + for row in _MEMORY_ELEMENT_ROWS.values() + ): + raise ValueError("duplicate critical data element in this catalog object") + _MEMORY_ELEMENT_ROWS[record.critical_data_element_id] = record.model_copy(deep=True) + + def insert_data_quality_rule(self, record: DataQualityRuleRecord) -> None: + """Insert a rule only when its CDE exists in the same tenant.""" + + with _MEMORY_LOCK: + parent = _MEMORY_ELEMENT_ROWS.get(record.critical_data_element_id) + if parent is None or parent.tenant_reference != record.tenant_reference: + raise KeyError("critical data element not found in this tenant") + if any( + row.critical_data_element_id == record.critical_data_element_id + and row.rule_code == record.rule_code + for row in _MEMORY_RULE_ROWS.values() + ): + raise ValueError("duplicate data-quality rule for this critical data element") + _MEMORY_RULE_ROWS[record.data_quality_rule_id] = record.model_copy(deep=True) + + def insert_data_quality_observation(self, record: DataQualityObservationRecord) -> None: + """Append an observation only when its rule exists in the same tenant.""" + + with _MEMORY_LOCK: + parent = _MEMORY_RULE_ROWS.get(record.data_quality_rule_id) + if parent is None or parent.tenant_reference != record.tenant_reference: + raise KeyError("data-quality rule not found in this tenant") + if any( + row.data_quality_rule_id == record.data_quality_rule_id + and row.source_observation_id == record.source_observation_id + for row in _MEMORY_OBSERVATION_ROWS.values() + ): + raise ValueError("duplicate data-quality observation for this rule") + _MEMORY_OBSERVATION_ROWS[record.data_quality_observation_id] = record.model_copy( + deep=True + ) + + def get_critical_data_element( + self, + *, + tenant_reference: str, + critical_data_element_id: str, + ) -> CriticalDataElementRecord | None: + """Return a deep copy of one tenant-owned CDE.""" + + with _MEMORY_LOCK: + row = _MEMORY_ELEMENT_ROWS.get(critical_data_element_id) + if row is None or row.tenant_reference != tenant_reference: + return None + return row.model_copy(deep=True) + + def get_data_quality_rule( + self, + *, + tenant_reference: str, + data_quality_rule_id: str, + ) -> DataQualityRuleRecord | None: + """Return a deep copy of one tenant-owned quality rule.""" + + with _MEMORY_LOCK: + row = _MEMORY_RULE_ROWS.get(data_quality_rule_id) + if row is None or row.tenant_reference != tenant_reference: + return None + return row.model_copy(deep=True) + + def profile_rows( + self, + *, + tenant_reference: str, + catalog_object_id: str, + ) -> ProfileRows: + """Return deep copies of all rows attached to one catalog object.""" + + with _MEMORY_LOCK: + owners = [ + row.model_copy(deep=True) + for row in _MEMORY_OWNER_ROWS.values() + if row.tenant_reference == tenant_reference + and row.catalog_object_id == catalog_object_id + ] + elements = [ + row.model_copy(deep=True) + for row in _MEMORY_ELEMENT_ROWS.values() + if row.tenant_reference == tenant_reference + and row.catalog_object_id == catalog_object_id + ] + rules = [ + row.model_copy(deep=True) + for row in _MEMORY_RULE_ROWS.values() + if row.tenant_reference == tenant_reference + and row.catalog_object_id == catalog_object_id + ] + observations = [ + row.model_copy(deep=True) + for row in _MEMORY_OBSERVATION_ROWS.values() + if row.tenant_reference == tenant_reference + and row.catalog_object_id == catalog_object_id + ] + owners.sort(key=lambda row: (row.valid_from, row.data_owner_assignment_id)) + elements.sort(key=lambda row: (row.element_key, row.critical_data_element_id)) + rules.sort(key=lambda row: (row.rule_code, row.data_quality_rule_id)) + observations.sort(key=lambda row: (row.observed_at, row.source_observation_id)) + return owners, elements, rules, observations + + +class RelationalDataManagementStore(DataManagementStore): + """SQLAlchemy mapping onto migration 0003 tables.""" + + def __init__(self, database_dsn: str) -> None: + """Open the configured relational database.""" + + self._engine = _open_engine(database_dsn) + + def insert_owner_assignment(self, record: DataOwnerAssignmentRecord) -> None: + """Insert one owner assignment.""" + + try: + with self._engine.begin() as conn: + conn.execute( + text( + """ + INSERT INTO data_owner_assignments ( + data_owner_assignment_id, catalog_object_id, tenant_reference, + owner_subject, owner_display_name, valid_from, valid_to, + evidence_reference, truth_status, recorded_at + ) VALUES ( + :data_owner_assignment_id, :catalog_object_id, :tenant_reference, + :owner_subject, :owner_display_name, :valid_from, :valid_to, + :evidence_reference, :truth_status, :recorded_at + ) + """ + ), + { + "data_owner_assignment_id": record.data_owner_assignment_id, + "catalog_object_id": record.catalog_object_id, + "tenant_reference": record.tenant_reference, + "owner_subject": record.owner_subject, + "owner_display_name": record.owner_display_name, + "valid_from": _sql_timestamp(record.valid_from), + "valid_to": _sql_timestamp(record.valid_to), + "evidence_reference": record.evidence_reference, + "truth_status": record.truth_status, + "recorded_at": _sql_timestamp(record.recorded_at), + }, + ) + except IntegrityError as exc: + _raise_unique_violation(exc, "data-owner assignment could not be stored") + + def insert_critical_data_element(self, record: CriticalDataElementRecord) -> None: + """Insert one critical data element.""" + + try: + with self._engine.begin() as conn: + conn.execute( + text( + """ + INSERT INTO critical_data_elements ( + critical_data_element_id, catalog_object_id, tenant_reference, + element_key, display_name, definition_text, + data_classification, evidence_reference, truth_status, recorded_at + ) VALUES ( + :critical_data_element_id, :catalog_object_id, :tenant_reference, + :element_key, :display_name, :definition_text, + :data_classification, :evidence_reference, :truth_status, :recorded_at + ) + """ + ), + { + "critical_data_element_id": record.critical_data_element_id, + "catalog_object_id": record.catalog_object_id, + "tenant_reference": record.tenant_reference, + "element_key": record.element_key, + "display_name": record.display_name, + "definition_text": record.definition_text, + "data_classification": record.data_classification, + "evidence_reference": record.evidence_reference, + "truth_status": record.truth_status, + "recorded_at": _sql_timestamp(record.recorded_at), + }, + ) + except IntegrityError as exc: + _raise_unique_violation(exc, "critical data element could not be stored") + + def insert_data_quality_rule(self, record: DataQualityRuleRecord) -> None: + """Insert one quality rule after checking the tenant-owned parent CDE.""" + + try: + with self._engine.begin() as conn: + parent = conn.execute( + text( + """ + SELECT critical_data_element_id + FROM critical_data_elements + WHERE critical_data_element_id = :critical_data_element_id + AND tenant_reference = :tenant_reference + """ + ), + { + "critical_data_element_id": record.critical_data_element_id, + "tenant_reference": record.tenant_reference, + }, + ).first() + if parent is None: + raise KeyError("critical data element not found in this tenant") + conn.execute( + text( + """ + INSERT INTO data_quality_rules ( + data_quality_rule_id, critical_data_element_id, + catalog_object_id, tenant_reference, rule_code, + rule_description, metric_code, threshold_operator, + threshold_value, unit_code, evidence_reference, + truth_status, recorded_at + ) VALUES ( + :data_quality_rule_id, :critical_data_element_id, + :catalog_object_id, :tenant_reference, :rule_code, + :rule_description, :metric_code, :threshold_operator, + :threshold_value, :unit_code, :evidence_reference, + :truth_status, :recorded_at + ) + """ + ), + { + "data_quality_rule_id": record.data_quality_rule_id, + "critical_data_element_id": record.critical_data_element_id, + "catalog_object_id": record.catalog_object_id, + "tenant_reference": record.tenant_reference, + "rule_code": record.rule_code, + "rule_description": record.rule_description, + "metric_code": record.metric_code, + "threshold_operator": record.threshold_operator, + "threshold_value": str(record.threshold_value), + "unit_code": record.unit_code, + "evidence_reference": record.evidence_reference, + "truth_status": record.truth_status, + "recorded_at": _sql_timestamp(record.recorded_at), + }, + ) + except IntegrityError as exc: + _raise_unique_violation(exc, "data-quality rule could not be stored") + + def insert_data_quality_observation(self, record: DataQualityObservationRecord) -> None: + """Insert one immutable observation after checking its tenant-owned rule.""" + + try: + with self._engine.begin() as conn: + parent = conn.execute( + text( + """ + SELECT data_quality_rule_id + FROM data_quality_rules + WHERE data_quality_rule_id = :data_quality_rule_id + AND tenant_reference = :tenant_reference + """ + ), + { + "data_quality_rule_id": record.data_quality_rule_id, + "tenant_reference": record.tenant_reference, + }, + ).first() + if parent is None: + raise KeyError("data-quality rule not found in this tenant") + conn.execute( + text( + """ + INSERT INTO data_quality_observations ( + data_quality_observation_id, data_quality_rule_id, + critical_data_element_id, catalog_object_id, + tenant_reference, source_observation_id, observed_value, + observed_at, quality_status, evidence_reference, + truth_status, recorded_at + ) VALUES ( + :data_quality_observation_id, :data_quality_rule_id, + :critical_data_element_id, :catalog_object_id, + :tenant_reference, :source_observation_id, :observed_value, + :observed_at, :quality_status, :evidence_reference, + :truth_status, :recorded_at + ) + """ + ), + { + "data_quality_observation_id": record.data_quality_observation_id, + "data_quality_rule_id": record.data_quality_rule_id, + "critical_data_element_id": record.critical_data_element_id, + "catalog_object_id": record.catalog_object_id, + "tenant_reference": record.tenant_reference, + "source_observation_id": record.source_observation_id, + "observed_value": str(record.observed_value), + "observed_at": _sql_timestamp(record.observed_at), + "quality_status": record.quality_status, + "evidence_reference": record.evidence_reference, + "truth_status": record.truth_status, + "recorded_at": _sql_timestamp(record.recorded_at), + }, + ) + except IntegrityError as exc: + _raise_unique_violation(exc, "data-quality observation could not be stored") + + def get_critical_data_element( + self, + *, + tenant_reference: str, + critical_data_element_id: str, + ) -> CriticalDataElementRecord | None: + """Load one tenant-owned CDE.""" + + with self._engine.begin() as conn: + row = conn.execute( + text( + """ + SELECT * FROM critical_data_elements + WHERE critical_data_element_id = :critical_data_element_id + AND tenant_reference = :tenant_reference + """ + ), + { + "critical_data_element_id": critical_data_element_id, + "tenant_reference": tenant_reference, + }, + ).mappings().first() + if row is None: + return None + return self._critical_data_element(row) + + def get_data_quality_rule( + self, + *, + tenant_reference: str, + data_quality_rule_id: str, + ) -> DataQualityRuleRecord | None: + """Load one tenant-owned quality rule.""" + + with self._engine.begin() as conn: + row = conn.execute( + text( + """ + SELECT * FROM data_quality_rules + WHERE data_quality_rule_id = :data_quality_rule_id + AND tenant_reference = :tenant_reference + """ + ), + { + "data_quality_rule_id": data_quality_rule_id, + "tenant_reference": tenant_reference, + }, + ).mappings().first() + if row is None: + return None + return self._data_quality_rule(row) + + def profile_rows( + self, + *, + tenant_reference: str, + catalog_object_id: str, + ) -> ProfileRows: + """Hydrate all 0003 rows for one tenant-owned catalog object.""" + + with self._engine.begin() as conn: + owners = [ + self._owner_assignment(row) + for row in conn.execute( + text( + """ + SELECT * FROM data_owner_assignments + WHERE tenant_reference = :tenant_reference + AND catalog_object_id = :catalog_object_id + ORDER BY valid_from, data_owner_assignment_id + """ + ), + { + "tenant_reference": tenant_reference, + "catalog_object_id": catalog_object_id, + }, + ).mappings() + ] + elements = [ + self._critical_data_element(row) + for row in conn.execute( + text( + """ + SELECT * FROM critical_data_elements + WHERE tenant_reference = :tenant_reference + AND catalog_object_id = :catalog_object_id + ORDER BY element_key, critical_data_element_id + """ + ), + { + "tenant_reference": tenant_reference, + "catalog_object_id": catalog_object_id, + }, + ).mappings() + ] + rules = [ + self._data_quality_rule(row) + for row in conn.execute( + text( + """ + SELECT * FROM data_quality_rules + WHERE tenant_reference = :tenant_reference + AND catalog_object_id = :catalog_object_id + ORDER BY rule_code, data_quality_rule_id + """ + ), + { + "tenant_reference": tenant_reference, + "catalog_object_id": catalog_object_id, + }, + ).mappings() + ] + observations = [ + self._data_quality_observation(row) + for row in conn.execute( + text( + """ + SELECT * FROM data_quality_observations + WHERE tenant_reference = :tenant_reference + AND catalog_object_id = :catalog_object_id + ORDER BY observed_at, source_observation_id + """ + ), + { + "tenant_reference": tenant_reference, + "catalog_object_id": catalog_object_id, + }, + ).mappings() + ] + return owners, elements, rules, observations + + @staticmethod + def _owner_assignment(row: Any) -> DataOwnerAssignmentRecord: + """Map one relational owner row into its contract.""" + + return DataOwnerAssignmentRecord( + data_owner_assignment_id=row["data_owner_assignment_id"], + catalog_object_id=row["catalog_object_id"], + tenant_reference=row["tenant_reference"], + owner_subject=row["owner_subject"], + owner_display_name=row["owner_display_name"], + valid_from=_as_datetime(row["valid_from"]), + valid_to=_as_datetime(row["valid_to"]) if row["valid_to"] is not None else None, + evidence_reference=row["evidence_reference"], + truth_status=row["truth_status"], + recorded_at=_as_datetime(row["recorded_at"]), + ) + + @staticmethod + def _critical_data_element(row: Any) -> CriticalDataElementRecord: + """Map one relational CDE row into its contract.""" + + return CriticalDataElementRecord( + critical_data_element_id=row["critical_data_element_id"], + catalog_object_id=row["catalog_object_id"], + tenant_reference=row["tenant_reference"], + element_key=row["element_key"], + display_name=row["display_name"], + definition_text=row["definition_text"], + data_classification=row["data_classification"], + evidence_reference=row["evidence_reference"], + truth_status=row["truth_status"], + recorded_at=_as_datetime(row["recorded_at"]), + ) + + @staticmethod + def _data_quality_rule(row: Any) -> DataQualityRuleRecord: + """Map one relational quality-rule row into its contract.""" + + return DataQualityRuleRecord( + data_quality_rule_id=row["data_quality_rule_id"], + critical_data_element_id=row["critical_data_element_id"], + catalog_object_id=row["catalog_object_id"], + tenant_reference=row["tenant_reference"], + rule_code=row["rule_code"], + rule_description=row["rule_description"], + metric_code=row["metric_code"], + threshold_operator=row["threshold_operator"], + threshold_value=_decimal(row["threshold_value"]), + unit_code=row["unit_code"], + evidence_reference=row["evidence_reference"], + truth_status=row["truth_status"], + recorded_at=_as_datetime(row["recorded_at"]), + ) + + @staticmethod + def _data_quality_observation(row: Any) -> DataQualityObservationRecord: + """Map one relational observation row into its contract.""" + + return DataQualityObservationRecord( + data_quality_observation_id=row["data_quality_observation_id"], + data_quality_rule_id=row["data_quality_rule_id"], + critical_data_element_id=row["critical_data_element_id"], + catalog_object_id=row["catalog_object_id"], + tenant_reference=row["tenant_reference"], + source_observation_id=row["source_observation_id"], + observed_value=_decimal(row["observed_value"]), + observed_at=_as_datetime(row["observed_at"]), + quality_status=row["quality_status"], + evidence_reference=row["evidence_reference"], + truth_status=row["truth_status"], + recorded_at=_as_datetime(row["recorded_at"]), + ) + + +_ACTIVE_STORE: DataManagementStore | None = None + + +def reset_memory_data_management() -> None: + """Clear every process-local evidence row.""" + + with _MEMORY_LOCK: + _MEMORY_OWNER_ROWS.clear() + _MEMORY_ELEMENT_ROWS.clear() + _MEMORY_RULE_ROWS.clear() + _MEMORY_OBSERVATION_ROWS.clear() + + +def snapshot_memory_data_management() -> dict[str, dict[str, Any]]: + """Copy process-local evidence rows for test isolation.""" + + with _MEMORY_LOCK: + return { + "owners": deepcopy(_MEMORY_OWNER_ROWS), + "elements": deepcopy(_MEMORY_ELEMENT_ROWS), + "rules": deepcopy(_MEMORY_RULE_ROWS), + "observations": deepcopy(_MEMORY_OBSERVATION_ROWS), + } + + +def restore_memory_data_management(snapshot: dict[str, dict[str, Any]]) -> None: + """Replace process-local rows from a prior snapshot.""" + + with _MEMORY_LOCK: + _MEMORY_OWNER_ROWS.clear() + _MEMORY_OWNER_ROWS.update(deepcopy(snapshot.get("owners", {}))) + _MEMORY_ELEMENT_ROWS.clear() + _MEMORY_ELEMENT_ROWS.update(deepcopy(snapshot.get("elements", {}))) + _MEMORY_RULE_ROWS.clear() + _MEMORY_RULE_ROWS.update(deepcopy(snapshot.get("rules", {}))) + _MEMORY_OBSERVATION_ROWS.clear() + _MEMORY_OBSERVATION_ROWS.update(deepcopy(snapshot.get("observations", {}))) + + +def build_data_management_store() -> DataManagementStore: + """Return the in-memory store or migration-0003 relational store.""" + + dsn = load_bootstrap().database_dsn + if not dsn: + return InMemoryDataManagementStore() + try: + return RelationalDataManagementStore(dsn) + except RuntimeError as exc: + raise RuntimeError( + "SDP_DATABASE_DSN is set but the data-management store could not open. " + f"{_GRAPH_EXTRA_HINT}" + ) from exc + + +def get_data_management_store() -> DataManagementStore: + """Return the process-wide evidence store, building it once.""" + + global _ACTIVE_STORE + if _ACTIVE_STORE is None: + _ACTIVE_STORE = build_data_management_store() + return _ACTIVE_STORE + + +def set_data_management_store(store: DataManagementStore | None) -> None: + """Replace or clear the process-wide evidence store for tests.""" + + global _ACTIVE_STORE + _ACTIVE_STORE = store From 474c7f0737fa958a249a4083bf7f6cff35cb2e26 Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Tue, 18 Aug 2026 03:37:42 -0700 Subject: [PATCH 06/13] feat: add data management evidence service --- src/sdp/data_management_evidence.py | 307 ++++++++++++++++++++++++++++ 1 file changed, 307 insertions(+) create mode 100644 src/sdp/data_management_evidence.py diff --git a/src/sdp/data_management_evidence.py b/src/sdp/data_management_evidence.py new file mode 100644 index 0000000..2d60d4a --- /dev/null +++ b/src/sdp/data_management_evidence.py @@ -0,0 +1,307 @@ +"""Actionable data-management evidence services above the catalog plane.""" + +from __future__ import annotations + +from datetime import datetime, timezone +from uuid import uuid4 + +from sdp_core.catalog_plane import CatalogObjectRecord, PlaneActor +from sdp_core.data_management_evidence import ( + CriticalDataElementDraft, + CriticalDataElementRecord, + DataManagementMutationEnvelope, + DataManagementProfile, + DataOwnerAssignmentDraft, + DataOwnerAssignmentRecord, + DataQualityObservationDraft, + DataQualityObservationRecord, + DataQualityRuleDraft, + DataQualityRuleRecord, +) + +from .catalog_plane_store import get_catalog_plane_store +from .data_management_store import get_data_management_store +from .policy import evaluate + + +def _now() -> datetime: + """Return an aware UTC timestamp.""" + + return datetime.now(timezone.utc) + + +def _record_id() -> str: + """Allocate an opaque record identifier.""" + + return str(uuid4()) + + +def _govern(actor: PlaneActor, *, mutate: bool) -> str: + """Authorize the evidence-plane operation and record policy evidence.""" + + if mutate and not actor.can_mutate(): + raise PermissionError("data-management writes require an admin Keyverse role") + if not mutate and not actor.can_read(): + raise PermissionError("data-management reads require a Keyverse reader role") + decision = evaluate( + subject=actor.subject, + resource="plane" if mutate else "catalog", + action="create" if mutate else "search", + purpose=actor.access_purpose, + ) + if decision.effect != "allow": + raise PermissionError(decision.reason) + return decision.decision_id + + +def _catalog_object(actor: PlaneActor, catalog_object_id: str) -> CatalogObjectRecord: + """Load one tenant-owned catalog object or fail without leaking existence.""" + + record = get_catalog_plane_store().get_catalog_object( + tenant_reference=actor.tenant_reference, + catalog_object_id=catalog_object_id, + ) + if record is None: + raise KeyError("catalog object not found in this tenant") + return record + + +def _catalog_dataset(actor: PlaneActor, catalog_object_id: str) -> CatalogObjectRecord: + """Require a tenant-owned catalog dataset for CDE evidence.""" + + record = _catalog_object(actor, catalog_object_id) + if record.object_kind != "catalog_dataset": + raise ValueError("critical data elements require a catalog_dataset parent") + return record + + +def assign_data_owner( + actor: PlaneActor, + catalog_object_id: str, + draft: DataOwnerAssignmentDraft, +) -> DataManagementMutationEnvelope: + """Attach an effective-dated owner decision to one catalog object.""" + + decision_id = _govern(actor, mutate=True) + _catalog_object(actor, catalog_object_id) + record = DataOwnerAssignmentRecord( + data_owner_assignment_id=_record_id(), + catalog_object_id=catalog_object_id, + tenant_reference=actor.tenant_reference, + recorded_at=_now(), + **draft.model_dump(), + ) + get_data_management_store().insert_owner_assignment(record) + return DataManagementMutationEnvelope( + status="data_owner_assigned", + tenant_reference=actor.tenant_reference, + catalog_object_id=catalog_object_id, + policy_decision_id=decision_id, + customer_next_action=( + f"다음으로 POST /plane/catalog-objects/{catalog_object_id}/critical-data-elements " + "에서 업무 의사결정에 중요한 Critical Data Element를 등록하세요." + ), + data_owner_assignment=record, + ) + + +def register_critical_data_element( + actor: PlaneActor, + catalog_object_id: str, + draft: CriticalDataElementDraft, +) -> DataManagementMutationEnvelope: + """Register one evidence-backed CDE under a catalog dataset.""" + + decision_id = _govern(actor, mutate=True) + _catalog_dataset(actor, catalog_object_id) + record = CriticalDataElementRecord( + critical_data_element_id=_record_id(), + catalog_object_id=catalog_object_id, + tenant_reference=actor.tenant_reference, + recorded_at=_now(), + **draft.model_dump(), + ) + get_data_management_store().insert_critical_data_element(record) + return DataManagementMutationEnvelope( + status="critical_data_element_registered", + tenant_reference=actor.tenant_reference, + catalog_object_id=catalog_object_id, + policy_decision_id=decision_id, + customer_next_action=( + f"다음으로 POST /plane/critical-data-elements/{record.critical_data_element_id}/" + "quality-rules 에 측정 가능한 Data Quality rule과 기준값을 등록하세요." + ), + critical_data_element=record, + ) + + +def define_data_quality_rule( + actor: PlaneActor, + critical_data_element_id: str, + draft: DataQualityRuleDraft, +) -> DataManagementMutationEnvelope: + """Define one evidence-backed rule for a tenant-owned CDE.""" + + decision_id = _govern(actor, mutate=True) + store = get_data_management_store() + element = store.get_critical_data_element( + tenant_reference=actor.tenant_reference, + critical_data_element_id=critical_data_element_id, + ) + if element is None: + raise KeyError("critical data element not found in this tenant") + record = DataQualityRuleRecord( + data_quality_rule_id=_record_id(), + critical_data_element_id=critical_data_element_id, + catalog_object_id=element.catalog_object_id, + tenant_reference=actor.tenant_reference, + recorded_at=_now(), + **draft.model_dump(), + ) + store.insert_data_quality_rule(record) + return DataManagementMutationEnvelope( + status="data_quality_rule_defined", + tenant_reference=actor.tenant_reference, + catalog_object_id=element.catalog_object_id, + policy_decision_id=decision_id, + customer_next_action=( + f"다음으로 POST /plane/quality-rules/{record.data_quality_rule_id}/observations " + "에서 실제 control run 또는 measurement evidence를 기록하세요." + ), + data_quality_rule=record, + ) + + +def record_data_quality_observation( + actor: PlaneActor, + data_quality_rule_id: str, + draft: DataQualityObservationDraft, +) -> DataManagementMutationEnvelope: + """Append one immutable observation for a tenant-owned quality rule.""" + + decision_id = _govern(actor, mutate=True) + store = get_data_management_store() + rule = store.get_data_quality_rule( + tenant_reference=actor.tenant_reference, + data_quality_rule_id=data_quality_rule_id, + ) + if rule is None: + raise KeyError("data-quality rule not found in this tenant") + record = DataQualityObservationRecord( + data_quality_observation_id=_record_id(), + data_quality_rule_id=data_quality_rule_id, + critical_data_element_id=rule.critical_data_element_id, + catalog_object_id=rule.catalog_object_id, + tenant_reference=actor.tenant_reference, + recorded_at=_now(), + **draft.model_dump(), + ) + store.insert_data_quality_observation(record) + return DataManagementMutationEnvelope( + status="data_quality_observation_recorded", + tenant_reference=actor.tenant_reference, + catalog_object_id=rule.catalog_object_id, + policy_decision_id=decision_id, + customer_next_action=( + f"GET /plane/catalog-objects/{rule.catalog_object_id}/data-management-profile " + "에서 owner, CDE, rule, observation evidence가 완결되었는지 확인하세요." + ), + data_quality_observation=record, + ) + + +def _profile_next_action( + catalog_object_id: str, + *, + owner_present: bool, + element_present: bool, + rule_present: bool, + observation_present: bool, + elements: list[CriticalDataElementRecord], + rules: list[DataQualityRuleRecord], +) -> str: + """Return the first precise buyer action needed to complete the profile.""" + + if not owner_present: + return ( + f"POST /plane/catalog-objects/{catalog_object_id}/data-owner-assignments 로 " + "책임 있는 Data Owner와 승인 evidence를 먼저 지정하세요." + ) + if not element_present: + return ( + f"POST /plane/catalog-objects/{catalog_object_id}/critical-data-elements 로 " + "업무 의사결정에 중요한 Critical Data Element를 등록하세요." + ) + if not rule_present: + return ( + f"POST /plane/critical-data-elements/{elements[0].critical_data_element_id}/" + "quality-rules 로 첫 Data Quality rule을 정의하세요." + ) + if not observation_present: + return ( + f"POST /plane/quality-rules/{rules[0].data_quality_rule_id}/observations 로 " + "실제 measurement 또는 control-run evidence를 기록하세요." + ) + return ( + "Evidence profile이 완결되었습니다. 최신 observation의 provenance와 status를 " + "검토하고, 새로운 control run이 발생하면 append-only evidence를 추가하세요." + ) + + +def build_data_management_profile( + actor: PlaneActor, + catalog_object_id: str, +) -> DataManagementProfile: + """Build an explainable evidence-completeness profile for one dataset.""" + + decision_id = _govern(actor, mutate=False) + _catalog_object(actor, catalog_object_id) + owners, elements, rules, observations = get_data_management_store().profile_rows( + tenant_reference=actor.tenant_reference, + catalog_object_id=catalog_object_id, + ) + current_time = _now() + owner_present = any( + row.truth_status == "authoritative" + and row.valid_from <= current_time + and (row.valid_to is None or current_time < row.valid_to) + for row in owners + ) + element_present = any(row.truth_status == "authoritative" for row in elements) + rule_present = any(row.truth_status == "authoritative" for row in rules) + observation_present = any( + row.truth_status in {"authoritative", "observed"} for row in observations + ) + factors = { + "data_owner_present": owner_present, + "critical_data_element_present": element_present, + "quality_rule_present": rule_present, + "quality_observation_present": observation_present, + } + counts = { + "data_owner_assignments": len(owners), + "critical_data_elements": len(elements), + "data_quality_rules": len(rules), + "data_quality_observations": len(observations), + } + return DataManagementProfile( + tenant_reference=actor.tenant_reference, + catalog_object_id=catalog_object_id, + policy_decision_id=decision_id, + evidence_complete=all(factors.values()), + factors=factors, + counts=counts, + customer_next_action=_profile_next_action( + catalog_object_id, + owner_present=owner_present, + element_present=element_present, + rule_present=rule_present, + observation_present=observation_present, + elements=elements, + rules=rules, + ), + data_owner_assignments=owners, + critical_data_elements=elements, + data_quality_rules=rules, + data_quality_observations=observations, + ) From a9b352c89544a35e2329c45784035089436229bc Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Tue, 18 Aug 2026 03:41:35 -0700 Subject: [PATCH 07/13] ci: apply and verify data management route patch --- .../apply-data-management-routes.yml | 190 ++++++++++++++++++ 1 file changed, 190 insertions(+) create mode 100644 .github/workflows/apply-data-management-routes.yml diff --git a/.github/workflows/apply-data-management-routes.yml b/.github/workflows/apply-data-management-routes.yml new file mode 100644 index 0000000..e6c7917 --- /dev/null +++ b/.github/workflows/apply-data-management-routes.yml @@ -0,0 +1,190 @@ +name: Apply data-management routes + +on: + push: + branches: [feat/data-management-evidence-profile] + +permissions: + contents: write + +concurrency: + group: apply-data-management-routes + cancel-in-progress: false + +jobs: + apply-and-verify: + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 + with: + ref: feat/data-management-evidence-profile + fetch-depth: 0 + - uses: actions/setup-python@ece7cb06caefa5fff74198d8649806c4678c61a1 # v6 + with: + python-version: '3.12' + - name: Apply bounded API registration patch + run: | + python - <<'PY' + from pathlib import Path + + api_path = Path("src/sdp/api.py") + source = api_path.read_text(encoding="utf-8") + endpoint_marker = '"/plane/catalog-objects/{catalog_object_id}/data-management-profile"' + + if endpoint_marker not in source: + core_import = '''from sdp_core.data_management_evidence import ( + CriticalDataElementDraft, + DataOwnerAssignmentDraft, + DataQualityObservationDraft, + DataQualityRuleDraft, + ) + ''' + import_anchor = "from sdp_core.catalog_plane import (\n" + if import_anchor not in source: + raise SystemExit("catalog-plane import anchor not found") + source = source.replace(import_anchor, core_import + import_anchor, 1) + + service_import = '''from .data_management_evidence import ( + assign_data_owner, + build_data_management_profile, + define_data_quality_rule, + record_data_quality_observation, + register_critical_data_element, + ) + ''' + service_anchor = "from .config import get_app_config\n" + if service_anchor not in source: + raise SystemExit("config import anchor not found") + source = source.replace(service_anchor, service_import + service_anchor, 1) + + route_anchor = '@app.get("/ontology/term/{term}/graph")\n' + if route_anchor not in source: + raise SystemExit("ontology route anchor not found") + routes = '''@app.post("/plane/catalog-objects/{catalog_object_id}/data-owner-assignments") + def plane_assign_data_owner( + catalog_object_id: str, + payload: DataOwnerAssignmentDraft, + request: Request, + ) -> dict[str, Any]: + """Assign an evidence-backed data owner to one catalog object.""" + + actor = _plane_actor(request) + try: + return assign_data_owner(actor, catalog_object_id, payload).model_dump() + except PermissionError as exc: + raise HTTPException(status_code=403, detail=str(exc)) from exc + except KeyError as exc: + raise HTTPException(status_code=404, detail=str(exc)) from exc + except ValueError as exc: + raise HTTPException(status_code=400, detail=str(exc)) from exc + + + @app.post("/plane/catalog-objects/{catalog_object_id}/critical-data-elements") + def plane_register_critical_data_element( + catalog_object_id: str, + payload: CriticalDataElementDraft, + request: Request, + ) -> dict[str, Any]: + """Register an evidence-backed CDE under a catalog dataset.""" + + actor = _plane_actor(request) + try: + return register_critical_data_element( + actor, + catalog_object_id, + payload, + ).model_dump() + except PermissionError as exc: + raise HTTPException(status_code=403, detail=str(exc)) from exc + except KeyError as exc: + raise HTTPException(status_code=404, detail=str(exc)) from exc + except ValueError as exc: + raise HTTPException(status_code=400, detail=str(exc)) from exc + + + @app.post("/plane/critical-data-elements/{critical_data_element_id}/quality-rules") + def plane_define_data_quality_rule( + critical_data_element_id: str, + payload: DataQualityRuleDraft, + request: Request, + ) -> dict[str, Any]: + """Define one evidence-backed quality rule for a CDE.""" + + actor = _plane_actor(request) + try: + return define_data_quality_rule( + actor, + critical_data_element_id, + payload, + ).model_dump() + except PermissionError as exc: + raise HTTPException(status_code=403, detail=str(exc)) from exc + except KeyError as exc: + raise HTTPException(status_code=404, detail=str(exc)) from exc + except ValueError as exc: + raise HTTPException(status_code=400, detail=str(exc)) from exc + + + @app.post("/plane/quality-rules/{data_quality_rule_id}/observations") + def plane_record_data_quality_observation( + data_quality_rule_id: str, + payload: DataQualityObservationDraft, + request: Request, + ) -> dict[str, Any]: + """Append one immutable evidence-backed quality observation.""" + + actor = _plane_actor(request) + try: + return record_data_quality_observation( + actor, + data_quality_rule_id, + payload, + ).model_dump() + except PermissionError as exc: + raise HTTPException(status_code=403, detail=str(exc)) from exc + except KeyError as exc: + raise HTTPException(status_code=404, detail=str(exc)) from exc + except ValueError as exc: + raise HTTPException(status_code=400, detail=str(exc)) from exc + + + @app.get("/plane/catalog-objects/{catalog_object_id}/data-management-profile") + def plane_data_management_profile( + catalog_object_id: str, + request: Request, + ) -> dict[str, Any]: + """Return the explainable evidence-completeness profile.""" + + actor = _plane_actor(request) + try: + return build_data_management_profile(actor, catalog_object_id).model_dump() + except PermissionError as exc: + raise HTTPException(status_code=403, detail=str(exc)) from exc + except KeyError as exc: + raise HTTPException(status_code=404, detail=str(exc)) from exc + except ValueError as exc: + raise HTTPException(status_code=400, detail=str(exc)) from exc + + + ''' + source = source.replace(route_anchor, routes + route_anchor, 1) + api_path.write_text(source, encoding="utf-8") + + workflow_path = Path(".github/workflows/apply-data-management-routes.yml") + workflow_path.unlink() + PY + - name: Install hash-pinned test dependencies + run: python -m pip install --disable-pip-version-check --no-cache-dir --require-hashes -r requirements-test.txt + - name: Verify focused data-management acceptance + env: + PYTHONPATH: src + run: | + python -m pytest tests/test_data_management_evidence.py tests/test_data_management_migration.py + python -m compileall -q src + - name: Commit verified patch and remove one-shot workflow + run: | + git config user.name "cwl-data-management-agent" + git config user.email "cwl-data-management-agent@users.noreply.github.com" + git add src/sdp/api.py .github/workflows/apply-data-management-routes.yml + git commit -m "feat: expose data management evidence routes" + git push origin HEAD:feat/data-management-evidence-profile From eb6f53859c56885e8ed3a73bb5c7ba7eaca1ad13 Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Tue, 18 Aug 2026 03:43:41 -0700 Subject: [PATCH 08/13] test: isolate data management evidence state --- tests/conftest.py | 18 +++++++++++++++++- 1 file changed, 17 insertions(+), 1 deletion(-) diff --git a/tests/conftest.py b/tests/conftest.py index 47e9906..04b5231 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -1,12 +1,28 @@ -"""Shared pytest configuration for explicit demo-only identity behavior.""" +"""Shared pytest configuration for identity and evidence-store isolation.""" from __future__ import annotations import pytest +from sdp.data_management_store import ( + restore_memory_data_management, + snapshot_memory_data_management, +) + @pytest.fixture(autouse=True) def _allow_demo_subject_header(monkeypatch: pytest.MonkeyPatch) -> None: """Opt tests into the demo subject-header path unless a test removes it.""" monkeypatch.setenv("SDP_ALLOW_UNVERIFIED_SUBJECT_HEADER", "true") + + +@pytest.fixture(autouse=True) +def _isolate_data_management_evidence() -> None: + """Restore process-local evidence rows after every test.""" + + snapshot = snapshot_memory_data_management() + try: + yield + finally: + restore_memory_data_management(snapshot) From ca0313064eeeef5e04b6435eeb45394aefc9d92f Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Tue, 18 Aug 2026 03:44:44 -0700 Subject: [PATCH 09/13] test: verify relational data management evidence parity --- tests/test_data_management_sql.py | 225 ++++++++++++++++++++++++++++++ 1 file changed, 225 insertions(+) create mode 100644 tests/test_data_management_sql.py diff --git a/tests/test_data_management_sql.py b/tests/test_data_management_sql.py new file mode 100644 index 0000000..ad05899 --- /dev/null +++ b/tests/test_data_management_sql.py @@ -0,0 +1,225 @@ +"""SQLite mapping tests for migration-0003 data-management evidence.""" + +from __future__ import annotations + +from datetime import datetime, timezone + +import pytest +from sqlalchemy import create_engine, text + +from sdp.catalog_plane import create_catalog_object +from sdp.catalog_plane_store import ( + RelationalCatalogPlaneStore, + apply_catalog_plane_sqlite_schema, + set_catalog_plane_store, +) +from sdp.data_management_evidence import ( + assign_data_owner, + build_data_management_profile, + define_data_quality_rule, + record_data_quality_observation, + register_critical_data_element, +) +from sdp.data_management_store import ( + RelationalDataManagementStore, + apply_data_management_sqlite_schema, + set_data_management_store, +) +from sdp_core.catalog_plane import CatalogObjectCreateRequest, PlaneActor +from sdp_core.data_management_evidence import ( + CriticalDataElementDraft, + DataOwnerAssignmentDraft, + DataQualityObservationDraft, + DataQualityRuleDraft, +) + + +def _admin_actor(tenant_reference: str = "demo") -> PlaneActor: + """Return a purpose-bound Keyverse admin actor.""" + + return PlaneActor( + subject="admin", + tenant_reference=tenant_reference, + roles=["admin"], + access_purpose="glossary_stewardship", + binding_source="test", + ) + + +def _reader_actor(tenant_reference: str = "demo") -> PlaneActor: + """Return a purpose-bound catalog reader.""" + + return PlaneActor( + subject="analyst", + tenant_reference=tenant_reference, + roles=["data-analyst"], + access_purpose="catalog_browse", + binding_source="test", + ) + + +def _sqlite_dsn(tmp_path) -> str: + """Create a SQLite file with the 0002 parent and 0003 evidence schemas.""" + + database_path = tmp_path / "data-management.sqlite" + dsn = f"sqlite:///{database_path}" + engine = create_engine(dsn, future=True) + apply_catalog_plane_sqlite_schema(engine) + apply_data_management_sqlite_schema(engine) + engine.dispose() + return dsn + + +def _create_dataset(actor: PlaneActor) -> str: + """Create one relational catalog dataset and return its identifier.""" + + envelope = create_catalog_object( + actor, + CatalogObjectCreateRequest( + object_kind="catalog_dataset", + object_slug="billing-settlement-evidence", + display_title="Billing settlement evidence", + definition_text="Usage, invoice, payment, refund, and settlement evidence.", + preferred_language="en", + steward_display_name="Billing Data Steward", + aliases=[], + document_kg_links=[], + concept_bindings=[], + score_references=[], + ), + ) + return envelope.catalog_object.catalog_object_id + + +def _populate_profile(actor: PlaneActor, catalog_object_id: str) -> None: + """Populate owner, CDE, quality rule, and immutable observation rows.""" + + assign_data_owner( + actor, + catalog_object_id, + DataOwnerAssignmentDraft( + owner_subject="billing-operations-owner", + owner_display_name="Billing Operations Owner", + valid_from=datetime(2026, 8, 18, tzinfo=timezone.utc), + evidence_reference="https://evidence.example.test/decisions/billing-owner", + truth_status="authoritative", + ), + ) + cde = register_critical_data_element( + actor, + catalog_object_id, + CriticalDataElementDraft( + element_key="settlement_amount", + display_name="Settlement amount", + definition_text="Cash amount paid out by the provider for one settlement.", + data_classification="restricted_financial", + evidence_reference="https://evidence.example.test/dictionary/settlement-amount", + truth_status="authoritative", + ), + ).critical_data_element + assert cde is not None + rule = define_data_quality_rule( + actor, + cde.critical_data_element_id, + DataQualityRuleDraft( + rule_code="settlement_matches_expected_amount", + rule_description="Settlement plus provider fee equals the captured invoice amount.", + metric_code="reconciliation_difference", + threshold_operator="equal_to", + threshold_value="0", + unit_code="KRW", + evidence_reference="https://evidence.example.test/controls/reconciliation", + truth_status="authoritative", + ), + ).data_quality_rule + assert rule is not None + record_data_quality_observation( + actor, + rule.data_quality_rule_id, + DataQualityObservationDraft( + source_observation_id="reconciliation_run_2026_08_18", + observed_value="0", + observed_at=datetime(2026, 8, 18, 1, tzinfo=timezone.utc), + quality_status="passed", + evidence_reference="https://evidence.example.test/runs/reconciliation-2026-08-18", + truth_status="observed", + ), + ) + + +def test_relational_profile_survives_store_reopen(tmp_path) -> None: + """Migration-0003 rows remain complete after both stores reopen the DSN.""" + + dsn = _sqlite_dsn(tmp_path) + set_catalog_plane_store(RelationalCatalogPlaneStore(dsn)) + set_data_management_store(RelationalDataManagementStore(dsn)) + try: + catalog_object_id = _create_dataset(_admin_actor()) + _populate_profile(_admin_actor(), catalog_object_id) + + set_catalog_plane_store(RelationalCatalogPlaneStore(dsn)) + set_data_management_store(RelationalDataManagementStore(dsn)) + profile = build_data_management_profile(_reader_actor(), catalog_object_id) + + assert profile.evidence_complete is True + assert profile.counts == { + "data_owner_assignments": 1, + "critical_data_elements": 1, + "data_quality_rules": 1, + "data_quality_observations": 1, + } + assert profile.data_owner_assignments[0].owner_display_name == "Billing Operations Owner" + assert profile.data_quality_observations[0].source_observation_id == "reconciliation_run_2026_08_18" + finally: + set_data_management_store(None) + set_catalog_plane_store(None) + + +def test_relational_store_enforces_natural_keys_and_tenant_scope(tmp_path) -> None: + """SQL uniqueness and tenant filters match the in-memory store.""" + + dsn = _sqlite_dsn(tmp_path) + set_catalog_plane_store(RelationalCatalogPlaneStore(dsn)) + set_data_management_store(RelationalDataManagementStore(dsn)) + try: + actor = _admin_actor() + catalog_object_id = _create_dataset(actor) + draft = DataOwnerAssignmentDraft( + owner_subject="billing-operations-owner", + owner_display_name="Billing Operations Owner", + valid_from=datetime(2026, 8, 18, tzinfo=timezone.utc), + evidence_reference="https://evidence.example.test/decisions/billing-owner", + truth_status="authoritative", + ) + assign_data_owner(actor, catalog_object_id, draft) + with pytest.raises(ValueError, match="duplicate data-owner assignment"): + assign_data_owner(actor, catalog_object_id, draft) + + profile = build_data_management_profile(_reader_actor("external"), catalog_object_id) + assert profile.catalog_object_id == catalog_object_id # pragma: no cover + except KeyError as exc: + assert "catalog object not found" in str(exc) + finally: + set_data_management_store(None) + set_catalog_plane_store(None) + + +def test_sqlite_schema_creates_all_0003_tables(tmp_path) -> None: + """The portable rewrite keeps the production migration table names.""" + + dsn = _sqlite_dsn(tmp_path) + engine = create_engine(dsn, future=True) + with engine.connect() as connection: + tables = { + row[0] + for row in connection.execute( + text("SELECT name FROM sqlite_master WHERE type = 'table'") + ) + } + engine.dispose() + assert { + "data_owner_assignments", + "critical_data_elements", + "data_quality_rules", + "data_quality_observations", + }.issubset(tables) From 56b0ad2fd3c74a140aea0e5ece3b6c8738a674fd Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Tue, 18 Aug 2026 03:45:50 -0700 Subject: [PATCH 10/13] test: require explicit cross-tenant denial --- tests/test_data_management_sql.py | 6 ++---- 1 file changed, 2 insertions(+), 4 deletions(-) diff --git a/tests/test_data_management_sql.py b/tests/test_data_management_sql.py index ad05899..ac7e820 100644 --- a/tests/test_data_management_sql.py +++ b/tests/test_data_management_sql.py @@ -195,10 +195,8 @@ def test_relational_store_enforces_natural_keys_and_tenant_scope(tmp_path) -> No with pytest.raises(ValueError, match="duplicate data-owner assignment"): assign_data_owner(actor, catalog_object_id, draft) - profile = build_data_management_profile(_reader_actor("external"), catalog_object_id) - assert profile.catalog_object_id == catalog_object_id # pragma: no cover - except KeyError as exc: - assert "catalog object not found" in str(exc) + with pytest.raises(KeyError, match="catalog object not found"): + build_data_management_profile(_reader_actor("external"), catalog_object_id) finally: set_data_management_store(None) set_catalog_plane_store(None) From 2e97e205d690619eb01a5e268c89b8e3dbcd85f9 Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Wed, 26 Aug 2026 01:30:50 +0900 Subject: [PATCH 11/13] feat: expose data management evidence routes (apply one-shot patch inline) Applies the route registration that the one-shot apply-data-management-routes workflow was meant to land, then removes that self-modifying workflow: - /plane/catalog-objects/{id}/data-owner-assignments (POST) - /plane/catalog-objects/{id}/critical-data-elements (POST) - /plane/critical-data-elements/{id}/quality-rules (POST) - /plane/quality-rules/{id}/observations (POST) - /plane/catalog-objects/{id}/data-management-profile (GET) Fixes surfaced by its verify step: - conftest now snapshot/restores the catalog-plane store so tenant-isolation tests cannot collide on object_slug with earlier tests. - migration contract test asserts whitespace-collapsed SQL and the actual column-level FK form instead of alignment-fragile literals. Full suite: 286 passed, 9 skipped. --- .../apply-data-management-routes.yml | 190 ------------------ src/sdp/api.py | 107 ++++++++++ tests/conftest.py | 12 +- tests/test_data_management_migration.py | 8 +- 4 files changed, 122 insertions(+), 195 deletions(-) delete mode 100644 .github/workflows/apply-data-management-routes.yml diff --git a/.github/workflows/apply-data-management-routes.yml b/.github/workflows/apply-data-management-routes.yml deleted file mode 100644 index e6c7917..0000000 --- a/.github/workflows/apply-data-management-routes.yml +++ /dev/null @@ -1,190 +0,0 @@ -name: Apply data-management routes - -on: - push: - branches: [feat/data-management-evidence-profile] - -permissions: - contents: write - -concurrency: - group: apply-data-management-routes - cancel-in-progress: false - -jobs: - apply-and-verify: - runs-on: ubuntu-latest - steps: - - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 - with: - ref: feat/data-management-evidence-profile - fetch-depth: 0 - - uses: actions/setup-python@ece7cb06caefa5fff74198d8649806c4678c61a1 # v6 - with: - python-version: '3.12' - - name: Apply bounded API registration patch - run: | - python - <<'PY' - from pathlib import Path - - api_path = Path("src/sdp/api.py") - source = api_path.read_text(encoding="utf-8") - endpoint_marker = '"/plane/catalog-objects/{catalog_object_id}/data-management-profile"' - - if endpoint_marker not in source: - core_import = '''from sdp_core.data_management_evidence import ( - CriticalDataElementDraft, - DataOwnerAssignmentDraft, - DataQualityObservationDraft, - DataQualityRuleDraft, - ) - ''' - import_anchor = "from sdp_core.catalog_plane import (\n" - if import_anchor not in source: - raise SystemExit("catalog-plane import anchor not found") - source = source.replace(import_anchor, core_import + import_anchor, 1) - - service_import = '''from .data_management_evidence import ( - assign_data_owner, - build_data_management_profile, - define_data_quality_rule, - record_data_quality_observation, - register_critical_data_element, - ) - ''' - service_anchor = "from .config import get_app_config\n" - if service_anchor not in source: - raise SystemExit("config import anchor not found") - source = source.replace(service_anchor, service_import + service_anchor, 1) - - route_anchor = '@app.get("/ontology/term/{term}/graph")\n' - if route_anchor not in source: - raise SystemExit("ontology route anchor not found") - routes = '''@app.post("/plane/catalog-objects/{catalog_object_id}/data-owner-assignments") - def plane_assign_data_owner( - catalog_object_id: str, - payload: DataOwnerAssignmentDraft, - request: Request, - ) -> dict[str, Any]: - """Assign an evidence-backed data owner to one catalog object.""" - - actor = _plane_actor(request) - try: - return assign_data_owner(actor, catalog_object_id, payload).model_dump() - except PermissionError as exc: - raise HTTPException(status_code=403, detail=str(exc)) from exc - except KeyError as exc: - raise HTTPException(status_code=404, detail=str(exc)) from exc - except ValueError as exc: - raise HTTPException(status_code=400, detail=str(exc)) from exc - - - @app.post("/plane/catalog-objects/{catalog_object_id}/critical-data-elements") - def plane_register_critical_data_element( - catalog_object_id: str, - payload: CriticalDataElementDraft, - request: Request, - ) -> dict[str, Any]: - """Register an evidence-backed CDE under a catalog dataset.""" - - actor = _plane_actor(request) - try: - return register_critical_data_element( - actor, - catalog_object_id, - payload, - ).model_dump() - except PermissionError as exc: - raise HTTPException(status_code=403, detail=str(exc)) from exc - except KeyError as exc: - raise HTTPException(status_code=404, detail=str(exc)) from exc - except ValueError as exc: - raise HTTPException(status_code=400, detail=str(exc)) from exc - - - @app.post("/plane/critical-data-elements/{critical_data_element_id}/quality-rules") - def plane_define_data_quality_rule( - critical_data_element_id: str, - payload: DataQualityRuleDraft, - request: Request, - ) -> dict[str, Any]: - """Define one evidence-backed quality rule for a CDE.""" - - actor = _plane_actor(request) - try: - return define_data_quality_rule( - actor, - critical_data_element_id, - payload, - ).model_dump() - except PermissionError as exc: - raise HTTPException(status_code=403, detail=str(exc)) from exc - except KeyError as exc: - raise HTTPException(status_code=404, detail=str(exc)) from exc - except ValueError as exc: - raise HTTPException(status_code=400, detail=str(exc)) from exc - - - @app.post("/plane/quality-rules/{data_quality_rule_id}/observations") - def plane_record_data_quality_observation( - data_quality_rule_id: str, - payload: DataQualityObservationDraft, - request: Request, - ) -> dict[str, Any]: - """Append one immutable evidence-backed quality observation.""" - - actor = _plane_actor(request) - try: - return record_data_quality_observation( - actor, - data_quality_rule_id, - payload, - ).model_dump() - except PermissionError as exc: - raise HTTPException(status_code=403, detail=str(exc)) from exc - except KeyError as exc: - raise HTTPException(status_code=404, detail=str(exc)) from exc - except ValueError as exc: - raise HTTPException(status_code=400, detail=str(exc)) from exc - - - @app.get("/plane/catalog-objects/{catalog_object_id}/data-management-profile") - def plane_data_management_profile( - catalog_object_id: str, - request: Request, - ) -> dict[str, Any]: - """Return the explainable evidence-completeness profile.""" - - actor = _plane_actor(request) - try: - return build_data_management_profile(actor, catalog_object_id).model_dump() - except PermissionError as exc: - raise HTTPException(status_code=403, detail=str(exc)) from exc - except KeyError as exc: - raise HTTPException(status_code=404, detail=str(exc)) from exc - except ValueError as exc: - raise HTTPException(status_code=400, detail=str(exc)) from exc - - - ''' - source = source.replace(route_anchor, routes + route_anchor, 1) - api_path.write_text(source, encoding="utf-8") - - workflow_path = Path(".github/workflows/apply-data-management-routes.yml") - workflow_path.unlink() - PY - - name: Install hash-pinned test dependencies - run: python -m pip install --disable-pip-version-check --no-cache-dir --require-hashes -r requirements-test.txt - - name: Verify focused data-management acceptance - env: - PYTHONPATH: src - run: | - python -m pytest tests/test_data_management_evidence.py tests/test_data_management_migration.py - python -m compileall -q src - - name: Commit verified patch and remove one-shot workflow - run: | - git config user.name "cwl-data-management-agent" - git config user.email "cwl-data-management-agent@users.noreply.github.com" - git add src/sdp/api.py .github/workflows/apply-data-management-routes.yml - git commit -m "feat: expose data management evidence routes" - git push origin HEAD:feat/data-management-evidence-profile diff --git a/src/sdp/api.py b/src/sdp/api.py index 38e4778..e366a1e 100644 --- a/src/sdp/api.py +++ b/src/sdp/api.py @@ -16,6 +16,12 @@ enterprise_rbac_matrix, enterprise_readiness_manifest, ) +from sdp_core.data_management_evidence import ( + CriticalDataElementDraft, + DataOwnerAssignmentDraft, + DataQualityObservationDraft, + DataQualityRuleDraft, +) from sdp_core.catalog_plane import ( CatalogObjectCreateRequest, ConceptBindingDraft, @@ -33,6 +39,13 @@ list_catalog_objects, query_catalog_objects, ) +from .data_management_evidence import ( + assign_data_owner, + build_data_management_profile, + define_data_quality_rule, + record_data_quality_observation, + register_critical_data_element, +) from .config import get_app_config from .console import render_enterprise_console from .tenant_binding import TenantBindingError, bind_keyverse_tenant @@ -826,6 +839,100 @@ def plane_query_catalog_objects( raise HTTPException(status_code=400, detail=str(exc)) from exc +@app.post("/plane/catalog-objects/{catalog_object_id}/data-owner-assignments") +def plane_assign_data_owner( + catalog_object_id: str, + payload: DataOwnerAssignmentDraft, + request: Request, +) -> dict[str, Any]: + """Assign an evidence-backed data owner to one catalog object.""" + + actor = _plane_actor(request) + try: + return assign_data_owner(actor, catalog_object_id, payload).model_dump() + except PermissionError as exc: + raise HTTPException(status_code=403, detail=str(exc)) from exc + except KeyError as exc: + raise HTTPException(status_code=404, detail=str(exc)) from exc + except ValueError as exc: + raise HTTPException(status_code=400, detail=str(exc)) from exc + + +@app.post("/plane/catalog-objects/{catalog_object_id}/critical-data-elements") +def plane_register_critical_data_element( + catalog_object_id: str, + payload: CriticalDataElementDraft, + request: Request, +) -> dict[str, Any]: + """Register an evidence-backed CDE under a catalog dataset.""" + + actor = _plane_actor(request) + try: + return register_critical_data_element(actor, catalog_object_id, payload).model_dump() + except PermissionError as exc: + raise HTTPException(status_code=403, detail=str(exc)) from exc + except KeyError as exc: + raise HTTPException(status_code=404, detail=str(exc)) from exc + except ValueError as exc: + raise HTTPException(status_code=400, detail=str(exc)) from exc + + +@app.post("/plane/critical-data-elements/{critical_data_element_id}/quality-rules") +def plane_define_data_quality_rule( + critical_data_element_id: str, + payload: DataQualityRuleDraft, + request: Request, +) -> dict[str, Any]: + """Define one evidence-backed quality rule for a CDE.""" + + actor = _plane_actor(request) + try: + return define_data_quality_rule(actor, critical_data_element_id, payload).model_dump() + except PermissionError as exc: + raise HTTPException(status_code=403, detail=str(exc)) from exc + except KeyError as exc: + raise HTTPException(status_code=404, detail=str(exc)) from exc + except ValueError as exc: + raise HTTPException(status_code=400, detail=str(exc)) from exc + + +@app.post("/plane/quality-rules/{data_quality_rule_id}/observations") +def plane_record_data_quality_observation( + data_quality_rule_id: str, + payload: DataQualityObservationDraft, + request: Request, +) -> dict[str, Any]: + """Append one immutable evidence-backed quality observation.""" + + actor = _plane_actor(request) + try: + return record_data_quality_observation(actor, data_quality_rule_id, payload).model_dump() + except PermissionError as exc: + raise HTTPException(status_code=403, detail=str(exc)) from exc + except KeyError as exc: + raise HTTPException(status_code=404, detail=str(exc)) from exc + except ValueError as exc: + raise HTTPException(status_code=400, detail=str(exc)) from exc + + +@app.get("/plane/catalog-objects/{catalog_object_id}/data-management-profile") +def plane_data_management_profile( + catalog_object_id: str, + request: Request, +) -> dict[str, Any]: + """Return the explainable evidence-completeness profile.""" + + actor = _plane_actor(request) + try: + return build_data_management_profile(actor, catalog_object_id).model_dump() + except PermissionError as exc: + raise HTTPException(status_code=403, detail=str(exc)) from exc + except KeyError as exc: + raise HTTPException(status_code=404, detail=str(exc)) from exc + except ValueError as exc: + raise HTTPException(status_code=400, detail=str(exc)) from exc + + @app.get("/ontology/term/{term}/graph") def ontology_term_graph(term: str) -> dict[str, Any]: # Backed by the persistent graph store now; falls back to the legacy diff --git a/tests/conftest.py b/tests/conftest.py index 04b5231..eae3f6e 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -4,6 +4,10 @@ import pytest +from sdp.catalog_plane_store import ( + restore_memory_catalog_plane, + snapshot_memory_catalog_plane, +) from sdp.data_management_store import ( restore_memory_data_management, snapshot_memory_data_management, @@ -19,10 +23,12 @@ def _allow_demo_subject_header(monkeypatch: pytest.MonkeyPatch) -> None: @pytest.fixture(autouse=True) def _isolate_data_management_evidence() -> None: - """Restore process-local evidence rows after every test.""" + """Restore process-local catalog-plane and evidence rows after every test.""" - snapshot = snapshot_memory_data_management() + plane_snapshot = snapshot_memory_catalog_plane() + evidence_snapshot = snapshot_memory_data_management() try: yield finally: - restore_memory_data_management(snapshot) + restore_memory_catalog_plane(plane_snapshot) + restore_memory_data_management(evidence_snapshot) diff --git a/tests/test_data_management_migration.py b/tests/test_data_management_migration.py index 7a65c0f..db03b05 100644 --- a/tests/test_data_management_migration.py +++ b/tests/test_data_management_migration.py @@ -32,11 +32,15 @@ def test_data_management_evidence_migration_preserves_authority_and_provenance() """Every governance fact carries tenant, truth, time, and HTTPS evidence fields.""" assert MIGRATION.exists(), "0003 data-management evidence migration is required" - sql = MIGRATION.read_text(encoding="utf-8") + # Collapse whitespace so the column contracts are checked independently of + # SQL alignment while still requiring one declaration per governed table. + sql = " ".join(MIGRATION.read_text(encoding="utf-8").split()) assert sql.count("tenant_reference TEXT NOT NULL") >= 4 assert sql.count("truth_status TEXT NOT NULL") >= 4 assert sql.count("evidence_reference TEXT NOT NULL") >= 4 assert "CHECK (truth_status IN ('authoritative', 'observed', 'inferred', 'proposed'))" in sql assert "CHECK (evidence_reference LIKE 'https://%')" in sql - assert "FOREIGN KEY (catalog_object_id) REFERENCES catalog_objects" in sql + # Each governed table carries a column-level FK into the catalog plane so + # evidence rows cannot outlive their tenant-scoped parent object. + assert sql.count("catalog_object_id TEXT NOT NULL REFERENCES catalog_objects") >= 4 From 44a8bf42c205c665a64e302e92de6c1caca630f0 Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Wed, 26 Aug 2026 02:55:02 +0900 Subject: [PATCH 12/13] chore: retrigger review scheduler on green checks From 6cba648a0385abb07305ffb8483d61b2243982de Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Wed, 26 Aug 2026 03:45:00 +0900 Subject: [PATCH 13/13] =?UTF-8?q?fix(data-mgmt):=20review=20hardening=20?= =?UTF-8?q?=E2=80=94=20FK=20pragma,=20classification=20CHECK,=20docs?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - data_management_store._open_engine now enables PRAGMA foreign_keys=ON on SQLite so unit-test engines enforce the same cross-tenant/parent-row guarantees as Postgres in production. - migration 0003 adds a DB-level CHECK for data_classification, matching the DataClassification Literal (public/internal/confidential/restricted_pii/ restricted_financial). - README API list + implementation-compliance matrix now include the five /plane evidence endpoints. - build_data_management_profile documents the contractual factor asymmetry: ownership requires an active effective-time window; CDE/rule require authoritative truth; observations accept observed truth by nature. Full suite: 262 passed. --- README.md | 5 +++++ docs/implementation-compliance.md | 3 +++ migrations/0003_data_management_evidence.sql | 5 ++++- src/sdp/data_management_evidence.py | 9 +++++++++ src/sdp/data_management_store.py | 21 +++++++++++++++++--- 5 files changed, 39 insertions(+), 4 deletions(-) diff --git a/README.md b/README.md index b74a30e..c368c86 100644 --- a/README.md +++ b/README.md @@ -124,6 +124,11 @@ glossary/catalog row가 유지됩니다. - `GET /plane/catalog-objects/{catalog_object_id}` — 단건 조회 - `POST /plane/catalog-objects/{catalog_object_id}/document-kg-links` — 문서 KG 참조 연결 - `POST /plane/catalog-objects/{catalog_object_id}/concept-bindings` — 온톨로지 개념 연결 +- `POST /plane/catalog-objects/{catalog_object_id}/data-owner-assignments` — 데이터 오너 지정 (evidence-backed) +- `POST /plane/catalog-objects/{catalog_object_id}/critical-data-elements` — CDE 등록 +- `POST /plane/critical-data-elements/{critical_data_element_id}/quality-rules` — 품질 규칙 정의 +- `POST /plane/quality-rules/{data_quality_rule_id}/observations` — 품질 observation 기록 (immutable) +- `GET /plane/catalog-objects/{catalog_object_id}/data-management-profile` — evidence 완결도 프로파일 - `GET /plane/query?q=...` — title/alias/concept/document-KG id 검색 Production 필수 헤더: `Authorization: Bearer `, diff --git a/docs/implementation-compliance.md b/docs/implementation-compliance.md index b24cb69..6ad26f3 100644 --- a/docs/implementation-compliance.md +++ b/docs/implementation-compliance.md @@ -74,9 +74,12 @@ ### CAT-010 Ontology/catalog plane above the document KG (issue #13) - 대응: `POST/GET /plane/catalog-objects`, `GET /plane/catalog-objects/{catalog_object_id}`, `POST .../document-kg-links`, `POST .../concept-bindings`, `GET /plane/query` +- 데이터 관리 증거 확장(issue #74): `POST .../data-owner-assignments`, `POST .../critical-data-elements`, `POST /plane/critical-data-elements/{id}/quality-rules`, `POST /plane/quality-rules/{id}/observations`, `GET .../data-management-profile` - 핵심 코드: - [src/sdp/catalog_plane.py](src/sdp/catalog_plane.py) - [src/sdp/catalog_plane_store.py](src/sdp/catalog_plane_store.py) + - [src/sdp/data_management_evidence.py](src/sdp/data_management_evidence.py) + - [src/sdp/data_management_store.py](src/sdp/data_management_store.py) - [src/sdp/tenant_binding.py](src/sdp/tenant_binding.py) - [src/sdp_core/catalog_plane.py](src/sdp_core/catalog_plane.py) - [migrations/0002_ontology_catalog_plane.sql](migrations/0002_ontology_catalog_plane.sql) diff --git a/migrations/0003_data_management_evidence.sql b/migrations/0003_data_management_evidence.sql index eac18e9..b3b664c 100644 --- a/migrations/0003_data_management_evidence.sql +++ b/migrations/0003_data_management_evidence.sql @@ -33,7 +33,10 @@ CREATE TABLE IF NOT EXISTS critical_data_elements ( element_key TEXT NOT NULL, display_name TEXT NOT NULL, definition_text TEXT NOT NULL, - data_classification TEXT NOT NULL, + data_classification TEXT NOT NULL + CHECK (data_classification IN ( + 'public', 'internal', 'confidential', 'restricted_pii', 'restricted_financial' + )), evidence_reference TEXT NOT NULL CHECK (evidence_reference LIKE 'https://%'), truth_status TEXT NOT NULL CHECK (truth_status IN ('authoritative', 'observed', 'inferred', 'proposed')), diff --git a/src/sdp/data_management_evidence.py b/src/sdp/data_management_evidence.py index 2d60d4a..fe452a1 100644 --- a/src/sdp/data_management_evidence.py +++ b/src/sdp/data_management_evidence.py @@ -261,6 +261,15 @@ def build_data_management_profile( catalog_object_id=catalog_object_id, ) current_time = _now() + # Factor asymmetry is contractual, not accidental: + # - data_owner_present requires an *authoritative* assignment whose + # effective-time window covers `current_time`, because authority lapses + # when the window closes. + # - CDE/rule presence requires authoritative truth; these are definitions, + # not measurements, so "observed" does not make a definition exist. + # - observation presence accepts authoritative or observed truth because an + # observation row IS a measurement record — its native truth status is + # "observed" by construction. owner_present = any( row.truth_status == "authoritative" and row.valid_from <= current_time diff --git a/src/sdp/data_management_store.py b/src/sdp/data_management_store.py index 53fb647..6326122 100644 --- a/src/sdp/data_management_store.py +++ b/src/sdp/data_management_store.py @@ -24,11 +24,12 @@ ) try: - from sqlalchemy import create_engine, text + from sqlalchemy import create_engine, event, text from sqlalchemy.engine import Engine from sqlalchemy.exc import IntegrityError except ImportError: # pragma: no cover - optional graph extra create_engine = None # type: ignore[assignment] + event = None # type: ignore[assignment] text = None # type: ignore[assignment] Engine = Any # type: ignore[misc,assignment] @@ -75,11 +76,25 @@ def _sql_timestamp(value: datetime | None) -> str | None: def _open_engine(database_dsn: str) -> Engine: - """Open a SQLAlchemy engine or fail loud when the graph extra is absent.""" + """Open a SQLAlchemy engine or fail loud when the graph extra is absent. + + SQLite engines get ``PRAGMA foreign_keys=ON`` so relational cross-tenant + and parent-row guarantees hold during unit tests exactly as they do under + Postgres in production. + """ if create_engine is None: raise RuntimeError(_GRAPH_EXTRA_HINT) - return create_engine(database_dsn, future=True) + engine = create_engine(database_dsn, future=True) + if engine.dialect.name == "sqlite": + + @event.listens_for(engine, "connect") + def _set_sqlite_pragma(dbapi_connection, _connection_record) -> None: + cursor = dbapi_connection.cursor() + cursor.execute("PRAGMA foreign_keys=ON") + cursor.close() + + return engine def _decimal(value: Any) -> Decimal: