postgres.py 35 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018
  1. """PostgreSQL implementation of the formal acquisition repository."""
  2. from __future__ import annotations
  3. from typing import Any
  4. from uuid import UUID
  5. import psycopg2.extras
  6. from acquisition.domain import (
  7. AcquisitionJob,
  8. AcquisitionRun,
  9. CandidateItem,
  10. ItemClassification,
  11. MediaAsset,
  12. Query,
  13. QueryBatch,
  14. )
  15. Json = psycopg2.extras.Json
  16. psycopg2.extras.register_uuid()
  17. class PostgresAcquisitionRepository:
  18. """Repository backed by the formal cloud PostgreSQL schema.
  19. The repository does not commit by itself; callers own transaction scope via
  20. core.db_session.transaction or an equivalent connection boundary.
  21. """
  22. def __init__(self, conn: Any):
  23. self.conn = conn
  24. def _one(self, sql: str, params: tuple[Any, ...]) -> dict[str, Any]:
  25. with self.conn.cursor(cursor_factory=psycopg2.extras.RealDictCursor) as cur:
  26. cur.execute(sql, params)
  27. row = cur.fetchone()
  28. if row is None:
  29. raise RuntimeError("expected one row, got none")
  30. return dict(row)
  31. def _all(self, sql: str, params: tuple[Any, ...]) -> list[dict[str, Any]]:
  32. with self.conn.cursor(cursor_factory=psycopg2.extras.RealDictCursor) as cur:
  33. cur.execute(sql, params)
  34. return [dict(row) for row in cur.fetchall()]
  35. def create_query_batch(
  36. self,
  37. *,
  38. name: str,
  39. source_type: str = "manual",
  40. generation_method: str | None = None,
  41. target_platforms: list[str] | None = None,
  42. status: str = "draft",
  43. metadata: dict[str, Any] | None = None,
  44. ) -> QueryBatch:
  45. row = self._one(
  46. """
  47. INSERT INTO query_batches(
  48. name, source_type, generation_method, target_platforms,
  49. status, metadata
  50. )
  51. VALUES (%s, %s, %s, %s, %s, %s)
  52. RETURNING *
  53. """,
  54. (
  55. name,
  56. source_type,
  57. generation_method,
  58. target_platforms or [],
  59. status,
  60. Json(metadata or {}),
  61. ),
  62. )
  63. return QueryBatch.model_validate(row)
  64. def add_query(
  65. self,
  66. *,
  67. batch_id: UUID | None,
  68. query_text: str,
  69. axes: dict[str, Any] | None = None,
  70. keep: bool | None = None,
  71. filter_reason: str | None = None,
  72. status: str = "draft",
  73. sort_order: int = 0,
  74. metadata: dict[str, Any] | None = None,
  75. ) -> Query:
  76. row = self._one(
  77. """
  78. INSERT INTO queries(
  79. batch_id, query_text, axes, keep, filter_reason,
  80. status, sort_order, metadata
  81. )
  82. VALUES (%s, %s, %s, %s, %s, %s, %s, %s)
  83. RETURNING *
  84. """,
  85. (
  86. batch_id,
  87. query_text,
  88. Json(axes or {}),
  89. keep,
  90. filter_reason,
  91. status,
  92. sort_order,
  93. Json(metadata or {}),
  94. ),
  95. )
  96. return Query.model_validate(row)
  97. def list_queries_for_batch(
  98. self,
  99. batch_id: UUID,
  100. *,
  101. keep: bool | None = None,
  102. ) -> list[Query]:
  103. if keep is None:
  104. rows = self._all(
  105. """
  106. SELECT * FROM queries
  107. WHERE batch_id = %s
  108. ORDER BY sort_order, created_at
  109. """,
  110. (batch_id,),
  111. )
  112. else:
  113. rows = self._all(
  114. """
  115. SELECT * FROM queries
  116. WHERE batch_id = %s AND keep IS NOT DISTINCT FROM %s
  117. ORDER BY sort_order, created_at
  118. """,
  119. (batch_id, keep),
  120. )
  121. return [Query.model_validate(row) for row in rows]
  122. def get_query_batch(self, batch_id: UUID) -> QueryBatch:
  123. row = self._one("SELECT * FROM query_batches WHERE id = %s", (batch_id,))
  124. return QueryBatch.model_validate(row)
  125. def create_acquisition_run(
  126. self,
  127. *,
  128. batch_id: UUID | None = None,
  129. run_key: str | None = None,
  130. status: str = "pending",
  131. note: str | None = None,
  132. metadata: dict[str, Any] | None = None,
  133. ) -> AcquisitionRun:
  134. row = self._one(
  135. """
  136. INSERT INTO acquisition_runs(batch_id, run_key, status, note, metadata, started_at)
  137. VALUES (%s, %s, %s, %s, %s, CASE WHEN %s = 'running' THEN now() ELSE NULL END)
  138. ON CONFLICT (run_key) DO UPDATE SET
  139. status = EXCLUDED.status,
  140. note = EXCLUDED.note,
  141. metadata = EXCLUDED.metadata,
  142. started_at = COALESCE(acquisition_runs.started_at, EXCLUDED.started_at)
  143. RETURNING *
  144. """,
  145. (batch_id, run_key, status, note, Json(metadata or {}), status),
  146. )
  147. return AcquisitionRun.model_validate(row)
  148. def ensure_acquisition_job(
  149. self,
  150. *,
  151. run_id: UUID,
  152. query_id: UUID | None,
  153. platform: str,
  154. search_limit: int | None = None,
  155. display_limit: int | None = None,
  156. status: str = "pending",
  157. metadata: dict[str, Any] | None = None,
  158. ) -> AcquisitionJob:
  159. row = self._one(
  160. """
  161. INSERT INTO acquisition_jobs(
  162. run_id, query_id, platform, search_limit,
  163. display_limit, status, metadata
  164. )
  165. VALUES (%s, %s, %s, %s, %s, %s, %s)
  166. ON CONFLICT (run_id, query_id, platform) DO UPDATE SET
  167. search_limit = EXCLUDED.search_limit,
  168. display_limit = EXCLUDED.display_limit,
  169. metadata = acquisition_jobs.metadata || EXCLUDED.metadata
  170. RETURNING *
  171. """,
  172. (
  173. run_id,
  174. query_id,
  175. platform,
  176. search_limit,
  177. display_limit,
  178. status,
  179. Json(metadata or {}),
  180. ),
  181. )
  182. return AcquisitionJob.model_validate(row)
  183. def update_acquisition_job(
  184. self,
  185. job_id: UUID,
  186. *,
  187. status: str,
  188. attempt_count: int | None = None,
  189. error_message: str | None = None,
  190. metadata: dict[str, Any] | None = None,
  191. ) -> AcquisitionJob:
  192. row = self._one(
  193. """
  194. UPDATE acquisition_jobs SET
  195. status = %s,
  196. attempt_count = COALESCE(%s, attempt_count),
  197. error_message = %s,
  198. metadata = CASE WHEN %s THEN %s ELSE metadata END
  199. WHERE id = %s
  200. RETURNING *
  201. """,
  202. (
  203. status,
  204. attempt_count,
  205. error_message,
  206. metadata is not None,
  207. Json(metadata or {}),
  208. job_id,
  209. ),
  210. )
  211. return AcquisitionJob.model_validate(row)
  212. def update_acquisition_run(
  213. self,
  214. run_id: UUID,
  215. *,
  216. status: str,
  217. error_message: str | None = None,
  218. metadata: dict[str, Any] | None = None,
  219. ) -> AcquisitionRun:
  220. row = self._one(
  221. """
  222. UPDATE acquisition_runs SET
  223. status = %s,
  224. error_message = %s,
  225. metadata = CASE WHEN %s THEN metadata || %s ELSE metadata END,
  226. finished_at = CASE WHEN %s IN ('done', 'partial', 'failed') THEN now() ELSE finished_at END
  227. WHERE id = %s
  228. RETURNING *
  229. """,
  230. (
  231. status,
  232. error_message,
  233. metadata is not None,
  234. Json(metadata or {}),
  235. status,
  236. run_id,
  237. ),
  238. )
  239. return AcquisitionRun.model_validate(row)
  240. def upsert_candidate_item(
  241. self,
  242. *,
  243. platform: str,
  244. job_id: UUID | None = None,
  245. query_id: UUID | None = None,
  246. platform_item_id: str | None = None,
  247. unique_key: str | None = None,
  248. canonical_url: str | None = None,
  249. content_type: str | None = None,
  250. content_mode: str | None = None,
  251. title: str | None = None,
  252. author_name: str | None = None,
  253. body_text: str | None = None,
  254. raw_summary: str | None = None,
  255. status: str = "candidate",
  256. source_payload: dict[str, Any] | None = None,
  257. metadata: dict[str, Any] | None = None,
  258. error_message: str | None = None,
  259. ) -> CandidateItem:
  260. existing_id = None
  261. if unique_key:
  262. row = self._one_or_none(
  263. """
  264. SELECT id FROM candidate_items
  265. WHERE unique_key = %s
  266. ORDER BY created_at DESC
  267. LIMIT 1
  268. """,
  269. (unique_key,),
  270. )
  271. existing_id = row["id"] if row else None
  272. if platform_item_id:
  273. if existing_id is None:
  274. row = self._one_or_none(
  275. """
  276. SELECT id FROM candidate_items
  277. WHERE platform = %s AND platform_item_id = %s
  278. ORDER BY created_at DESC
  279. LIMIT 1
  280. """,
  281. (platform, platform_item_id),
  282. )
  283. existing_id = row["id"] if row else None
  284. if existing_id:
  285. row = self._one(
  286. """
  287. UPDATE candidate_items SET
  288. job_id = %s,
  289. query_id = %s,
  290. unique_key = COALESCE(%s, unique_key),
  291. canonical_url = %s,
  292. content_type = %s,
  293. content_mode = %s,
  294. title = %s,
  295. author_name = %s,
  296. body_text = %s,
  297. raw_summary = %s,
  298. status = %s,
  299. metadata = %s,
  300. error_message = %s
  301. WHERE id = %s
  302. RETURNING *
  303. """,
  304. (
  305. job_id,
  306. query_id,
  307. unique_key,
  308. canonical_url,
  309. content_type,
  310. content_mode,
  311. title,
  312. author_name,
  313. body_text,
  314. raw_summary,
  315. status,
  316. Json(metadata or {}),
  317. error_message,
  318. existing_id,
  319. ),
  320. )
  321. else:
  322. row = self._one(
  323. """
  324. INSERT INTO candidate_items(
  325. job_id, query_id, platform, platform_item_id, unique_key, canonical_url,
  326. content_type, content_mode, title, author_name, body_text, raw_summary, status,
  327. source_payload, metadata, error_message
  328. )
  329. VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
  330. RETURNING *
  331. """,
  332. (
  333. job_id,
  334. query_id,
  335. platform,
  336. platform_item_id,
  337. unique_key,
  338. canonical_url,
  339. content_type,
  340. content_mode,
  341. title,
  342. author_name,
  343. body_text,
  344. raw_summary,
  345. status,
  346. Json(source_payload or {}),
  347. Json(metadata or {}),
  348. error_message,
  349. ),
  350. )
  351. return CandidateItem.model_validate(row)
  352. def _one_or_none(self, sql: str, params: tuple[Any, ...]) -> dict[str, Any] | None:
  353. with self.conn.cursor(cursor_factory=psycopg2.extras.RealDictCursor) as cur:
  354. cur.execute(sql, params)
  355. row = cur.fetchone()
  356. return dict(row) if row else None
  357. def get_candidate_item_by_unique_key(self, unique_key: str) -> CandidateItem | None:
  358. row = self._one_or_none(
  359. """
  360. SELECT * FROM candidate_items
  361. WHERE unique_key = %s
  362. ORDER BY created_at DESC
  363. LIMIT 1
  364. """,
  365. (unique_key,),
  366. )
  367. return CandidateItem.model_validate(row) if row else None
  368. def attach_existing_candidate_item(
  369. self,
  370. item_id: UUID,
  371. *,
  372. job_id: UUID,
  373. query_id: UUID,
  374. metadata: dict[str, Any] | None = None,
  375. ) -> CandidateItem:
  376. row = self._one(
  377. """
  378. UPDATE candidate_items SET
  379. job_id = %s,
  380. query_id = %s,
  381. metadata = metadata || %s
  382. WHERE id = %s
  383. RETURNING *
  384. """,
  385. (
  386. job_id,
  387. query_id,
  388. Json(metadata or {}),
  389. item_id,
  390. ),
  391. )
  392. return CandidateItem.model_validate(row)
  393. def add_media_asset(
  394. self,
  395. *,
  396. item_id: UUID,
  397. media_type: str,
  398. source_url: str | None = None,
  399. oss_url: str | None = None,
  400. cdn_url: str | None = None,
  401. position: int = 0,
  402. status: str = "pending",
  403. source_payload: dict[str, Any] | None = None,
  404. metadata: dict[str, Any] | None = None,
  405. ) -> MediaAsset:
  406. row = self._one(
  407. """
  408. INSERT INTO media_assets(
  409. item_id, media_type, source_url, oss_url, cdn_url,
  410. position, status, source_payload, metadata
  411. )
  412. VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s)
  413. RETURNING *
  414. """,
  415. (
  416. item_id,
  417. media_type,
  418. source_url,
  419. oss_url,
  420. cdn_url,
  421. position,
  422. status,
  423. Json(source_payload or {}),
  424. Json(metadata or {}),
  425. ),
  426. )
  427. return MediaAsset.model_validate(row)
  428. def add_item_classification(
  429. self,
  430. *,
  431. item_id: UUID,
  432. is_creation_knowledge: bool | None = None,
  433. label: str | None = None,
  434. confidence: float | None = None,
  435. reason: str | None = None,
  436. model_name: str | None = None,
  437. prompt_version: str | None = None,
  438. result_payload: dict[str, Any] | None = None,
  439. status: str = "pending",
  440. error_message: str | None = None,
  441. ) -> ItemClassification:
  442. row = self._one(
  443. """
  444. INSERT INTO item_classifications(
  445. item_id, is_creation_knowledge, label, confidence, reason,
  446. model_name, prompt_version, result_payload, status, error_message
  447. )
  448. VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
  449. RETURNING *
  450. """,
  451. (
  452. item_id,
  453. is_creation_knowledge,
  454. label,
  455. confidence,
  456. reason,
  457. model_name,
  458. prompt_version,
  459. Json(result_payload or {}),
  460. status,
  461. error_message,
  462. ),
  463. )
  464. return ItemClassification.model_validate(row)
  465. def get_run_summary(self, run_id: UUID) -> dict[str, Any]:
  466. summary = self._one(
  467. """
  468. SELECT
  469. ar.id,
  470. ar.run_key,
  471. ar.batch_id,
  472. ar.status,
  473. ar.started_at,
  474. ar.finished_at,
  475. COUNT(DISTINCT q.id)::int AS query_count,
  476. COUNT(DISTINCT aj.id)::int AS job_count,
  477. COUNT(DISTINCT ci.id)::int AS candidate_count,
  478. COUNT(DISTINCT ic.id) FILTER (
  479. WHERE ic.is_creation_knowledge IS TRUE
  480. )::int AS creation_hit_count,
  481. COUNT(DISTINCT ci.id) FILTER (
  482. WHERE ic.is_creation_knowledge IS TRUE
  483. AND dr.status = 'decoded'
  484. AND kp.id IS NOT NULL
  485. )::int AS decoded_count,
  486. COUNT(DISTINCT pd.id) FILTER (
  487. WHERE ic.is_creation_knowledge IS TRUE
  488. AND dr.status = 'decoded'
  489. AND kp.id IS NOT NULL
  490. )::int AS payload_count
  491. FROM acquisition_runs ar
  492. LEFT JOIN acquisition_jobs aj ON aj.run_id = ar.id
  493. LEFT JOIN queries q ON q.id = aj.query_id
  494. LEFT JOIN candidate_items ci ON ci.job_id = aj.id
  495. LEFT JOIN item_classifications ic ON ic.item_id = ci.id
  496. LEFT JOIN decode_results dr ON dr.item_id = ci.id
  497. LEFT JOIN knowledge_particles kp ON kp.item_id = ci.id
  498. LEFT JOIN payload_drafts pd ON pd.item_id = ci.id
  499. WHERE ar.id = %s
  500. GROUP BY ar.id
  501. """,
  502. (run_id,),
  503. )
  504. summary["queries"] = self._all(
  505. """
  506. SELECT
  507. q.id AS query_id,
  508. q.query_text,
  509. COUNT(DISTINCT aj.id)::int AS job_count,
  510. COUNT(DISTINCT ci.id)::int AS candidate_count,
  511. COUNT(DISTINCT ic.id) FILTER (
  512. WHERE ic.is_creation_knowledge IS TRUE
  513. )::int AS creation_hit_count,
  514. COUNT(DISTINCT ci.id) FILTER (
  515. WHERE ic.is_creation_knowledge IS TRUE
  516. AND dr.status = 'decoded'
  517. AND kp.id IS NOT NULL
  518. )::int AS decoded_count,
  519. COUNT(DISTINCT pd.id) FILTER (
  520. WHERE ic.is_creation_knowledge IS TRUE
  521. AND dr.status = 'decoded'
  522. AND kp.id IS NOT NULL
  523. )::int AS payload_count,
  524. jsonb_object_agg(
  525. aj.platform,
  526. jsonb_build_object(
  527. 'status', aj.status,
  528. 'attempt_count', aj.attempt_count,
  529. 'display_limit', aj.display_limit,
  530. 'search_limit', aj.search_limit,
  531. 'error_message', aj.error_message
  532. )
  533. ) FILTER (WHERE aj.id IS NOT NULL) AS platforms
  534. FROM queries q
  535. JOIN acquisition_jobs aj ON aj.query_id = q.id
  536. LEFT JOIN candidate_items ci ON ci.job_id = aj.id
  537. LEFT JOIN item_classifications ic ON ic.item_id = ci.id
  538. LEFT JOIN decode_results dr ON dr.item_id = ci.id
  539. LEFT JOIN knowledge_particles kp ON kp.item_id = ci.id
  540. LEFT JOIN payload_drafts pd ON pd.item_id = ci.id
  541. WHERE aj.run_id = %s
  542. GROUP BY q.id, q.query_text, q.sort_order
  543. ORDER BY q.sort_order, q.query_text
  544. """,
  545. (run_id,),
  546. )
  547. return summary
  548. def get_query_detail(self, *, run_id: UUID, query_id: UUID) -> dict[str, Any]:
  549. query = self._one("SELECT * FROM queries WHERE id = %s", (query_id,))
  550. jobs = self._all(
  551. """
  552. SELECT * FROM acquisition_jobs
  553. WHERE run_id = %s AND query_id = %s
  554. ORDER BY platform
  555. """,
  556. (run_id, query_id),
  557. )
  558. items = self._all(
  559. """
  560. SELECT ci.* FROM candidate_items ci
  561. JOIN acquisition_jobs aj ON aj.id = ci.job_id
  562. WHERE aj.run_id = %s AND ci.query_id = %s
  563. ORDER BY ci.platform, ci.created_at
  564. """,
  565. (run_id, query_id),
  566. )
  567. item_ids = [row["id"] for row in items]
  568. media: list[dict[str, Any]] = []
  569. classifications: list[dict[str, Any]] = []
  570. decode_summaries: list[dict[str, Any]] = []
  571. if item_ids:
  572. media = self._all(
  573. """
  574. SELECT * FROM media_assets
  575. WHERE item_id = ANY(%s)
  576. ORDER BY item_id, position
  577. """,
  578. (item_ids,),
  579. )
  580. classifications = self._all(
  581. """
  582. SELECT * FROM item_classifications
  583. WHERE item_id = ANY(%s)
  584. ORDER BY created_at DESC
  585. """,
  586. (item_ids,),
  587. )
  588. decode_summaries = self._all(
  589. """
  590. SELECT DISTINCT ON (dr.item_id)
  591. dr.item_id,
  592. dr.status AS decode_status,
  593. COUNT(DISTINCT kp.id)::int AS particle_count,
  594. COUNT(DISTINCT pd.id)::int AS payload_count
  595. FROM decode_results dr
  596. LEFT JOIN knowledge_particles kp ON kp.item_id = dr.item_id
  597. LEFT JOIN payload_drafts pd ON pd.item_id = dr.item_id
  598. WHERE dr.item_id = ANY(%s)
  599. GROUP BY dr.item_id, dr.status, dr.created_at
  600. ORDER BY dr.item_id, dr.created_at DESC
  601. """,
  602. (item_ids,),
  603. )
  604. return {
  605. "query": query,
  606. "jobs": jobs,
  607. "items": items,
  608. "media_assets": media,
  609. "classifications": classifications,
  610. "decode_summaries": decode_summaries,
  611. }
  612. def get_query_detail_for_batch(self, *, batch_id: UUID, query_id: UUID) -> dict[str, Any]:
  613. query = self._one(
  614. "SELECT * FROM queries WHERE id = %s AND batch_id = %s",
  615. (query_id, batch_id),
  616. )
  617. jobs = self._all(
  618. """
  619. SELECT aj.* FROM acquisition_jobs aj
  620. JOIN acquisition_runs ar ON ar.id = aj.run_id
  621. WHERE ar.batch_id = %s AND aj.query_id = %s
  622. ORDER BY aj.created_at, aj.platform
  623. """,
  624. (batch_id, query_id),
  625. )
  626. items = self._all(
  627. """
  628. SELECT ci.* FROM candidate_items ci
  629. JOIN acquisition_jobs aj ON aj.id = ci.job_id
  630. JOIN acquisition_runs ar ON ar.id = aj.run_id
  631. WHERE ar.batch_id = %s AND ci.query_id = %s
  632. ORDER BY ci.platform, ci.created_at
  633. """,
  634. (batch_id, query_id),
  635. )
  636. item_ids = [row["id"] for row in items]
  637. media: list[dict[str, Any]] = []
  638. classifications: list[dict[str, Any]] = []
  639. decode_summaries: list[dict[str, Any]] = []
  640. if item_ids:
  641. media = self._all(
  642. """
  643. SELECT * FROM media_assets
  644. WHERE item_id = ANY(%s)
  645. ORDER BY item_id, position
  646. """,
  647. (item_ids,),
  648. )
  649. classifications = self._all(
  650. """
  651. SELECT * FROM item_classifications
  652. WHERE item_id = ANY(%s)
  653. ORDER BY created_at DESC
  654. """,
  655. (item_ids,),
  656. )
  657. decode_summaries = self._all(
  658. """
  659. SELECT DISTINCT ON (dr.item_id)
  660. dr.item_id,
  661. dr.status AS decode_status,
  662. COUNT(DISTINCT kp.id)::int AS particle_count,
  663. COUNT(DISTINCT pd.id)::int AS payload_count
  664. FROM decode_results dr
  665. LEFT JOIN knowledge_particles kp ON kp.item_id = dr.item_id
  666. LEFT JOIN payload_drafts pd ON pd.item_id = dr.item_id
  667. WHERE dr.item_id = ANY(%s)
  668. GROUP BY dr.item_id, dr.status, dr.created_at
  669. ORDER BY dr.item_id, dr.created_at DESC
  670. """,
  671. (item_ids,),
  672. )
  673. return {
  674. "query": query,
  675. "jobs": jobs,
  676. "items": items,
  677. "media_assets": media,
  678. "classifications": classifications,
  679. "decode_summaries": decode_summaries,
  680. }
  681. def get_latest_query_result_list(self, query_id: UUID) -> dict[str, Any]:
  682. query = self._one("SELECT * FROM queries WHERE id = %s", (query_id,))
  683. run = self._one_or_none(
  684. """
  685. SELECT ar.* FROM acquisition_runs ar
  686. JOIN acquisition_jobs aj ON aj.run_id = ar.id
  687. WHERE aj.query_id = %s
  688. ORDER BY aj.created_at DESC, ar.created_at DESC
  689. LIMIT 1
  690. """,
  691. (query_id,),
  692. )
  693. if run is None:
  694. return {
  695. "query": query,
  696. "run": None,
  697. "jobs": [],
  698. "items": [],
  699. "media_assets": [],
  700. "classifications": [],
  701. "decode_summaries": [],
  702. }
  703. jobs = self._all(
  704. """
  705. SELECT aj.* FROM acquisition_jobs aj
  706. WHERE aj.query_id = %s
  707. ORDER BY aj.created_at, aj.platform
  708. """,
  709. (query_id,),
  710. )
  711. items = self._all(
  712. """
  713. SELECT
  714. ci.id,
  715. ci.query_id,
  716. ci.job_id,
  717. ci.platform,
  718. ci.title,
  719. LEFT(ci.raw_summary, 700) AS raw_summary,
  720. ci.status,
  721. ci.content_mode,
  722. ci.metadata,
  723. ci.created_at,
  724. ci.updated_at
  725. FROM candidate_items ci
  726. JOIN acquisition_jobs aj ON aj.id = ci.job_id
  727. WHERE aj.query_id = %s
  728. ORDER BY ci.platform, ci.created_at
  729. """,
  730. (query_id,),
  731. )
  732. item_ids = [row["id"] for row in items]
  733. media: list[dict[str, Any]] = []
  734. classifications: list[dict[str, Any]] = []
  735. decode_summaries: list[dict[str, Any]] = []
  736. if item_ids:
  737. media = self._all(
  738. """
  739. SELECT DISTINCT ON (item_id)
  740. id, item_id, media_type, source_url, oss_url, cdn_url, position, status
  741. FROM media_assets
  742. WHERE item_id = ANY(%s)
  743. ORDER BY
  744. item_id,
  745. CASE WHEN media_type IN ('cover', 'image', 'frame') THEN 0 ELSE 1 END,
  746. position,
  747. created_at
  748. """,
  749. (item_ids,),
  750. )
  751. classifications = self._all(
  752. """
  753. SELECT DISTINCT ON (item_id)
  754. id, item_id, is_creation_knowledge, label, confidence, status, error_message
  755. FROM item_classifications
  756. WHERE item_id = ANY(%s)
  757. ORDER BY item_id, created_at DESC
  758. """,
  759. (item_ids,),
  760. )
  761. decode_summaries = self._all(
  762. """
  763. SELECT DISTINCT ON (dr.item_id)
  764. dr.item_id,
  765. dr.status AS decode_status,
  766. COUNT(DISTINCT kp.id)::int AS particle_count,
  767. COUNT(DISTINCT pd.id)::int AS payload_count
  768. FROM decode_results dr
  769. LEFT JOIN knowledge_particles kp ON kp.item_id = dr.item_id
  770. LEFT JOIN payload_drafts pd ON pd.item_id = dr.item_id
  771. WHERE dr.item_id = ANY(%s)
  772. GROUP BY dr.item_id, dr.status, dr.created_at
  773. ORDER BY dr.item_id, dr.created_at DESC
  774. """,
  775. (item_ids,),
  776. )
  777. return {
  778. "query": query,
  779. "run": run,
  780. "jobs": jobs,
  781. "items": items,
  782. "media_assets": media,
  783. "classifications": classifications,
  784. "decode_summaries": decode_summaries,
  785. }
  786. def get_latest_singleton_overview(self) -> dict[str, Any]:
  787. batch = self._one_or_none(
  788. """
  789. SELECT qb.* FROM query_batches qb
  790. JOIN acquisition_runs ar ON ar.batch_id = qb.id
  791. GROUP BY qb.id
  792. ORDER BY MAX(ar.created_at) DESC
  793. LIMIT 1
  794. """,
  795. (),
  796. )
  797. if batch is None:
  798. return {"batch": None, "run": None, "queries": [], "decoded_items": []}
  799. run = self._one_or_none(
  800. """
  801. SELECT * FROM acquisition_runs
  802. WHERE batch_id = %s
  803. ORDER BY created_at DESC
  804. LIMIT 1
  805. """,
  806. (batch["id"],),
  807. )
  808. if run is None:
  809. queries = self._all(
  810. """
  811. SELECT
  812. q.id AS query_id,
  813. q.query_text,
  814. q.metadata->>'family_key' AS family_key,
  815. q.sort_order,
  816. 0::int AS candidate_count,
  817. 0::int AS creation_hit_count,
  818. 0::int AS decoded_count,
  819. 0::int AS payload_count
  820. FROM queries q
  821. WHERE q.batch_id = %s
  822. ORDER BY q.sort_order, q.created_at
  823. """,
  824. (batch["id"],),
  825. )
  826. return {"batch": batch, "run": None, "queries": queries, "decoded_items": []}
  827. latest_queries_cte = """
  828. WITH query_activity AS (
  829. SELECT
  830. q.id AS query_id,
  831. q.query_text,
  832. q.created_at AS query_created_at,
  833. MAX(ar.created_at) AS activity_at
  834. FROM queries q
  835. LEFT JOIN acquisition_jobs aj ON aj.query_id = q.id
  836. LEFT JOIN acquisition_runs ar ON ar.id = aj.run_id
  837. GROUP BY q.id, q.query_text, q.created_at
  838. ),
  839. query_stats AS (
  840. SELECT
  841. qa.query_id,
  842. qa.query_text,
  843. qa.query_created_at,
  844. qa.activity_at,
  845. COUNT(DISTINCT pd.id) FILTER (
  846. WHERE ic.is_creation_knowledge IS TRUE
  847. AND dr.status = 'decoded'
  848. AND kp.id IS NOT NULL
  849. )::int AS payload_count
  850. FROM query_activity qa
  851. LEFT JOIN acquisition_jobs aj ON aj.query_id = qa.query_id
  852. LEFT JOIN candidate_items ci ON ci.job_id = aj.id
  853. LEFT JOIN item_classifications ic ON ic.item_id = ci.id
  854. LEFT JOIN decode_results dr ON dr.item_id = ci.id
  855. LEFT JOIN knowledge_particles kp ON kp.item_id = ci.id
  856. LEFT JOIN payload_drafts pd ON pd.item_id = ci.id
  857. GROUP BY qa.query_id, qa.query_text, qa.query_created_at, qa.activity_at
  858. ),
  859. latest_queries AS (
  860. SELECT DISTINCT ON (query_text)
  861. query_id
  862. FROM query_stats
  863. ORDER BY query_text, (payload_count > 0) DESC, activity_at DESC NULLS LAST, query_created_at DESC
  864. )
  865. """
  866. queries = self._all(
  867. latest_queries_cte
  868. + """
  869. SELECT
  870. q.id AS query_id,
  871. q.query_text,
  872. q.metadata->>'family_key' AS family_key,
  873. q.sort_order,
  874. COUNT(DISTINCT ci.id)::int AS candidate_count,
  875. COUNT(DISTINCT ic.id) FILTER (
  876. WHERE ic.is_creation_knowledge IS TRUE
  877. )::int AS creation_hit_count,
  878. COUNT(DISTINCT ci.id) FILTER (
  879. WHERE ic.is_creation_knowledge IS TRUE
  880. AND dr.status = 'decoded'
  881. AND kp.id IS NOT NULL
  882. )::int AS decoded_count,
  883. COUNT(DISTINCT pd.id) FILTER (
  884. WHERE ic.is_creation_knowledge IS TRUE
  885. AND dr.status = 'decoded'
  886. AND kp.id IS NOT NULL
  887. )::int AS payload_count
  888. FROM latest_queries lq
  889. JOIN queries q ON q.id = lq.query_id
  890. LEFT JOIN acquisition_jobs aj ON aj.query_id = q.id
  891. LEFT JOIN candidate_items ci ON ci.job_id = aj.id
  892. LEFT JOIN item_classifications ic ON ic.item_id = ci.id
  893. LEFT JOIN decode_results dr ON dr.item_id = ci.id
  894. LEFT JOIN knowledge_particles kp ON kp.item_id = ci.id
  895. LEFT JOIN payload_drafts pd ON pd.item_id = ci.id
  896. GROUP BY q.id, q.query_text, q.metadata, q.sort_order
  897. ORDER BY q.sort_order, q.created_at
  898. """,
  899. (),
  900. )
  901. decoded_items = self._all(
  902. latest_queries_cte
  903. + """
  904. SELECT
  905. ci.id AS item_id,
  906. ci.query_id,
  907. ci.title,
  908. ci.platform,
  909. dr.status AS decode_status,
  910. COUNT(DISTINCT kp.id)::int AS particle_count,
  911. COUNT(DISTINCT pd.id)::int AS payload_count
  912. FROM candidate_items ci
  913. JOIN latest_queries lq ON lq.query_id = ci.query_id
  914. JOIN acquisition_jobs aj ON aj.id = ci.job_id
  915. JOIN decode_results dr ON dr.item_id = ci.id
  916. LEFT JOIN knowledge_particles kp ON kp.item_id = ci.id
  917. LEFT JOIN payload_drafts pd ON pd.item_id = ci.id
  918. GROUP BY ci.id, ci.query_id, ci.title, ci.platform, dr.status, dr.created_at
  919. ORDER BY dr.created_at DESC
  920. """,
  921. (),
  922. )
  923. return {
  924. "batch": batch,
  925. "run": run,
  926. "queries": queries,
  927. "decoded_items": decoded_items,
  928. }
  929. def list_creation_candidate_items(
  930. self,
  931. *,
  932. run_id: UUID | None = None,
  933. limit: int = 100,
  934. ) -> list[CandidateItem]:
  935. if run_id is None:
  936. rows = self._all(
  937. """
  938. SELECT ci.* FROM candidate_items ci
  939. JOIN item_classifications ic ON ic.item_id = ci.id
  940. WHERE ic.is_creation_knowledge IS TRUE
  941. AND NOT EXISTS (
  942. SELECT 1 FROM decode_results dr
  943. WHERE dr.item_id = ci.id
  944. AND dr.status IN ('decoded', 'rejected', 'skipped')
  945. )
  946. ORDER BY ci.created_at
  947. LIMIT %s
  948. """,
  949. (limit,),
  950. )
  951. else:
  952. rows = self._all(
  953. """
  954. SELECT ci.* FROM candidate_items ci
  955. JOIN acquisition_jobs aj ON aj.id = ci.job_id
  956. JOIN item_classifications ic ON ic.item_id = ci.id
  957. WHERE aj.run_id = %s AND ic.is_creation_knowledge IS TRUE
  958. AND NOT EXISTS (
  959. SELECT 1 FROM decode_results dr
  960. WHERE dr.item_id = ci.id
  961. AND dr.status IN ('decoded', 'rejected', 'skipped')
  962. )
  963. ORDER BY ci.created_at
  964. LIMIT %s
  965. """,
  966. (run_id, limit),
  967. )
  968. return [CandidateItem.model_validate(row) for row in rows]
  969. def get_candidate_item(self, item_id: UUID) -> CandidateItem:
  970. row = self._one("SELECT * FROM candidate_items WHERE id = %s", (item_id,))
  971. return CandidateItem.model_validate(row)
  972. def list_media_assets_for_item(self, item_id: UUID) -> list[MediaAsset]:
  973. rows = self._all(
  974. """
  975. SELECT * FROM media_assets
  976. WHERE item_id = %s
  977. ORDER BY position
  978. """,
  979. (item_id,),
  980. )
  981. return [MediaAsset.model_validate(row) for row in rows]