data_research.py 15 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405
  1. from __future__ import annotations
  2. from typing import Any
  3. from sqlalchemy.dialects.postgresql import JSONB, UUID
  4. from app import db
  5. from app.core.common.identifiers import new_governance_uid
  6. from app.core.common.timezone_utils import now_china_naive
  7. JOB_STATUSES = (
  8. "created",
  9. "queued",
  10. "extracting",
  11. "normalizing",
  12. "matching",
  13. "awaiting_review",
  14. "published",
  15. "partial",
  16. "failed",
  17. "cancelled",
  18. )
  19. def _iso(value: Any) -> str | None:
  20. return value.isoformat() if value is not None else None
  21. class IngestionSource(db.Model):
  22. __tablename__ = "ingestion_sources"
  23. __table_args__ = (
  24. db.CheckConstraint(
  25. "source_type IN ('database','file','ddl')",
  26. name="ck_ingestion_source_type",
  27. ),
  28. db.CheckConstraint(
  29. "status IN ('active','disabled')",
  30. name="ck_ingestion_source_status",
  31. ),
  32. {"schema": "public"},
  33. )
  34. uid = db.Column(UUID(as_uuid=False), primary_key=True, default=new_governance_uid)
  35. source_type = db.Column(db.String(20), nullable=False)
  36. name = db.Column(db.String(300), nullable=False)
  37. config = db.Column(JSONB, nullable=False, default=dict)
  38. permission_scope = db.Column(JSONB, nullable=False, default=dict)
  39. status = db.Column(db.String(20), nullable=False, default="active")
  40. created_by = db.Column(db.String(100))
  41. created_at = db.Column(db.DateTime, nullable=False, default=now_china_naive)
  42. updated_at = db.Column(db.DateTime, nullable=False, default=now_china_naive)
  43. def to_dict(self) -> dict[str, Any]:
  44. return {
  45. "uid": str(self.uid),
  46. "source_type": self.source_type,
  47. "name": self.name,
  48. "permission_scope": dict(self.permission_scope or {}),
  49. "status": self.status,
  50. "created_by": self.created_by,
  51. "created_at": _iso(self.created_at),
  52. "updated_at": _iso(self.updated_at),
  53. }
  54. class SourceArtifact(db.Model):
  55. __tablename__ = "source_artifacts"
  56. __table_args__ = (
  57. db.CheckConstraint("size_bytes >= 0", name="ck_source_artifact_size"),
  58. db.UniqueConstraint(
  59. "source_uid",
  60. "content_hash",
  61. "parser_version",
  62. name="uq_source_artifact_parser_content",
  63. ),
  64. {"schema": "public"},
  65. )
  66. uid = db.Column(UUID(as_uuid=False), primary_key=True, default=new_governance_uid)
  67. source_uid = db.Column(
  68. UUID(as_uuid=False),
  69. db.ForeignKey("public.ingestion_sources.uid"),
  70. nullable=False,
  71. )
  72. filename = db.Column(db.String(500), nullable=False)
  73. media_type = db.Column(db.String(200), nullable=False)
  74. size_bytes = db.Column(db.BigInteger, nullable=False)
  75. content_hash = db.Column(db.String(64), nullable=False)
  76. storage_ref = db.Column(db.Text, nullable=False)
  77. parser_version = db.Column(db.String(100), nullable=False)
  78. created_at = db.Column(db.DateTime, nullable=False, default=now_china_naive)
  79. def to_dict(self) -> dict[str, Any]:
  80. return {
  81. "uid": str(self.uid),
  82. "source_uid": str(self.source_uid),
  83. "filename": self.filename,
  84. "media_type": self.media_type,
  85. "size_bytes": int(self.size_bytes),
  86. "content_hash": self.content_hash,
  87. "parser_version": self.parser_version,
  88. "created_at": _iso(self.created_at),
  89. }
  90. class IngestionJob(db.Model):
  91. __tablename__ = "ingestion_jobs"
  92. __table_args__ = (
  93. db.CheckConstraint(
  94. "status IN (" + ",".join(f"'{value}'" for value in JOB_STATUSES) + ")",
  95. name="ck_ingestion_job_status",
  96. ),
  97. db.UniqueConstraint(
  98. "idempotency_key",
  99. name="uq_ingestion_job_idempotency_key",
  100. ),
  101. {"schema": "public"},
  102. )
  103. uid = db.Column(UUID(as_uuid=False), primary_key=True, default=new_governance_uid)
  104. source_uid = db.Column(
  105. UUID(as_uuid=False),
  106. db.ForeignKey("public.ingestion_sources.uid"),
  107. nullable=False,
  108. )
  109. artifact_uid = db.Column(
  110. UUID(as_uuid=False),
  111. db.ForeignKey("public.source_artifacts.uid"),
  112. )
  113. job_type = db.Column(db.String(30), nullable=False)
  114. status = db.Column(db.String(30), nullable=False, default="created")
  115. idempotency_key = db.Column(db.String(64), nullable=False, unique=True)
  116. parser_version = db.Column(db.String(100), nullable=False)
  117. parameters = db.Column(JSONB, nullable=False, default=dict)
  118. statistics = db.Column(JSONB, nullable=False, default=dict)
  119. last_error = db.Column(db.String(1000))
  120. actor_uid = db.Column(db.String(100))
  121. force_rerun = db.Column(db.Boolean, nullable=False, default=False)
  122. created_at = db.Column(db.DateTime, nullable=False, default=now_china_naive)
  123. updated_at = db.Column(db.DateTime, nullable=False, default=now_china_naive)
  124. started_at = db.Column(db.DateTime)
  125. finished_at = db.Column(db.DateTime)
  126. def to_dict(self) -> dict[str, Any]:
  127. return {
  128. "uid": str(self.uid),
  129. "source_uid": str(self.source_uid),
  130. "artifact_uid": str(self.artifact_uid) if self.artifact_uid else None,
  131. "job_type": self.job_type,
  132. "status": self.status,
  133. "parser_version": self.parser_version,
  134. "parameters": dict(self.parameters or {}),
  135. "statistics": dict(self.statistics or {}),
  136. "last_error": self.last_error,
  137. "actor_uid": self.actor_uid,
  138. "created_at": _iso(self.created_at),
  139. "updated_at": _iso(self.updated_at),
  140. "started_at": _iso(self.started_at),
  141. "finished_at": _iso(self.finished_at),
  142. }
  143. class EvidenceFragment(db.Model):
  144. __tablename__ = "evidence_fragments"
  145. __table_args__ = (
  146. db.CheckConstraint(
  147. "confidence IS NULL OR (confidence >= 0 AND confidence <= 1)",
  148. name="ck_evidence_confidence",
  149. ),
  150. {"schema": "public"},
  151. )
  152. uid = db.Column(UUID(as_uuid=False), primary_key=True, default=new_governance_uid)
  153. job_uid = db.Column(
  154. UUID(as_uuid=False),
  155. db.ForeignKey("public.ingestion_jobs.uid", ondelete="CASCADE"),
  156. nullable=False,
  157. )
  158. artifact_uid = db.Column(
  159. UUID(as_uuid=False),
  160. db.ForeignKey("public.source_artifacts.uid"),
  161. )
  162. locator = db.Column(JSONB, nullable=False)
  163. excerpt = db.Column(db.Text)
  164. confidence = db.Column(db.Float)
  165. created_at = db.Column(db.DateTime, nullable=False, default=now_china_naive)
  166. def to_dict(self) -> dict[str, Any]:
  167. return {
  168. "uid": str(self.uid),
  169. "job_uid": str(self.job_uid),
  170. "artifact_uid": str(self.artifact_uid) if self.artifact_uid else None,
  171. "locator": dict(self.locator or {}),
  172. "excerpt": self.excerpt,
  173. "confidence": self.confidence,
  174. "created_at": _iso(self.created_at),
  175. }
  176. class ExtractionCandidate(db.Model):
  177. __tablename__ = "extraction_candidates"
  178. __table_args__ = (
  179. db.CheckConstraint(
  180. "confidence >= 0 AND confidence <= 1",
  181. name="ck_candidate_confidence",
  182. ),
  183. db.CheckConstraint(
  184. "status IN ('candidate','accepted','rejected','ignored')",
  185. name="ck_candidate_status",
  186. ),
  187. {"schema": "public"},
  188. )
  189. uid = db.Column(UUID(as_uuid=False), primary_key=True, default=new_governance_uid)
  190. job_uid = db.Column(
  191. UUID(as_uuid=False),
  192. db.ForeignKey("public.ingestion_jobs.uid", ondelete="CASCADE"),
  193. nullable=False,
  194. )
  195. candidate_type = db.Column(db.String(40), nullable=False)
  196. normalized_data = db.Column(JSONB, nullable=False)
  197. evidence_uids = db.Column(JSONB, nullable=False, default=list)
  198. confidence = db.Column(db.Float, nullable=False)
  199. status = db.Column(db.String(20), nullable=False, default="candidate")
  200. created_at = db.Column(db.DateTime, nullable=False, default=now_china_naive)
  201. updated_at = db.Column(db.DateTime, nullable=False, default=now_china_naive)
  202. def to_dict(self) -> dict[str, Any]:
  203. return {
  204. "uid": str(self.uid),
  205. "job_uid": str(self.job_uid),
  206. "candidate_type": self.candidate_type,
  207. "normalized_data": dict(self.normalized_data or {}),
  208. "evidence_uids": list(self.evidence_uids or []),
  209. "confidence": float(self.confidence),
  210. "status": self.status,
  211. "created_at": _iso(self.created_at),
  212. "updated_at": _iso(self.updated_at),
  213. }
  214. class DataElement(db.Model):
  215. __tablename__ = "data_elements"
  216. __table_args__ = (
  217. db.CheckConstraint(
  218. "status IN ('draft','in_review','published','deprecated','retired')",
  219. name="ck_data_element_status",
  220. ),
  221. {"schema": "public"},
  222. )
  223. uid = db.Column(UUID(as_uuid=False), primary_key=True, default=new_governance_uid)
  224. code = db.Column(db.String(120), nullable=False, unique=True)
  225. current_version = db.Column(db.Integer, nullable=False, default=1)
  226. status = db.Column(db.String(20), nullable=False, default="draft")
  227. business_domain_uids = db.Column(JSONB, nullable=False, default=list)
  228. created_by = db.Column(db.String(100))
  229. created_at = db.Column(db.DateTime, nullable=False, default=now_china_naive)
  230. updated_at = db.Column(db.DateTime, nullable=False, default=now_china_naive)
  231. versions = db.relationship(
  232. "DataElementVersion",
  233. lazy="selectin",
  234. order_by="DataElementVersion.version",
  235. cascade="all, delete-orphan",
  236. )
  237. class DataElementVersion(db.Model):
  238. __tablename__ = "data_element_versions"
  239. __table_args__ = (
  240. db.UniqueConstraint(
  241. "data_element_uid",
  242. "version",
  243. name="uq_data_element_version",
  244. ),
  245. {"schema": "public"},
  246. )
  247. uid = db.Column(UUID(as_uuid=False), primary_key=True, default=new_governance_uid)
  248. data_element_uid = db.Column(
  249. UUID(as_uuid=False),
  250. db.ForeignKey("public.data_elements.uid"),
  251. nullable=False,
  252. )
  253. version = db.Column(db.Integer, nullable=False)
  254. status = db.Column(db.String(20), nullable=False)
  255. snapshot = db.Column(JSONB, nullable=False)
  256. evidence_uids = db.Column(JSONB, nullable=False, default=list)
  257. created_by = db.Column(db.String(100))
  258. created_at = db.Column(db.DateTime, nullable=False, default=now_china_naive)
  259. class CandidateDecisionRecord(db.Model):
  260. __tablename__ = "candidate_decisions"
  261. __table_args__ = (
  262. db.CheckConstraint(
  263. "action IN ('reuse','create','map','ignore')",
  264. name="ck_candidate_decision_action",
  265. ),
  266. db.UniqueConstraint("candidate_uid", name="uq_candidate_decision"),
  267. {"schema": "public"},
  268. )
  269. uid = db.Column(UUID(as_uuid=False), primary_key=True, default=new_governance_uid)
  270. candidate_uid = db.Column(
  271. UUID(as_uuid=False),
  272. db.ForeignKey("public.extraction_candidates.uid"),
  273. nullable=False,
  274. )
  275. action = db.Column(db.String(20), nullable=False)
  276. data_element_uid = db.Column(
  277. UUID(as_uuid=False),
  278. db.ForeignKey("public.data_elements.uid"),
  279. )
  280. evidence_uids = db.Column(JSONB, nullable=False, default=list)
  281. payload = db.Column(JSONB, nullable=False, default=dict)
  282. actor_uid = db.Column(db.String(100))
  283. reason = db.Column(db.String(1000))
  284. created_at = db.Column(db.DateTime, nullable=False, default=now_china_naive)
  285. class OntologyModel(db.Model):
  286. __tablename__ = "ontologies"
  287. __table_args__ = ({"schema": "public"},)
  288. uid = db.Column(UUID(as_uuid=False), primary_key=True, default=new_governance_uid)
  289. code = db.Column(db.String(120), nullable=False, unique=True)
  290. name = db.Column(db.String(300), nullable=False)
  291. owner_uid = db.Column(db.String(100), nullable=False)
  292. status = db.Column(db.String(20), nullable=False, default="draft")
  293. draft_revision = db.Column(db.Integer, nullable=False, default=0)
  294. active_version_uid = db.Column(UUID(as_uuid=False))
  295. created_by = db.Column(db.String(100))
  296. created_at = db.Column(db.DateTime, nullable=False, default=now_china_naive)
  297. updated_at = db.Column(db.DateTime, nullable=False, default=now_china_naive)
  298. class OntologyVersionModel(db.Model):
  299. __tablename__ = "ontology_versions"
  300. __table_args__ = (
  301. db.UniqueConstraint("ontology_uid", "version", name="uq_ontology_version"),
  302. {"schema": "public"},
  303. )
  304. uid = db.Column(UUID(as_uuid=False), primary_key=True, default=new_governance_uid)
  305. ontology_uid = db.Column(
  306. UUID(as_uuid=False), db.ForeignKey("public.ontologies.uid"), nullable=False
  307. )
  308. version = db.Column(db.Integer, nullable=False)
  309. parent_version_uid = db.Column(
  310. UUID(as_uuid=False), db.ForeignKey("public.ontology_versions.uid")
  311. )
  312. status = db.Column(db.String(20), nullable=False, default="draft")
  313. graph_document = db.Column(JSONB, nullable=False)
  314. content_hash = db.Column(db.String(64), nullable=False)
  315. created_by = db.Column(db.String(100))
  316. created_at = db.Column(db.DateTime, nullable=False, default=now_china_naive)
  317. published_at = db.Column(db.DateTime)
  318. class OntologyDomainLinkModel(db.Model):
  319. __tablename__ = "ontology_domain_links"
  320. __table_args__ = ({"schema": "public"},)
  321. ontology_uid = db.Column(
  322. UUID(as_uuid=False), db.ForeignKey("public.ontologies.uid"), primary_key=True
  323. )
  324. domain_uid = db.Column(UUID(as_uuid=False), primary_key=True)
  325. role = db.Column(db.String(20), nullable=False)
  326. class OntologyChangeSetModel(db.Model):
  327. __tablename__ = "ontology_change_sets"
  328. __table_args__ = ({"schema": "public"},)
  329. uid = db.Column(UUID(as_uuid=False), primary_key=True, default=new_governance_uid)
  330. ontology_uid = db.Column(
  331. UUID(as_uuid=False), db.ForeignKey("public.ontologies.uid"), nullable=False
  332. )
  333. base_version_uid = db.Column(UUID(as_uuid=False))
  334. status = db.Column(db.String(20), nullable=False, default="draft")
  335. changes = db.Column(JSONB, nullable=False, default=list)
  336. decisions = db.Column(JSONB, nullable=False, default=list)
  337. created_by = db.Column(db.String(100))
  338. created_at = db.Column(db.DateTime, nullable=False, default=now_china_naive)
  339. class OntologyPublishRunModel(db.Model):
  340. __tablename__ = "ontology_publish_runs"
  341. __table_args__ = ({"schema": "public"},)
  342. uid = db.Column(UUID(as_uuid=False), primary_key=True, default=new_governance_uid)
  343. ontology_uid = db.Column(
  344. UUID(as_uuid=False), db.ForeignKey("public.ontologies.uid"), nullable=False
  345. )
  346. version_uid = db.Column(
  347. UUID(as_uuid=False), db.ForeignKey("public.ontology_versions.uid"), nullable=False
  348. )
  349. idempotency_key = db.Column(db.String(128), nullable=False, unique=True)
  350. status = db.Column(db.String(20), nullable=False)
  351. validation_result = db.Column(JSONB, nullable=False, default=dict)
  352. actor_uid = db.Column(db.String(100))
  353. created_at = db.Column(db.DateTime, nullable=False, default=now_china_naive)
  354. finished_at = db.Column(db.DateTime)