from __future__ import annotations from typing import Any from sqlalchemy.dialects.postgresql import JSONB, UUID from app import db from app.core.common.identifiers import new_governance_uid from app.core.common.timezone_utils import now_china, now_china_naive JOB_STATUSES = ( "created", "queued", "extracting", "normalizing", "matching", "awaiting_review", "published", "partial", "failed", "cancelled", ) def _iso(value: Any) -> str | None: return value.isoformat() if value is not None else None class IngestionSource(db.Model): __tablename__ = "ingestion_sources" __table_args__ = ( db.CheckConstraint( "source_type IN ('database','file','ddl')", name="ck_ingestion_source_type", ), db.CheckConstraint( "status IN ('active','disabled')", name="ck_ingestion_source_status", ), {"schema": "public"}, ) uid = db.Column(UUID(as_uuid=False), primary_key=True, default=new_governance_uid) source_type = db.Column(db.String(20), nullable=False) name = db.Column(db.String(300), nullable=False) config = db.Column(JSONB, nullable=False, default=dict) permission_scope = db.Column(JSONB, nullable=False, default=dict) status = db.Column(db.String(20), nullable=False, default="active") created_by = db.Column(db.String(100)) created_at = db.Column(db.DateTime, nullable=False, default=now_china_naive) updated_at = db.Column(db.DateTime, nullable=False, default=now_china_naive) def to_dict(self) -> dict[str, Any]: return { "uid": str(self.uid), "source_type": self.source_type, "name": self.name, "permission_scope": dict(self.permission_scope or {}), "status": self.status, "created_by": self.created_by, "created_at": _iso(self.created_at), "updated_at": _iso(self.updated_at), } class SourceArtifact(db.Model): __tablename__ = "source_artifacts" __table_args__ = ( db.CheckConstraint("size_bytes >= 0", name="ck_source_artifact_size"), db.UniqueConstraint( "source_uid", "content_hash", "parser_version", name="uq_source_artifact_parser_content", ), {"schema": "public"}, ) uid = db.Column(UUID(as_uuid=False), primary_key=True, default=new_governance_uid) source_uid = db.Column( UUID(as_uuid=False), db.ForeignKey("public.ingestion_sources.uid"), nullable=False, ) filename = db.Column(db.String(500), nullable=False) media_type = db.Column(db.String(200), nullable=False) size_bytes = db.Column(db.BigInteger, nullable=False) content_hash = db.Column(db.String(64), nullable=False) storage_ref = db.Column(db.Text, nullable=False) parser_version = db.Column(db.String(100), nullable=False) created_at = db.Column(db.DateTime, nullable=False, default=now_china_naive) def to_dict(self) -> dict[str, Any]: return { "uid": str(self.uid), "source_uid": str(self.source_uid), "filename": self.filename, "media_type": self.media_type, "size_bytes": int(self.size_bytes), "content_hash": self.content_hash, "parser_version": self.parser_version, "created_at": _iso(self.created_at), } class IngestionJob(db.Model): __tablename__ = "ingestion_jobs" __table_args__ = ( db.CheckConstraint( "status IN (" + ",".join(f"'{value}'" for value in JOB_STATUSES) + ")", name="ck_ingestion_job_status", ), db.UniqueConstraint( "idempotency_key", name="uq_ingestion_job_idempotency_key", ), db.CheckConstraint( "attempt_count >= 0", name="ck_ingestion_job_attempt_count", ), {"schema": "public"}, ) uid = db.Column(UUID(as_uuid=False), primary_key=True, default=new_governance_uid) source_uid = db.Column( UUID(as_uuid=False), db.ForeignKey("public.ingestion_sources.uid"), nullable=False, ) artifact_uid = db.Column( UUID(as_uuid=False), db.ForeignKey("public.source_artifacts.uid"), ) job_type = db.Column(db.String(30), nullable=False) status = db.Column(db.String(30), nullable=False, default="created") idempotency_key = db.Column(db.String(64), nullable=False, unique=True) parser_version = db.Column(db.String(100), nullable=False) parameters = db.Column(JSONB, nullable=False, default=dict) statistics = db.Column(JSONB, nullable=False, default=dict) last_error = db.Column(db.String(1000)) attempt_count = db.Column(db.Integer, nullable=False, default=0) failure_stage = db.Column(db.String(30)) actor_uid = db.Column(db.String(100)) force_rerun = db.Column(db.Boolean, nullable=False, default=False) created_at = db.Column(db.DateTime, nullable=False, default=now_china_naive) updated_at = db.Column(db.DateTime, nullable=False, default=now_china_naive) started_at = db.Column(db.DateTime) finished_at = db.Column(db.DateTime) def to_dict(self) -> dict[str, Any]: return { "uid": str(self.uid), "source_uid": str(self.source_uid), "artifact_uid": str(self.artifact_uid) if self.artifact_uid else None, "job_type": self.job_type, "status": self.status, "parser_version": self.parser_version, "parameters": dict(self.parameters or {}), "statistics": dict(self.statistics or {}), "last_error": self.last_error, "attempt_count": int(self.attempt_count or 0), "failure_stage": self.failure_stage, "actor_uid": self.actor_uid, "created_at": _iso(self.created_at), "updated_at": _iso(self.updated_at), "started_at": _iso(self.started_at), "finished_at": _iso(self.finished_at), } class CatalogSnapshot(db.Model): __tablename__ = "catalog_snapshots" __table_args__ = ( db.CheckConstraint( "attempt > 0", name="ck_catalog_snapshot_attempt", ), db.UniqueConstraint( "job_uid", "attempt", name="uq_catalog_snapshot_job_attempt", ), {"schema": "public"}, ) uid = db.Column( UUID(as_uuid=False), primary_key=True, default=new_governance_uid, ) job_uid = db.Column( UUID(as_uuid=False), db.ForeignKey("public.ingestion_jobs.uid", ondelete="CASCADE"), nullable=False, ) source_uid = db.Column( UUID(as_uuid=False), db.ForeignKey("public.ingestion_sources.uid"), nullable=False, ) attempt = db.Column(db.Integer, nullable=False) database_type = db.Column(db.String(20), nullable=False) content_hash = db.Column(db.String(64), nullable=False) snapshot = db.Column(JSONB, nullable=False) created_at = db.Column( db.DateTime, nullable=False, default=now_china_naive, ) def to_dict(self) -> dict[str, Any]: return { "uid": str(self.uid), "job_uid": str(self.job_uid), "source_uid": str(self.source_uid), "attempt": int(self.attempt), "database_type": self.database_type, "content_hash": self.content_hash, "snapshot": dict(self.snapshot or {}), "created_at": _iso(self.created_at), } class DeviceAsset(db.Model): __tablename__ = "device_assets" __table_args__ = ( db.CheckConstraint( "asset_type IN (" "'device','component','measurement_point'," "'alarm','maintenance_record'" ")", name="ck_device_asset_type", ), db.CheckConstraint( "status IN ('active','retired')", name="ck_device_asset_status", ), db.CheckConstraint( "current_version > 0", name="ck_device_asset_current_version", ), {"schema": "public"}, ) uid = db.Column( UUID(as_uuid=False), primary_key=True, default=new_governance_uid, ) asset_type = db.Column(db.String(40), nullable=False) name = db.Column(db.String(300), nullable=False) status = db.Column(db.String(20), nullable=False, default="active") current_version = db.Column(db.Integer, nullable=False, default=1) content_hash = db.Column(db.String(64), nullable=False) location = db.Column(db.String(300)) organization = db.Column(db.String(300)) responsible_person = db.Column(db.String(300)) attributes = db.Column(JSONB, nullable=False, default=dict) created_by = db.Column(db.String(100)) updated_by = db.Column(db.String(100)) created_at = db.Column( db.DateTime(timezone=True), nullable=False, default=now_china, ) updated_at = db.Column( db.DateTime(timezone=True), nullable=False, default=now_china, ) class DeviceAssetSourceMapping(db.Model): __tablename__ = "device_asset_source_mappings" __table_args__ = ( db.CheckConstraint( "asset_type IN (" "'device','component','measurement_point'," "'alarm','maintenance_record'" ")", name="ck_device_asset_mapping_type", ), db.UniqueConstraint( "source_uid", "source_entity", "asset_type", "source_code", name="uq_device_asset_source_identity", ), {"schema": "public"}, ) uid = db.Column( UUID(as_uuid=False), primary_key=True, default=new_governance_uid, ) asset_uid = db.Column( UUID(as_uuid=False), db.ForeignKey("public.device_assets.uid", ondelete="CASCADE"), nullable=False, ) source_uid = db.Column( UUID(as_uuid=False), db.ForeignKey("public.ingestion_sources.uid"), nullable=False, ) source_entity = db.Column(db.String(300), nullable=False) asset_type = db.Column(db.String(40), nullable=False) source_code = db.Column(db.String(300), nullable=False) source_updated_at = db.Column(db.DateTime(timezone=True)) first_seen_at = db.Column( db.DateTime(timezone=True), nullable=False, default=now_china, ) last_seen_at = db.Column( db.DateTime(timezone=True), nullable=False, default=now_china, ) class DeviceAssetVersion(db.Model): __tablename__ = "device_asset_versions" __table_args__ = ( db.CheckConstraint( "version > 0", name="ck_device_asset_version_number", ), db.UniqueConstraint( "asset_uid", "version", name="uq_device_asset_version", ), {"schema": "public"}, ) uid = db.Column( UUID(as_uuid=False), primary_key=True, default=new_governance_uid, ) asset_uid = db.Column( UUID(as_uuid=False), db.ForeignKey("public.device_assets.uid", ondelete="CASCADE"), nullable=False, ) version = db.Column(db.Integer, nullable=False) content_hash = db.Column(db.String(64), nullable=False) snapshot = db.Column(JSONB, nullable=False) source_mapping_uid = db.Column( UUID(as_uuid=False), db.ForeignKey( "public.device_asset_source_mappings.uid", ondelete="RESTRICT", ), nullable=False, ) actor_uid = db.Column(db.String(100)) created_at = db.Column( db.DateTime(timezone=True), nullable=False, default=now_china, ) class EvidenceFragment(db.Model): __tablename__ = "evidence_fragments" __table_args__ = ( db.CheckConstraint( "confidence IS NULL OR (confidence >= 0 AND confidence <= 1)", name="ck_evidence_confidence", ), {"schema": "public"}, ) uid = db.Column(UUID(as_uuid=False), primary_key=True, default=new_governance_uid) job_uid = db.Column( UUID(as_uuid=False), db.ForeignKey("public.ingestion_jobs.uid", ondelete="CASCADE"), nullable=False, ) artifact_uid = db.Column( UUID(as_uuid=False), db.ForeignKey("public.source_artifacts.uid"), ) locator = db.Column(JSONB, nullable=False) excerpt = db.Column(db.Text) confidence = db.Column(db.Float) created_at = db.Column(db.DateTime, nullable=False, default=now_china_naive) def to_dict(self) -> dict[str, Any]: return { "uid": str(self.uid), "job_uid": str(self.job_uid), "artifact_uid": str(self.artifact_uid) if self.artifact_uid else None, "locator": dict(self.locator or {}), "excerpt": self.excerpt, "confidence": self.confidence, "created_at": _iso(self.created_at), } class ExtractionCandidate(db.Model): __tablename__ = "extraction_candidates" __table_args__ = ( db.CheckConstraint( "confidence >= 0 AND confidence <= 1", name="ck_candidate_confidence", ), db.CheckConstraint( "status IN ('candidate','accepted','rejected','ignored')", name="ck_candidate_status", ), {"schema": "public"}, ) uid = db.Column(UUID(as_uuid=False), primary_key=True, default=new_governance_uid) job_uid = db.Column( UUID(as_uuid=False), db.ForeignKey("public.ingestion_jobs.uid", ondelete="CASCADE"), nullable=False, ) candidate_type = db.Column(db.String(40), nullable=False) normalized_data = db.Column(JSONB, nullable=False) evidence_uids = db.Column(JSONB, nullable=False, default=list) confidence = db.Column(db.Float, nullable=False) status = db.Column(db.String(20), nullable=False, default="candidate") created_at = db.Column(db.DateTime, nullable=False, default=now_china_naive) updated_at = db.Column(db.DateTime, nullable=False, default=now_china_naive) def to_dict(self) -> dict[str, Any]: return { "uid": str(self.uid), "job_uid": str(self.job_uid), "candidate_type": self.candidate_type, "normalized_data": dict(self.normalized_data or {}), "evidence_uids": list(self.evidence_uids or []), "confidence": float(self.confidence), "status": self.status, "created_at": _iso(self.created_at), "updated_at": _iso(self.updated_at), } class DataElement(db.Model): __tablename__ = "data_elements" __table_args__ = ( db.CheckConstraint( "status IN ('draft','in_review','published','deprecated','retired')", name="ck_data_element_status", ), {"schema": "public"}, ) uid = db.Column(UUID(as_uuid=False), primary_key=True, default=new_governance_uid) code = db.Column(db.String(120), nullable=False, unique=True) current_version = db.Column(db.Integer, nullable=False, default=1) status = db.Column(db.String(20), nullable=False, default="draft") business_domain_uids = db.Column(JSONB, nullable=False, default=list) created_by = db.Column(db.String(100)) created_at = db.Column(db.DateTime, nullable=False, default=now_china_naive) updated_at = db.Column(db.DateTime, nullable=False, default=now_china_naive) versions = db.relationship( "DataElementVersion", lazy="selectin", order_by="DataElementVersion.version", cascade="all, delete-orphan", ) class DataElementVersion(db.Model): __tablename__ = "data_element_versions" __table_args__ = ( db.UniqueConstraint( "data_element_uid", "version", name="uq_data_element_version", ), {"schema": "public"}, ) uid = db.Column(UUID(as_uuid=False), primary_key=True, default=new_governance_uid) data_element_uid = db.Column( UUID(as_uuid=False), db.ForeignKey("public.data_elements.uid"), nullable=False, ) version = db.Column(db.Integer, nullable=False) status = db.Column(db.String(20), nullable=False) snapshot = db.Column(JSONB, nullable=False) evidence_uids = db.Column(JSONB, nullable=False, default=list) created_by = db.Column(db.String(100)) created_at = db.Column(db.DateTime, nullable=False, default=now_china_naive) class CandidateDecisionRecord(db.Model): __tablename__ = "candidate_decisions" __table_args__ = ( db.CheckConstraint( "action IN ('reuse','create','map','ignore')", name="ck_candidate_decision_action", ), db.UniqueConstraint("candidate_uid", name="uq_candidate_decision"), {"schema": "public"}, ) uid = db.Column(UUID(as_uuid=False), primary_key=True, default=new_governance_uid) candidate_uid = db.Column( UUID(as_uuid=False), db.ForeignKey("public.extraction_candidates.uid"), nullable=False, ) action = db.Column(db.String(20), nullable=False) data_element_uid = db.Column( UUID(as_uuid=False), db.ForeignKey("public.data_elements.uid"), ) evidence_uids = db.Column(JSONB, nullable=False, default=list) payload = db.Column(JSONB, nullable=False, default=dict) actor_uid = db.Column(db.String(100)) reason = db.Column(db.String(1000)) created_at = db.Column(db.DateTime, nullable=False, default=now_china_naive) class OntologyModel(db.Model): __tablename__ = "ontologies" __table_args__ = ({"schema": "public"},) uid = db.Column(UUID(as_uuid=False), primary_key=True, default=new_governance_uid) code = db.Column(db.String(120), nullable=False, unique=True) name = db.Column(db.String(300), nullable=False) owner_uid = db.Column(db.String(100), nullable=False) status = db.Column(db.String(20), nullable=False, default="draft") draft_revision = db.Column(db.Integer, nullable=False, default=0) active_version_uid = db.Column(UUID(as_uuid=False)) created_by = db.Column(db.String(100)) created_at = db.Column(db.DateTime, nullable=False, default=now_china_naive) updated_at = db.Column(db.DateTime, nullable=False, default=now_china_naive) class OntologyVersionModel(db.Model): __tablename__ = "ontology_versions" __table_args__ = ( db.UniqueConstraint("ontology_uid", "version", name="uq_ontology_version"), {"schema": "public"}, ) uid = db.Column(UUID(as_uuid=False), primary_key=True, default=new_governance_uid) ontology_uid = db.Column( UUID(as_uuid=False), db.ForeignKey("public.ontologies.uid"), nullable=False ) version = db.Column(db.Integer, nullable=False) parent_version_uid = db.Column( UUID(as_uuid=False), db.ForeignKey("public.ontology_versions.uid") ) status = db.Column(db.String(20), nullable=False, default="draft") graph_document = db.Column(JSONB, nullable=False) content_hash = db.Column(db.String(64), nullable=False) created_by = db.Column(db.String(100)) created_at = db.Column(db.DateTime, nullable=False, default=now_china_naive) published_at = db.Column(db.DateTime) class OntologyDomainLinkModel(db.Model): __tablename__ = "ontology_domain_links" __table_args__ = ({"schema": "public"},) ontology_uid = db.Column( UUID(as_uuid=False), db.ForeignKey("public.ontologies.uid"), primary_key=True ) domain_uid = db.Column(UUID(as_uuid=False), primary_key=True) role = db.Column(db.String(20), nullable=False) class OntologyChangeSetModel(db.Model): __tablename__ = "ontology_change_sets" __table_args__ = ({"schema": "public"},) uid = db.Column(UUID(as_uuid=False), primary_key=True, default=new_governance_uid) ontology_uid = db.Column( UUID(as_uuid=False), db.ForeignKey("public.ontologies.uid"), nullable=False ) base_version_uid = db.Column(UUID(as_uuid=False)) status = db.Column(db.String(20), nullable=False, default="draft") changes = db.Column(JSONB, nullable=False, default=list) decisions = db.Column(JSONB, nullable=False, default=list) created_by = db.Column(db.String(100)) created_at = db.Column(db.DateTime, nullable=False, default=now_china_naive) class OntologyPublishRunModel(db.Model): __tablename__ = "ontology_publish_runs" __table_args__ = ({"schema": "public"},) uid = db.Column(UUID(as_uuid=False), primary_key=True, default=new_governance_uid) ontology_uid = db.Column( UUID(as_uuid=False), db.ForeignKey("public.ontologies.uid"), nullable=False ) version_uid = db.Column( UUID(as_uuid=False), db.ForeignKey("public.ontology_versions.uid"), nullable=False ) idempotency_key = db.Column(db.String(128), nullable=False, unique=True) status = db.Column(db.String(20), nullable=False) validation_result = db.Column(JSONB, nullable=False, default=dict) actor_uid = db.Column(db.String(100)) created_at = db.Column(db.DateTime, nullable=False, default=now_china_naive) finished_at = db.Column(db.DateTime)