data_research.py 8.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241
  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. }