postgres.py 36 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044
  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
  706. aj.id,
  707. aj.run_id,
  708. aj.query_id,
  709. aj.platform,
  710. aj.status,
  711. aj.attempt_count,
  712. aj.display_limit,
  713. aj.search_limit,
  714. aj.error_message,
  715. aj.started_at,
  716. aj.finished_at,
  717. aj.created_at,
  718. aj.updated_at
  719. FROM acquisition_jobs aj
  720. WHERE aj.query_id = %s
  721. ORDER BY aj.created_at, aj.platform
  722. """,
  723. (query_id,),
  724. )
  725. items = self._all(
  726. """
  727. SELECT
  728. ci.id,
  729. ci.query_id,
  730. ci.job_id,
  731. ci.platform,
  732. ci.title,
  733. LEFT(ci.raw_summary, 700) AS raw_summary,
  734. ci.status,
  735. ci.content_mode,
  736. jsonb_strip_nulls(jsonb_build_object(
  737. 'page_index', ci.metadata -> 'page_index',
  738. 'page_rank', ci.metadata -> 'page_rank',
  739. 'source_cursor', ci.metadata -> 'source_cursor',
  740. 'content_mode', ci.metadata -> 'content_mode',
  741. 'search_provider', ci.metadata -> 'search_provider',
  742. 'detail_provider', ci.metadata -> 'detail_provider',
  743. 'acquisition_match_status', ci.metadata -> 'acquisition_match_status',
  744. 'matched_unique_key', ci.metadata -> 'matched_unique_key',
  745. 'skip_reason', ci.metadata -> 'skip_reason',
  746. 'unsupported_raw_type', ci.metadata -> 'unsupported_raw_type',
  747. 'video_url_missing', ci.metadata -> 'video_url_missing'
  748. )) AS metadata,
  749. ci.created_at,
  750. ci.updated_at
  751. FROM candidate_items ci
  752. JOIN acquisition_jobs aj ON aj.id = ci.job_id
  753. WHERE aj.query_id = %s
  754. ORDER BY ci.platform, ci.created_at
  755. """,
  756. (query_id,),
  757. )
  758. item_ids = [row["id"] for row in items]
  759. media: list[dict[str, Any]] = []
  760. classifications: list[dict[str, Any]] = []
  761. decode_summaries: list[dict[str, Any]] = []
  762. if item_ids:
  763. media = self._all(
  764. """
  765. SELECT DISTINCT ON (item_id)
  766. id, item_id, media_type, source_url, oss_url, cdn_url, position, status
  767. FROM media_assets
  768. WHERE item_id = ANY(%s)
  769. ORDER BY
  770. item_id,
  771. CASE WHEN media_type IN ('cover', 'image', 'frame') THEN 0 ELSE 1 END,
  772. position,
  773. created_at
  774. """,
  775. (item_ids,),
  776. )
  777. classifications = self._all(
  778. """
  779. SELECT DISTINCT ON (item_id)
  780. id, item_id, is_creation_knowledge, label, confidence, status, error_message
  781. FROM item_classifications
  782. WHERE item_id = ANY(%s)
  783. ORDER BY item_id, created_at DESC
  784. """,
  785. (item_ids,),
  786. )
  787. decode_summaries = self._all(
  788. """
  789. SELECT DISTINCT ON (dr.item_id)
  790. dr.item_id,
  791. dr.status AS decode_status,
  792. COUNT(DISTINCT kp.id)::int AS particle_count,
  793. COUNT(DISTINCT pd.id)::int AS payload_count
  794. FROM decode_results dr
  795. LEFT JOIN knowledge_particles kp ON kp.item_id = dr.item_id
  796. LEFT JOIN payload_drafts pd ON pd.item_id = dr.item_id
  797. WHERE dr.item_id = ANY(%s)
  798. GROUP BY dr.item_id, dr.status, dr.created_at
  799. ORDER BY dr.item_id, dr.created_at DESC
  800. """,
  801. (item_ids,),
  802. )
  803. return {
  804. "query": query,
  805. "run": run,
  806. "jobs": jobs,
  807. "items": items,
  808. "media_assets": media,
  809. "classifications": classifications,
  810. "decode_summaries": decode_summaries,
  811. }
  812. def get_latest_singleton_overview(self) -> dict[str, Any]:
  813. batch = self._one_or_none(
  814. """
  815. SELECT qb.* FROM query_batches qb
  816. JOIN acquisition_runs ar ON ar.batch_id = qb.id
  817. GROUP BY qb.id
  818. ORDER BY MAX(ar.created_at) DESC
  819. LIMIT 1
  820. """,
  821. (),
  822. )
  823. if batch is None:
  824. return {"batch": None, "run": None, "queries": [], "decoded_items": []}
  825. run = self._one_or_none(
  826. """
  827. SELECT * FROM acquisition_runs
  828. WHERE batch_id = %s
  829. ORDER BY created_at DESC
  830. LIMIT 1
  831. """,
  832. (batch["id"],),
  833. )
  834. if run is None:
  835. queries = self._all(
  836. """
  837. SELECT
  838. q.id AS query_id,
  839. q.query_text,
  840. q.metadata->>'family_key' AS family_key,
  841. q.sort_order,
  842. 0::int AS candidate_count,
  843. 0::int AS creation_hit_count,
  844. 0::int AS decoded_count,
  845. 0::int AS payload_count
  846. FROM queries q
  847. WHERE q.batch_id = %s
  848. ORDER BY q.sort_order, q.created_at
  849. """,
  850. (batch["id"],),
  851. )
  852. return {"batch": batch, "run": None, "queries": queries, "decoded_items": []}
  853. latest_queries_cte = """
  854. WITH query_activity AS (
  855. SELECT
  856. q.id AS query_id,
  857. q.query_text,
  858. q.created_at AS query_created_at,
  859. MAX(ar.created_at) AS activity_at
  860. FROM queries q
  861. LEFT JOIN acquisition_jobs aj ON aj.query_id = q.id
  862. LEFT JOIN acquisition_runs ar ON ar.id = aj.run_id
  863. GROUP BY q.id, q.query_text, q.created_at
  864. ),
  865. query_stats AS (
  866. SELECT
  867. qa.query_id,
  868. qa.query_text,
  869. qa.query_created_at,
  870. qa.activity_at,
  871. COUNT(DISTINCT pd.id) FILTER (
  872. WHERE ic.is_creation_knowledge IS TRUE
  873. AND dr.status = 'decoded'
  874. AND kp.id IS NOT NULL
  875. )::int AS payload_count
  876. FROM query_activity qa
  877. LEFT JOIN acquisition_jobs aj ON aj.query_id = qa.query_id
  878. LEFT JOIN candidate_items ci ON ci.job_id = aj.id
  879. LEFT JOIN item_classifications ic ON ic.item_id = ci.id
  880. LEFT JOIN decode_results dr ON dr.item_id = ci.id
  881. LEFT JOIN knowledge_particles kp ON kp.item_id = ci.id
  882. LEFT JOIN payload_drafts pd ON pd.item_id = ci.id
  883. GROUP BY qa.query_id, qa.query_text, qa.query_created_at, qa.activity_at
  884. ),
  885. latest_queries AS (
  886. SELECT DISTINCT ON (query_text)
  887. query_id
  888. FROM query_stats
  889. ORDER BY query_text, (payload_count > 0) DESC, activity_at DESC NULLS LAST, query_created_at DESC
  890. )
  891. """
  892. queries = self._all(
  893. latest_queries_cte
  894. + """
  895. SELECT
  896. q.id AS query_id,
  897. q.query_text,
  898. q.metadata->>'family_key' AS family_key,
  899. q.sort_order,
  900. COUNT(DISTINCT ci.id)::int AS candidate_count,
  901. COUNT(DISTINCT ic.id) FILTER (
  902. WHERE ic.is_creation_knowledge IS TRUE
  903. )::int AS creation_hit_count,
  904. COUNT(DISTINCT ci.id) FILTER (
  905. WHERE ic.is_creation_knowledge IS TRUE
  906. AND dr.status = 'decoded'
  907. AND kp.id IS NOT NULL
  908. )::int AS decoded_count,
  909. COUNT(DISTINCT pd.id) FILTER (
  910. WHERE ic.is_creation_knowledge IS TRUE
  911. AND dr.status = 'decoded'
  912. AND kp.id IS NOT NULL
  913. )::int AS payload_count
  914. FROM latest_queries lq
  915. JOIN queries q ON q.id = lq.query_id
  916. LEFT JOIN acquisition_jobs aj ON aj.query_id = q.id
  917. LEFT JOIN candidate_items ci ON ci.job_id = aj.id
  918. LEFT JOIN item_classifications ic ON ic.item_id = ci.id
  919. LEFT JOIN decode_results dr ON dr.item_id = ci.id
  920. LEFT JOIN knowledge_particles kp ON kp.item_id = ci.id
  921. LEFT JOIN payload_drafts pd ON pd.item_id = ci.id
  922. GROUP BY q.id, q.query_text, q.metadata, q.sort_order
  923. ORDER BY q.sort_order, q.created_at
  924. """,
  925. (),
  926. )
  927. decoded_items = self._all(
  928. latest_queries_cte
  929. + """
  930. SELECT
  931. ci.id AS item_id,
  932. ci.query_id,
  933. ci.title,
  934. ci.platform,
  935. dr.status AS decode_status,
  936. COUNT(DISTINCT kp.id)::int AS particle_count,
  937. COUNT(DISTINCT pd.id)::int AS payload_count
  938. FROM candidate_items ci
  939. JOIN latest_queries lq ON lq.query_id = ci.query_id
  940. JOIN acquisition_jobs aj ON aj.id = ci.job_id
  941. JOIN decode_results dr ON dr.item_id = ci.id
  942. LEFT JOIN knowledge_particles kp ON kp.item_id = ci.id
  943. LEFT JOIN payload_drafts pd ON pd.item_id = ci.id
  944. GROUP BY ci.id, ci.query_id, ci.title, ci.platform, dr.status, dr.created_at
  945. ORDER BY dr.created_at DESC
  946. """,
  947. (),
  948. )
  949. return {
  950. "batch": batch,
  951. "run": run,
  952. "queries": queries,
  953. "decoded_items": decoded_items,
  954. }
  955. def list_creation_candidate_items(
  956. self,
  957. *,
  958. run_id: UUID | None = None,
  959. limit: int = 100,
  960. ) -> list[CandidateItem]:
  961. if run_id is None:
  962. rows = self._all(
  963. """
  964. SELECT ci.* FROM candidate_items ci
  965. JOIN item_classifications ic ON ic.item_id = ci.id
  966. WHERE ic.is_creation_knowledge IS TRUE
  967. AND NOT EXISTS (
  968. SELECT 1 FROM decode_results dr
  969. WHERE dr.item_id = ci.id
  970. AND dr.status IN ('decoded', 'rejected', 'skipped')
  971. )
  972. ORDER BY ci.created_at
  973. LIMIT %s
  974. """,
  975. (limit,),
  976. )
  977. else:
  978. rows = self._all(
  979. """
  980. SELECT ci.* FROM candidate_items ci
  981. JOIN acquisition_jobs aj ON aj.id = ci.job_id
  982. JOIN item_classifications ic ON ic.item_id = ci.id
  983. WHERE aj.run_id = %s AND ic.is_creation_knowledge IS TRUE
  984. AND NOT EXISTS (
  985. SELECT 1 FROM decode_results dr
  986. WHERE dr.item_id = ci.id
  987. AND dr.status IN ('decoded', 'rejected', 'skipped')
  988. )
  989. ORDER BY ci.created_at
  990. LIMIT %s
  991. """,
  992. (run_id, limit),
  993. )
  994. return [CandidateItem.model_validate(row) for row in rows]
  995. def get_candidate_item(self, item_id: UUID) -> CandidateItem:
  996. row = self._one("SELECT * FROM candidate_items WHERE id = %s", (item_id,))
  997. return CandidateItem.model_validate(row)
  998. def list_media_assets_for_item(self, item_id: UUID) -> list[MediaAsset]:
  999. rows = self._all(
  1000. """
  1001. SELECT * FROM media_assets
  1002. WHERE item_id = %s
  1003. ORDER BY position
  1004. """,
  1005. (item_id,),
  1006. )
  1007. return [MediaAsset.model_validate(row) for row in rows]