data_research.py 17 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466
  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. db.CheckConstraint(
  102. "attempt_count >= 0",
  103. name="ck_ingestion_job_attempt_count",
  104. ),
  105. {"schema": "public"},
  106. )
  107. uid = db.Column(UUID(as_uuid=False), primary_key=True, default=new_governance_uid)
  108. source_uid = db.Column(
  109. UUID(as_uuid=False),
  110. db.ForeignKey("public.ingestion_sources.uid"),
  111. nullable=False,
  112. )
  113. artifact_uid = db.Column(
  114. UUID(as_uuid=False),
  115. db.ForeignKey("public.source_artifacts.uid"),
  116. )
  117. job_type = db.Column(db.String(30), nullable=False)
  118. status = db.Column(db.String(30), nullable=False, default="created")
  119. idempotency_key = db.Column(db.String(64), nullable=False, unique=True)
  120. parser_version = db.Column(db.String(100), nullable=False)
  121. parameters = db.Column(JSONB, nullable=False, default=dict)
  122. statistics = db.Column(JSONB, nullable=False, default=dict)
  123. last_error = db.Column(db.String(1000))
  124. attempt_count = db.Column(db.Integer, nullable=False, default=0)
  125. failure_stage = db.Column(db.String(30))
  126. actor_uid = db.Column(db.String(100))
  127. force_rerun = db.Column(db.Boolean, nullable=False, default=False)
  128. created_at = db.Column(db.DateTime, nullable=False, default=now_china_naive)
  129. updated_at = db.Column(db.DateTime, nullable=False, default=now_china_naive)
  130. started_at = db.Column(db.DateTime)
  131. finished_at = db.Column(db.DateTime)
  132. def to_dict(self) -> dict[str, Any]:
  133. return {
  134. "uid": str(self.uid),
  135. "source_uid": str(self.source_uid),
  136. "artifact_uid": str(self.artifact_uid) if self.artifact_uid else None,
  137. "job_type": self.job_type,
  138. "status": self.status,
  139. "parser_version": self.parser_version,
  140. "parameters": dict(self.parameters or {}),
  141. "statistics": dict(self.statistics or {}),
  142. "last_error": self.last_error,
  143. "attempt_count": int(self.attempt_count or 0),
  144. "failure_stage": self.failure_stage,
  145. "actor_uid": self.actor_uid,
  146. "created_at": _iso(self.created_at),
  147. "updated_at": _iso(self.updated_at),
  148. "started_at": _iso(self.started_at),
  149. "finished_at": _iso(self.finished_at),
  150. }
  151. class CatalogSnapshot(db.Model):
  152. __tablename__ = "catalog_snapshots"
  153. __table_args__ = (
  154. db.CheckConstraint(
  155. "attempt > 0",
  156. name="ck_catalog_snapshot_attempt",
  157. ),
  158. db.UniqueConstraint(
  159. "job_uid",
  160. "attempt",
  161. name="uq_catalog_snapshot_job_attempt",
  162. ),
  163. {"schema": "public"},
  164. )
  165. uid = db.Column(
  166. UUID(as_uuid=False),
  167. primary_key=True,
  168. default=new_governance_uid,
  169. )
  170. job_uid = db.Column(
  171. UUID(as_uuid=False),
  172. db.ForeignKey("public.ingestion_jobs.uid", ondelete="CASCADE"),
  173. nullable=False,
  174. )
  175. source_uid = db.Column(
  176. UUID(as_uuid=False),
  177. db.ForeignKey("public.ingestion_sources.uid"),
  178. nullable=False,
  179. )
  180. attempt = db.Column(db.Integer, nullable=False)
  181. database_type = db.Column(db.String(20), nullable=False)
  182. content_hash = db.Column(db.String(64), nullable=False)
  183. snapshot = db.Column(JSONB, nullable=False)
  184. created_at = db.Column(
  185. db.DateTime,
  186. nullable=False,
  187. default=now_china_naive,
  188. )
  189. def to_dict(self) -> dict[str, Any]:
  190. return {
  191. "uid": str(self.uid),
  192. "job_uid": str(self.job_uid),
  193. "source_uid": str(self.source_uid),
  194. "attempt": int(self.attempt),
  195. "database_type": self.database_type,
  196. "content_hash": self.content_hash,
  197. "snapshot": dict(self.snapshot or {}),
  198. "created_at": _iso(self.created_at),
  199. }
  200. class EvidenceFragment(db.Model):
  201. __tablename__ = "evidence_fragments"
  202. __table_args__ = (
  203. db.CheckConstraint(
  204. "confidence IS NULL OR (confidence >= 0 AND confidence <= 1)",
  205. name="ck_evidence_confidence",
  206. ),
  207. {"schema": "public"},
  208. )
  209. uid = db.Column(UUID(as_uuid=False), primary_key=True, default=new_governance_uid)
  210. job_uid = db.Column(
  211. UUID(as_uuid=False),
  212. db.ForeignKey("public.ingestion_jobs.uid", ondelete="CASCADE"),
  213. nullable=False,
  214. )
  215. artifact_uid = db.Column(
  216. UUID(as_uuid=False),
  217. db.ForeignKey("public.source_artifacts.uid"),
  218. )
  219. locator = db.Column(JSONB, nullable=False)
  220. excerpt = db.Column(db.Text)
  221. confidence = db.Column(db.Float)
  222. created_at = db.Column(db.DateTime, nullable=False, default=now_china_naive)
  223. def to_dict(self) -> dict[str, Any]:
  224. return {
  225. "uid": str(self.uid),
  226. "job_uid": str(self.job_uid),
  227. "artifact_uid": str(self.artifact_uid) if self.artifact_uid else None,
  228. "locator": dict(self.locator or {}),
  229. "excerpt": self.excerpt,
  230. "confidence": self.confidence,
  231. "created_at": _iso(self.created_at),
  232. }
  233. class ExtractionCandidate(db.Model):
  234. __tablename__ = "extraction_candidates"
  235. __table_args__ = (
  236. db.CheckConstraint(
  237. "confidence >= 0 AND confidence <= 1",
  238. name="ck_candidate_confidence",
  239. ),
  240. db.CheckConstraint(
  241. "status IN ('candidate','accepted','rejected','ignored')",
  242. name="ck_candidate_status",
  243. ),
  244. {"schema": "public"},
  245. )
  246. uid = db.Column(UUID(as_uuid=False), primary_key=True, default=new_governance_uid)
  247. job_uid = db.Column(
  248. UUID(as_uuid=False),
  249. db.ForeignKey("public.ingestion_jobs.uid", ondelete="CASCADE"),
  250. nullable=False,
  251. )
  252. candidate_type = db.Column(db.String(40), nullable=False)
  253. normalized_data = db.Column(JSONB, nullable=False)
  254. evidence_uids = db.Column(JSONB, nullable=False, default=list)
  255. confidence = db.Column(db.Float, nullable=False)
  256. status = db.Column(db.String(20), nullable=False, default="candidate")
  257. created_at = db.Column(db.DateTime, nullable=False, default=now_china_naive)
  258. updated_at = db.Column(db.DateTime, nullable=False, default=now_china_naive)
  259. def to_dict(self) -> dict[str, Any]:
  260. return {
  261. "uid": str(self.uid),
  262. "job_uid": str(self.job_uid),
  263. "candidate_type": self.candidate_type,
  264. "normalized_data": dict(self.normalized_data or {}),
  265. "evidence_uids": list(self.evidence_uids or []),
  266. "confidence": float(self.confidence),
  267. "status": self.status,
  268. "created_at": _iso(self.created_at),
  269. "updated_at": _iso(self.updated_at),
  270. }
  271. class DataElement(db.Model):
  272. __tablename__ = "data_elements"
  273. __table_args__ = (
  274. db.CheckConstraint(
  275. "status IN ('draft','in_review','published','deprecated','retired')",
  276. name="ck_data_element_status",
  277. ),
  278. {"schema": "public"},
  279. )
  280. uid = db.Column(UUID(as_uuid=False), primary_key=True, default=new_governance_uid)
  281. code = db.Column(db.String(120), nullable=False, unique=True)
  282. current_version = db.Column(db.Integer, nullable=False, default=1)
  283. status = db.Column(db.String(20), nullable=False, default="draft")
  284. business_domain_uids = db.Column(JSONB, nullable=False, default=list)
  285. created_by = db.Column(db.String(100))
  286. created_at = db.Column(db.DateTime, nullable=False, default=now_china_naive)
  287. updated_at = db.Column(db.DateTime, nullable=False, default=now_china_naive)
  288. versions = db.relationship(
  289. "DataElementVersion",
  290. lazy="selectin",
  291. order_by="DataElementVersion.version",
  292. cascade="all, delete-orphan",
  293. )
  294. class DataElementVersion(db.Model):
  295. __tablename__ = "data_element_versions"
  296. __table_args__ = (
  297. db.UniqueConstraint(
  298. "data_element_uid",
  299. "version",
  300. name="uq_data_element_version",
  301. ),
  302. {"schema": "public"},
  303. )
  304. uid = db.Column(UUID(as_uuid=False), primary_key=True, default=new_governance_uid)
  305. data_element_uid = db.Column(
  306. UUID(as_uuid=False),
  307. db.ForeignKey("public.data_elements.uid"),
  308. nullable=False,
  309. )
  310. version = db.Column(db.Integer, nullable=False)
  311. status = db.Column(db.String(20), nullable=False)
  312. snapshot = db.Column(JSONB, nullable=False)
  313. evidence_uids = db.Column(JSONB, nullable=False, default=list)
  314. created_by = db.Column(db.String(100))
  315. created_at = db.Column(db.DateTime, nullable=False, default=now_china_naive)
  316. class CandidateDecisionRecord(db.Model):
  317. __tablename__ = "candidate_decisions"
  318. __table_args__ = (
  319. db.CheckConstraint(
  320. "action IN ('reuse','create','map','ignore')",
  321. name="ck_candidate_decision_action",
  322. ),
  323. db.UniqueConstraint("candidate_uid", name="uq_candidate_decision"),
  324. {"schema": "public"},
  325. )
  326. uid = db.Column(UUID(as_uuid=False), primary_key=True, default=new_governance_uid)
  327. candidate_uid = db.Column(
  328. UUID(as_uuid=False),
  329. db.ForeignKey("public.extraction_candidates.uid"),
  330. nullable=False,
  331. )
  332. action = db.Column(db.String(20), nullable=False)
  333. data_element_uid = db.Column(
  334. UUID(as_uuid=False),
  335. db.ForeignKey("public.data_elements.uid"),
  336. )
  337. evidence_uids = db.Column(JSONB, nullable=False, default=list)
  338. payload = db.Column(JSONB, nullable=False, default=dict)
  339. actor_uid = db.Column(db.String(100))
  340. reason = db.Column(db.String(1000))
  341. created_at = db.Column(db.DateTime, nullable=False, default=now_china_naive)
  342. class OntologyModel(db.Model):
  343. __tablename__ = "ontologies"
  344. __table_args__ = ({"schema": "public"},)
  345. uid = db.Column(UUID(as_uuid=False), primary_key=True, default=new_governance_uid)
  346. code = db.Column(db.String(120), nullable=False, unique=True)
  347. name = db.Column(db.String(300), nullable=False)
  348. owner_uid = db.Column(db.String(100), nullable=False)
  349. status = db.Column(db.String(20), nullable=False, default="draft")
  350. draft_revision = db.Column(db.Integer, nullable=False, default=0)
  351. active_version_uid = db.Column(UUID(as_uuid=False))
  352. created_by = db.Column(db.String(100))
  353. created_at = db.Column(db.DateTime, nullable=False, default=now_china_naive)
  354. updated_at = db.Column(db.DateTime, nullable=False, default=now_china_naive)
  355. class OntologyVersionModel(db.Model):
  356. __tablename__ = "ontology_versions"
  357. __table_args__ = (
  358. db.UniqueConstraint("ontology_uid", "version", name="uq_ontology_version"),
  359. {"schema": "public"},
  360. )
  361. uid = db.Column(UUID(as_uuid=False), primary_key=True, default=new_governance_uid)
  362. ontology_uid = db.Column(
  363. UUID(as_uuid=False), db.ForeignKey("public.ontologies.uid"), nullable=False
  364. )
  365. version = db.Column(db.Integer, nullable=False)
  366. parent_version_uid = db.Column(
  367. UUID(as_uuid=False), db.ForeignKey("public.ontology_versions.uid")
  368. )
  369. status = db.Column(db.String(20), nullable=False, default="draft")
  370. graph_document = db.Column(JSONB, nullable=False)
  371. content_hash = db.Column(db.String(64), nullable=False)
  372. created_by = db.Column(db.String(100))
  373. created_at = db.Column(db.DateTime, nullable=False, default=now_china_naive)
  374. published_at = db.Column(db.DateTime)
  375. class OntologyDomainLinkModel(db.Model):
  376. __tablename__ = "ontology_domain_links"
  377. __table_args__ = ({"schema": "public"},)
  378. ontology_uid = db.Column(
  379. UUID(as_uuid=False), db.ForeignKey("public.ontologies.uid"), primary_key=True
  380. )
  381. domain_uid = db.Column(UUID(as_uuid=False), primary_key=True)
  382. role = db.Column(db.String(20), nullable=False)
  383. class OntologyChangeSetModel(db.Model):
  384. __tablename__ = "ontology_change_sets"
  385. __table_args__ = ({"schema": "public"},)
  386. uid = db.Column(UUID(as_uuid=False), primary_key=True, default=new_governance_uid)
  387. ontology_uid = db.Column(
  388. UUID(as_uuid=False), db.ForeignKey("public.ontologies.uid"), nullable=False
  389. )
  390. base_version_uid = db.Column(UUID(as_uuid=False))
  391. status = db.Column(db.String(20), nullable=False, default="draft")
  392. changes = db.Column(JSONB, nullable=False, default=list)
  393. decisions = db.Column(JSONB, nullable=False, default=list)
  394. created_by = db.Column(db.String(100))
  395. created_at = db.Column(db.DateTime, nullable=False, default=now_china_naive)
  396. class OntologyPublishRunModel(db.Model):
  397. __tablename__ = "ontology_publish_runs"
  398. __table_args__ = ({"schema": "public"},)
  399. uid = db.Column(UUID(as_uuid=False), primary_key=True, default=new_governance_uid)
  400. ontology_uid = db.Column(
  401. UUID(as_uuid=False), db.ForeignKey("public.ontologies.uid"), nullable=False
  402. )
  403. version_uid = db.Column(
  404. UUID(as_uuid=False), db.ForeignKey("public.ontology_versions.uid"), nullable=False
  405. )
  406. idempotency_key = db.Column(db.String(128), nullable=False, unique=True)
  407. status = db.Column(db.String(20), nullable=False)
  408. validation_result = db.Column(JSONB, nullable=False, default=dict)
  409. actor_uid = db.Column(db.String(100))
  410. created_at = db.Column(db.DateTime, nullable=False, default=now_china_naive)
  411. finished_at = db.Column(db.DateTime)