postgres.py 27 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820
  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. FROM acquisition_runs ar
  482. LEFT JOIN acquisition_jobs aj ON aj.run_id = ar.id
  483. LEFT JOIN queries q ON q.id = aj.query_id
  484. LEFT JOIN candidate_items ci ON ci.job_id = aj.id
  485. LEFT JOIN item_classifications ic ON ic.item_id = ci.id
  486. WHERE ar.id = %s
  487. GROUP BY ar.id
  488. """,
  489. (run_id,),
  490. )
  491. summary["queries"] = self._all(
  492. """
  493. SELECT
  494. q.id AS query_id,
  495. q.query_text,
  496. COUNT(DISTINCT aj.id)::int AS job_count,
  497. COUNT(DISTINCT ci.id)::int AS candidate_count,
  498. COUNT(DISTINCT ic.id) FILTER (
  499. WHERE ic.is_creation_knowledge IS TRUE
  500. )::int AS creation_hit_count,
  501. jsonb_object_agg(
  502. aj.platform,
  503. jsonb_build_object(
  504. 'status', aj.status,
  505. 'attempt_count', aj.attempt_count,
  506. 'display_limit', aj.display_limit,
  507. 'search_limit', aj.search_limit,
  508. 'error_message', aj.error_message
  509. )
  510. ) FILTER (WHERE aj.id IS NOT NULL) AS platforms
  511. FROM queries q
  512. JOIN acquisition_jobs aj ON aj.query_id = q.id
  513. LEFT JOIN candidate_items ci ON ci.job_id = aj.id
  514. LEFT JOIN item_classifications ic ON ic.item_id = ci.id
  515. WHERE aj.run_id = %s
  516. GROUP BY q.id, q.query_text, q.sort_order
  517. ORDER BY q.sort_order, q.query_text
  518. """,
  519. (run_id,),
  520. )
  521. return summary
  522. def get_query_detail(self, *, run_id: UUID, query_id: UUID) -> dict[str, Any]:
  523. query = self._one("SELECT * FROM queries WHERE id = %s", (query_id,))
  524. jobs = self._all(
  525. """
  526. SELECT * FROM acquisition_jobs
  527. WHERE run_id = %s AND query_id = %s
  528. ORDER BY platform
  529. """,
  530. (run_id, query_id),
  531. )
  532. items = self._all(
  533. """
  534. SELECT ci.* FROM candidate_items ci
  535. JOIN acquisition_jobs aj ON aj.id = ci.job_id
  536. WHERE aj.run_id = %s AND ci.query_id = %s
  537. ORDER BY ci.platform, ci.created_at
  538. """,
  539. (run_id, query_id),
  540. )
  541. item_ids = [row["id"] for row in items]
  542. media: list[dict[str, Any]] = []
  543. classifications: list[dict[str, Any]] = []
  544. decode_summaries: list[dict[str, Any]] = []
  545. if item_ids:
  546. media = self._all(
  547. """
  548. SELECT * FROM media_assets
  549. WHERE item_id = ANY(%s)
  550. ORDER BY item_id, position
  551. """,
  552. (item_ids,),
  553. )
  554. classifications = self._all(
  555. """
  556. SELECT * FROM item_classifications
  557. WHERE item_id = ANY(%s)
  558. ORDER BY created_at DESC
  559. """,
  560. (item_ids,),
  561. )
  562. decode_summaries = self._all(
  563. """
  564. SELECT DISTINCT ON (dr.item_id)
  565. dr.item_id,
  566. dr.status AS decode_status,
  567. COUNT(DISTINCT pd.id)::int AS payload_count
  568. FROM decode_results dr
  569. LEFT JOIN payload_drafts pd ON pd.item_id = dr.item_id
  570. WHERE dr.item_id = ANY(%s)
  571. GROUP BY dr.item_id, dr.status, dr.created_at
  572. ORDER BY dr.item_id, dr.created_at DESC
  573. """,
  574. (item_ids,),
  575. )
  576. return {
  577. "query": query,
  578. "jobs": jobs,
  579. "items": items,
  580. "media_assets": media,
  581. "classifications": classifications,
  582. "decode_summaries": decode_summaries,
  583. }
  584. def get_query_detail_for_batch(self, *, batch_id: UUID, query_id: UUID) -> dict[str, Any]:
  585. query = self._one(
  586. "SELECT * FROM queries WHERE id = %s AND batch_id = %s",
  587. (query_id, batch_id),
  588. )
  589. jobs = self._all(
  590. """
  591. SELECT aj.* FROM acquisition_jobs aj
  592. JOIN acquisition_runs ar ON ar.id = aj.run_id
  593. WHERE ar.batch_id = %s AND aj.query_id = %s
  594. ORDER BY aj.created_at, aj.platform
  595. """,
  596. (batch_id, query_id),
  597. )
  598. items = self._all(
  599. """
  600. SELECT ci.* FROM candidate_items ci
  601. JOIN acquisition_jobs aj ON aj.id = ci.job_id
  602. JOIN acquisition_runs ar ON ar.id = aj.run_id
  603. WHERE ar.batch_id = %s AND ci.query_id = %s
  604. ORDER BY ci.platform, ci.created_at
  605. """,
  606. (batch_id, query_id),
  607. )
  608. item_ids = [row["id"] for row in items]
  609. media: list[dict[str, Any]] = []
  610. classifications: list[dict[str, Any]] = []
  611. decode_summaries: list[dict[str, Any]] = []
  612. if item_ids:
  613. media = self._all(
  614. """
  615. SELECT * FROM media_assets
  616. WHERE item_id = ANY(%s)
  617. ORDER BY item_id, position
  618. """,
  619. (item_ids,),
  620. )
  621. classifications = self._all(
  622. """
  623. SELECT * FROM item_classifications
  624. WHERE item_id = ANY(%s)
  625. ORDER BY created_at DESC
  626. """,
  627. (item_ids,),
  628. )
  629. decode_summaries = self._all(
  630. """
  631. SELECT DISTINCT ON (dr.item_id)
  632. dr.item_id,
  633. dr.status AS decode_status,
  634. COUNT(DISTINCT pd.id)::int AS payload_count
  635. FROM decode_results dr
  636. LEFT JOIN payload_drafts pd ON pd.item_id = dr.item_id
  637. WHERE dr.item_id = ANY(%s)
  638. GROUP BY dr.item_id, dr.status, dr.created_at
  639. ORDER BY dr.item_id, dr.created_at DESC
  640. """,
  641. (item_ids,),
  642. )
  643. return {
  644. "query": query,
  645. "jobs": jobs,
  646. "items": items,
  647. "media_assets": media,
  648. "classifications": classifications,
  649. "decode_summaries": decode_summaries,
  650. }
  651. def get_latest_singleton_overview(self) -> dict[str, Any]:
  652. batch = self._one_or_none(
  653. """
  654. SELECT qb.* FROM query_batches qb
  655. JOIN acquisition_runs ar ON ar.batch_id = qb.id
  656. GROUP BY qb.id
  657. ORDER BY MAX(ar.created_at) DESC
  658. LIMIT 1
  659. """,
  660. (),
  661. )
  662. if batch is None:
  663. return {"batch": None, "run": None, "queries": [], "decoded_items": []}
  664. run = self._one_or_none(
  665. """
  666. SELECT * FROM acquisition_runs
  667. WHERE batch_id = %s
  668. ORDER BY created_at DESC
  669. LIMIT 1
  670. """,
  671. (batch["id"],),
  672. )
  673. if run is None:
  674. queries = self._all(
  675. """
  676. SELECT
  677. q.id AS query_id,
  678. q.query_text,
  679. q.metadata->>'family_key' AS family_key,
  680. q.sort_order,
  681. 0::int AS candidate_count,
  682. 0::int AS creation_hit_count
  683. FROM queries q
  684. WHERE q.batch_id = %s
  685. ORDER BY q.sort_order, q.created_at
  686. """,
  687. (batch["id"],),
  688. )
  689. return {"batch": batch, "run": None, "queries": queries, "decoded_items": []}
  690. queries = self._all(
  691. """
  692. SELECT
  693. q.id AS query_id,
  694. q.query_text,
  695. q.metadata->>'family_key' AS family_key,
  696. q.sort_order,
  697. COUNT(DISTINCT ci.id)::int AS candidate_count,
  698. COUNT(DISTINCT ic.id) FILTER (
  699. WHERE ic.is_creation_knowledge IS TRUE
  700. )::int AS creation_hit_count,
  701. COUNT(DISTINCT dr.id)::int AS decoded_count,
  702. COUNT(DISTINCT pd.id)::int AS payload_count
  703. FROM queries q
  704. LEFT JOIN acquisition_jobs aj ON aj.query_id = q.id
  705. LEFT JOIN acquisition_runs ar ON ar.id = aj.run_id AND ar.batch_id = q.batch_id
  706. LEFT JOIN candidate_items ci ON ci.job_id = aj.id
  707. LEFT JOIN item_classifications ic ON ic.item_id = ci.id
  708. LEFT JOIN decode_results dr ON dr.item_id = ci.id
  709. LEFT JOIN payload_drafts pd ON pd.item_id = ci.id
  710. WHERE q.batch_id = %s
  711. GROUP BY q.id, q.query_text, q.metadata, q.sort_order
  712. ORDER BY q.sort_order, q.created_at
  713. """,
  714. (batch["id"],),
  715. )
  716. decoded_items = self._all(
  717. """
  718. SELECT
  719. ci.id AS item_id,
  720. ci.query_id,
  721. ci.title,
  722. ci.platform,
  723. dr.status AS decode_status,
  724. COUNT(DISTINCT pd.id)::int AS payload_count
  725. FROM candidate_items ci
  726. JOIN acquisition_jobs aj ON aj.id = ci.job_id
  727. JOIN acquisition_runs ar ON ar.id = aj.run_id
  728. JOIN decode_results dr ON dr.item_id = ci.id
  729. LEFT JOIN payload_drafts pd ON pd.item_id = ci.id
  730. WHERE ar.batch_id = %s
  731. GROUP BY ci.id, ci.query_id, ci.title, ci.platform, dr.status, dr.created_at
  732. ORDER BY dr.created_at DESC
  733. """,
  734. (batch["id"],),
  735. )
  736. return {
  737. "batch": batch,
  738. "run": run,
  739. "queries": queries,
  740. "decoded_items": decoded_items,
  741. }
  742. def list_creation_candidate_items(
  743. self,
  744. *,
  745. run_id: UUID | None = None,
  746. limit: int = 100,
  747. ) -> list[CandidateItem]:
  748. if run_id is None:
  749. rows = self._all(
  750. """
  751. SELECT ci.* FROM candidate_items ci
  752. JOIN item_classifications ic ON ic.item_id = ci.id
  753. WHERE ic.is_creation_knowledge IS TRUE
  754. ORDER BY ci.created_at
  755. LIMIT %s
  756. """,
  757. (limit,),
  758. )
  759. else:
  760. rows = self._all(
  761. """
  762. SELECT ci.* FROM candidate_items ci
  763. JOIN acquisition_jobs aj ON aj.id = ci.job_id
  764. JOIN item_classifications ic ON ic.item_id = ci.id
  765. WHERE aj.run_id = %s AND ic.is_creation_knowledge IS TRUE
  766. ORDER BY ci.created_at
  767. LIMIT %s
  768. """,
  769. (run_id, limit),
  770. )
  771. return [CandidateItem.model_validate(row) for row in rows]
  772. def get_candidate_item(self, item_id: UUID) -> CandidateItem:
  773. row = self._one("SELECT * FROM candidate_items WHERE id = %s", (item_id,))
  774. return CandidateItem.model_validate(row)
  775. def list_media_assets_for_item(self, item_id: UUID) -> list[MediaAsset]:
  776. rows = self._all(
  777. """
  778. SELECT * FROM media_assets
  779. WHERE item_id = %s
  780. ORDER BY position
  781. """,
  782. (item_id,),
  783. )
  784. return [MediaAsset.model_validate(row) for row in rows]