| 1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018 |
- """PostgreSQL implementation of the formal acquisition repository."""
- from __future__ import annotations
- from typing import Any
- from uuid import UUID
- import psycopg2.extras
- from acquisition.domain import (
- AcquisitionJob,
- AcquisitionRun,
- CandidateItem,
- ItemClassification,
- MediaAsset,
- Query,
- QueryBatch,
- )
- Json = psycopg2.extras.Json
- psycopg2.extras.register_uuid()
- class PostgresAcquisitionRepository:
- """Repository backed by the formal cloud PostgreSQL schema.
- The repository does not commit by itself; callers own transaction scope via
- core.db_session.transaction or an equivalent connection boundary.
- """
- def __init__(self, conn: Any):
- self.conn = conn
- def _one(self, sql: str, params: tuple[Any, ...]) -> dict[str, Any]:
- with self.conn.cursor(cursor_factory=psycopg2.extras.RealDictCursor) as cur:
- cur.execute(sql, params)
- row = cur.fetchone()
- if row is None:
- raise RuntimeError("expected one row, got none")
- return dict(row)
- def _all(self, sql: str, params: tuple[Any, ...]) -> list[dict[str, Any]]:
- with self.conn.cursor(cursor_factory=psycopg2.extras.RealDictCursor) as cur:
- cur.execute(sql, params)
- return [dict(row) for row in cur.fetchall()]
- def create_query_batch(
- self,
- *,
- name: str,
- source_type: str = "manual",
- generation_method: str | None = None,
- target_platforms: list[str] | None = None,
- status: str = "draft",
- metadata: dict[str, Any] | None = None,
- ) -> QueryBatch:
- row = self._one(
- """
- INSERT INTO query_batches(
- name, source_type, generation_method, target_platforms,
- status, metadata
- )
- VALUES (%s, %s, %s, %s, %s, %s)
- RETURNING *
- """,
- (
- name,
- source_type,
- generation_method,
- target_platforms or [],
- status,
- Json(metadata or {}),
- ),
- )
- return QueryBatch.model_validate(row)
- def add_query(
- self,
- *,
- batch_id: UUID | None,
- query_text: str,
- axes: dict[str, Any] | None = None,
- keep: bool | None = None,
- filter_reason: str | None = None,
- status: str = "draft",
- sort_order: int = 0,
- metadata: dict[str, Any] | None = None,
- ) -> Query:
- row = self._one(
- """
- INSERT INTO queries(
- batch_id, query_text, axes, keep, filter_reason,
- status, sort_order, metadata
- )
- VALUES (%s, %s, %s, %s, %s, %s, %s, %s)
- RETURNING *
- """,
- (
- batch_id,
- query_text,
- Json(axes or {}),
- keep,
- filter_reason,
- status,
- sort_order,
- Json(metadata or {}),
- ),
- )
- return Query.model_validate(row)
- def list_queries_for_batch(
- self,
- batch_id: UUID,
- *,
- keep: bool | None = None,
- ) -> list[Query]:
- if keep is None:
- rows = self._all(
- """
- SELECT * FROM queries
- WHERE batch_id = %s
- ORDER BY sort_order, created_at
- """,
- (batch_id,),
- )
- else:
- rows = self._all(
- """
- SELECT * FROM queries
- WHERE batch_id = %s AND keep IS NOT DISTINCT FROM %s
- ORDER BY sort_order, created_at
- """,
- (batch_id, keep),
- )
- return [Query.model_validate(row) for row in rows]
- def get_query_batch(self, batch_id: UUID) -> QueryBatch:
- row = self._one("SELECT * FROM query_batches WHERE id = %s", (batch_id,))
- return QueryBatch.model_validate(row)
- def create_acquisition_run(
- self,
- *,
- batch_id: UUID | None = None,
- run_key: str | None = None,
- status: str = "pending",
- note: str | None = None,
- metadata: dict[str, Any] | None = None,
- ) -> AcquisitionRun:
- row = self._one(
- """
- INSERT INTO acquisition_runs(batch_id, run_key, status, note, metadata, started_at)
- VALUES (%s, %s, %s, %s, %s, CASE WHEN %s = 'running' THEN now() ELSE NULL END)
- ON CONFLICT (run_key) DO UPDATE SET
- status = EXCLUDED.status,
- note = EXCLUDED.note,
- metadata = EXCLUDED.metadata,
- started_at = COALESCE(acquisition_runs.started_at, EXCLUDED.started_at)
- RETURNING *
- """,
- (batch_id, run_key, status, note, Json(metadata or {}), status),
- )
- return AcquisitionRun.model_validate(row)
- def ensure_acquisition_job(
- self,
- *,
- run_id: UUID,
- query_id: UUID | None,
- platform: str,
- search_limit: int | None = None,
- display_limit: int | None = None,
- status: str = "pending",
- metadata: dict[str, Any] | None = None,
- ) -> AcquisitionJob:
- row = self._one(
- """
- INSERT INTO acquisition_jobs(
- run_id, query_id, platform, search_limit,
- display_limit, status, metadata
- )
- VALUES (%s, %s, %s, %s, %s, %s, %s)
- ON CONFLICT (run_id, query_id, platform) DO UPDATE SET
- search_limit = EXCLUDED.search_limit,
- display_limit = EXCLUDED.display_limit,
- metadata = acquisition_jobs.metadata || EXCLUDED.metadata
- RETURNING *
- """,
- (
- run_id,
- query_id,
- platform,
- search_limit,
- display_limit,
- status,
- Json(metadata or {}),
- ),
- )
- return AcquisitionJob.model_validate(row)
- def update_acquisition_job(
- self,
- job_id: UUID,
- *,
- status: str,
- attempt_count: int | None = None,
- error_message: str | None = None,
- metadata: dict[str, Any] | None = None,
- ) -> AcquisitionJob:
- row = self._one(
- """
- UPDATE acquisition_jobs SET
- status = %s,
- attempt_count = COALESCE(%s, attempt_count),
- error_message = %s,
- metadata = CASE WHEN %s THEN %s ELSE metadata END
- WHERE id = %s
- RETURNING *
- """,
- (
- status,
- attempt_count,
- error_message,
- metadata is not None,
- Json(metadata or {}),
- job_id,
- ),
- )
- return AcquisitionJob.model_validate(row)
- def update_acquisition_run(
- self,
- run_id: UUID,
- *,
- status: str,
- error_message: str | None = None,
- metadata: dict[str, Any] | None = None,
- ) -> AcquisitionRun:
- row = self._one(
- """
- UPDATE acquisition_runs SET
- status = %s,
- error_message = %s,
- metadata = CASE WHEN %s THEN metadata || %s ELSE metadata END,
- finished_at = CASE WHEN %s IN ('done', 'partial', 'failed') THEN now() ELSE finished_at END
- WHERE id = %s
- RETURNING *
- """,
- (
- status,
- error_message,
- metadata is not None,
- Json(metadata or {}),
- status,
- run_id,
- ),
- )
- return AcquisitionRun.model_validate(row)
- def upsert_candidate_item(
- self,
- *,
- platform: str,
- job_id: UUID | None = None,
- query_id: UUID | None = None,
- platform_item_id: str | None = None,
- unique_key: str | None = None,
- canonical_url: str | None = None,
- content_type: str | None = None,
- content_mode: str | None = None,
- title: str | None = None,
- author_name: str | None = None,
- body_text: str | None = None,
- raw_summary: str | None = None,
- status: str = "candidate",
- source_payload: dict[str, Any] | None = None,
- metadata: dict[str, Any] | None = None,
- error_message: str | None = None,
- ) -> CandidateItem:
- existing_id = None
- if unique_key:
- row = self._one_or_none(
- """
- SELECT id FROM candidate_items
- WHERE unique_key = %s
- ORDER BY created_at DESC
- LIMIT 1
- """,
- (unique_key,),
- )
- existing_id = row["id"] if row else None
- if platform_item_id:
- if existing_id is None:
- row = self._one_or_none(
- """
- SELECT id FROM candidate_items
- WHERE platform = %s AND platform_item_id = %s
- ORDER BY created_at DESC
- LIMIT 1
- """,
- (platform, platform_item_id),
- )
- existing_id = row["id"] if row else None
- if existing_id:
- row = self._one(
- """
- UPDATE candidate_items SET
- job_id = %s,
- query_id = %s,
- unique_key = COALESCE(%s, unique_key),
- canonical_url = %s,
- content_type = %s,
- content_mode = %s,
- title = %s,
- author_name = %s,
- body_text = %s,
- raw_summary = %s,
- status = %s,
- metadata = %s,
- error_message = %s
- WHERE id = %s
- RETURNING *
- """,
- (
- job_id,
- query_id,
- unique_key,
- canonical_url,
- content_type,
- content_mode,
- title,
- author_name,
- body_text,
- raw_summary,
- status,
- Json(metadata or {}),
- error_message,
- existing_id,
- ),
- )
- else:
- row = self._one(
- """
- INSERT INTO candidate_items(
- job_id, query_id, platform, platform_item_id, unique_key, canonical_url,
- content_type, content_mode, title, author_name, body_text, raw_summary, status,
- source_payload, metadata, error_message
- )
- VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
- RETURNING *
- """,
- (
- job_id,
- query_id,
- platform,
- platform_item_id,
- unique_key,
- canonical_url,
- content_type,
- content_mode,
- title,
- author_name,
- body_text,
- raw_summary,
- status,
- Json(source_payload or {}),
- Json(metadata or {}),
- error_message,
- ),
- )
- return CandidateItem.model_validate(row)
- def _one_or_none(self, sql: str, params: tuple[Any, ...]) -> dict[str, Any] | None:
- with self.conn.cursor(cursor_factory=psycopg2.extras.RealDictCursor) as cur:
- cur.execute(sql, params)
- row = cur.fetchone()
- return dict(row) if row else None
- def get_candidate_item_by_unique_key(self, unique_key: str) -> CandidateItem | None:
- row = self._one_or_none(
- """
- SELECT * FROM candidate_items
- WHERE unique_key = %s
- ORDER BY created_at DESC
- LIMIT 1
- """,
- (unique_key,),
- )
- return CandidateItem.model_validate(row) if row else None
- def attach_existing_candidate_item(
- self,
- item_id: UUID,
- *,
- job_id: UUID,
- query_id: UUID,
- metadata: dict[str, Any] | None = None,
- ) -> CandidateItem:
- row = self._one(
- """
- UPDATE candidate_items SET
- job_id = %s,
- query_id = %s,
- metadata = metadata || %s
- WHERE id = %s
- RETURNING *
- """,
- (
- job_id,
- query_id,
- Json(metadata or {}),
- item_id,
- ),
- )
- return CandidateItem.model_validate(row)
- def add_media_asset(
- self,
- *,
- item_id: UUID,
- media_type: str,
- source_url: str | None = None,
- oss_url: str | None = None,
- cdn_url: str | None = None,
- position: int = 0,
- status: str = "pending",
- source_payload: dict[str, Any] | None = None,
- metadata: dict[str, Any] | None = None,
- ) -> MediaAsset:
- row = self._one(
- """
- INSERT INTO media_assets(
- item_id, media_type, source_url, oss_url, cdn_url,
- position, status, source_payload, metadata
- )
- VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s)
- RETURNING *
- """,
- (
- item_id,
- media_type,
- source_url,
- oss_url,
- cdn_url,
- position,
- status,
- Json(source_payload or {}),
- Json(metadata or {}),
- ),
- )
- return MediaAsset.model_validate(row)
- def add_item_classification(
- self,
- *,
- item_id: UUID,
- is_creation_knowledge: bool | None = None,
- label: str | None = None,
- confidence: float | None = None,
- reason: str | None = None,
- model_name: str | None = None,
- prompt_version: str | None = None,
- result_payload: dict[str, Any] | None = None,
- status: str = "pending",
- error_message: str | None = None,
- ) -> ItemClassification:
- row = self._one(
- """
- INSERT INTO item_classifications(
- item_id, is_creation_knowledge, label, confidence, reason,
- model_name, prompt_version, result_payload, status, error_message
- )
- VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
- RETURNING *
- """,
- (
- item_id,
- is_creation_knowledge,
- label,
- confidence,
- reason,
- model_name,
- prompt_version,
- Json(result_payload or {}),
- status,
- error_message,
- ),
- )
- return ItemClassification.model_validate(row)
- def get_run_summary(self, run_id: UUID) -> dict[str, Any]:
- summary = self._one(
- """
- SELECT
- ar.id,
- ar.run_key,
- ar.batch_id,
- ar.status,
- ar.started_at,
- ar.finished_at,
- COUNT(DISTINCT q.id)::int AS query_count,
- COUNT(DISTINCT aj.id)::int AS job_count,
- COUNT(DISTINCT ci.id)::int AS candidate_count,
- COUNT(DISTINCT ic.id) FILTER (
- WHERE ic.is_creation_knowledge IS TRUE
- )::int AS creation_hit_count,
- COUNT(DISTINCT ci.id) FILTER (
- WHERE ic.is_creation_knowledge IS TRUE
- AND dr.status = 'decoded'
- AND kp.id IS NOT NULL
- )::int AS decoded_count,
- COUNT(DISTINCT pd.id) FILTER (
- WHERE ic.is_creation_knowledge IS TRUE
- AND dr.status = 'decoded'
- AND kp.id IS NOT NULL
- )::int AS payload_count
- FROM acquisition_runs ar
- LEFT JOIN acquisition_jobs aj ON aj.run_id = ar.id
- LEFT JOIN queries q ON q.id = aj.query_id
- LEFT JOIN candidate_items ci ON ci.job_id = aj.id
- LEFT JOIN item_classifications ic ON ic.item_id = ci.id
- LEFT JOIN decode_results dr ON dr.item_id = ci.id
- LEFT JOIN knowledge_particles kp ON kp.item_id = ci.id
- LEFT JOIN payload_drafts pd ON pd.item_id = ci.id
- WHERE ar.id = %s
- GROUP BY ar.id
- """,
- (run_id,),
- )
- summary["queries"] = self._all(
- """
- SELECT
- q.id AS query_id,
- q.query_text,
- COUNT(DISTINCT aj.id)::int AS job_count,
- COUNT(DISTINCT ci.id)::int AS candidate_count,
- COUNT(DISTINCT ic.id) FILTER (
- WHERE ic.is_creation_knowledge IS TRUE
- )::int AS creation_hit_count,
- COUNT(DISTINCT ci.id) FILTER (
- WHERE ic.is_creation_knowledge IS TRUE
- AND dr.status = 'decoded'
- AND kp.id IS NOT NULL
- )::int AS decoded_count,
- COUNT(DISTINCT pd.id) FILTER (
- WHERE ic.is_creation_knowledge IS TRUE
- AND dr.status = 'decoded'
- AND kp.id IS NOT NULL
- )::int AS payload_count,
- jsonb_object_agg(
- aj.platform,
- jsonb_build_object(
- 'status', aj.status,
- 'attempt_count', aj.attempt_count,
- 'display_limit', aj.display_limit,
- 'search_limit', aj.search_limit,
- 'error_message', aj.error_message
- )
- ) FILTER (WHERE aj.id IS NOT NULL) AS platforms
- FROM queries q
- JOIN acquisition_jobs aj ON aj.query_id = q.id
- LEFT JOIN candidate_items ci ON ci.job_id = aj.id
- LEFT JOIN item_classifications ic ON ic.item_id = ci.id
- LEFT JOIN decode_results dr ON dr.item_id = ci.id
- LEFT JOIN knowledge_particles kp ON kp.item_id = ci.id
- LEFT JOIN payload_drafts pd ON pd.item_id = ci.id
- WHERE aj.run_id = %s
- GROUP BY q.id, q.query_text, q.sort_order
- ORDER BY q.sort_order, q.query_text
- """,
- (run_id,),
- )
- return summary
- def get_query_detail(self, *, run_id: UUID, query_id: UUID) -> dict[str, Any]:
- query = self._one("SELECT * FROM queries WHERE id = %s", (query_id,))
- jobs = self._all(
- """
- SELECT * FROM acquisition_jobs
- WHERE run_id = %s AND query_id = %s
- ORDER BY platform
- """,
- (run_id, query_id),
- )
- items = self._all(
- """
- SELECT ci.* FROM candidate_items ci
- JOIN acquisition_jobs aj ON aj.id = ci.job_id
- WHERE aj.run_id = %s AND ci.query_id = %s
- ORDER BY ci.platform, ci.created_at
- """,
- (run_id, query_id),
- )
- item_ids = [row["id"] for row in items]
- media: list[dict[str, Any]] = []
- classifications: list[dict[str, Any]] = []
- decode_summaries: list[dict[str, Any]] = []
- if item_ids:
- media = self._all(
- """
- SELECT * FROM media_assets
- WHERE item_id = ANY(%s)
- ORDER BY item_id, position
- """,
- (item_ids,),
- )
- classifications = self._all(
- """
- SELECT * FROM item_classifications
- WHERE item_id = ANY(%s)
- ORDER BY created_at DESC
- """,
- (item_ids,),
- )
- decode_summaries = self._all(
- """
- SELECT DISTINCT ON (dr.item_id)
- dr.item_id,
- dr.status AS decode_status,
- COUNT(DISTINCT kp.id)::int AS particle_count,
- COUNT(DISTINCT pd.id)::int AS payload_count
- FROM decode_results dr
- LEFT JOIN knowledge_particles kp ON kp.item_id = dr.item_id
- LEFT JOIN payload_drafts pd ON pd.item_id = dr.item_id
- WHERE dr.item_id = ANY(%s)
- GROUP BY dr.item_id, dr.status, dr.created_at
- ORDER BY dr.item_id, dr.created_at DESC
- """,
- (item_ids,),
- )
- return {
- "query": query,
- "jobs": jobs,
- "items": items,
- "media_assets": media,
- "classifications": classifications,
- "decode_summaries": decode_summaries,
- }
- def get_query_detail_for_batch(self, *, batch_id: UUID, query_id: UUID) -> dict[str, Any]:
- query = self._one(
- "SELECT * FROM queries WHERE id = %s AND batch_id = %s",
- (query_id, batch_id),
- )
- jobs = self._all(
- """
- SELECT aj.* FROM acquisition_jobs aj
- JOIN acquisition_runs ar ON ar.id = aj.run_id
- WHERE ar.batch_id = %s AND aj.query_id = %s
- ORDER BY aj.created_at, aj.platform
- """,
- (batch_id, query_id),
- )
- items = self._all(
- """
- SELECT ci.* FROM candidate_items ci
- JOIN acquisition_jobs aj ON aj.id = ci.job_id
- JOIN acquisition_runs ar ON ar.id = aj.run_id
- WHERE ar.batch_id = %s AND ci.query_id = %s
- ORDER BY ci.platform, ci.created_at
- """,
- (batch_id, query_id),
- )
- item_ids = [row["id"] for row in items]
- media: list[dict[str, Any]] = []
- classifications: list[dict[str, Any]] = []
- decode_summaries: list[dict[str, Any]] = []
- if item_ids:
- media = self._all(
- """
- SELECT * FROM media_assets
- WHERE item_id = ANY(%s)
- ORDER BY item_id, position
- """,
- (item_ids,),
- )
- classifications = self._all(
- """
- SELECT * FROM item_classifications
- WHERE item_id = ANY(%s)
- ORDER BY created_at DESC
- """,
- (item_ids,),
- )
- decode_summaries = self._all(
- """
- SELECT DISTINCT ON (dr.item_id)
- dr.item_id,
- dr.status AS decode_status,
- COUNT(DISTINCT kp.id)::int AS particle_count,
- COUNT(DISTINCT pd.id)::int AS payload_count
- FROM decode_results dr
- LEFT JOIN knowledge_particles kp ON kp.item_id = dr.item_id
- LEFT JOIN payload_drafts pd ON pd.item_id = dr.item_id
- WHERE dr.item_id = ANY(%s)
- GROUP BY dr.item_id, dr.status, dr.created_at
- ORDER BY dr.item_id, dr.created_at DESC
- """,
- (item_ids,),
- )
- return {
- "query": query,
- "jobs": jobs,
- "items": items,
- "media_assets": media,
- "classifications": classifications,
- "decode_summaries": decode_summaries,
- }
- def get_latest_query_result_list(self, query_id: UUID) -> dict[str, Any]:
- query = self._one("SELECT * FROM queries WHERE id = %s", (query_id,))
- run = self._one_or_none(
- """
- SELECT ar.* FROM acquisition_runs ar
- JOIN acquisition_jobs aj ON aj.run_id = ar.id
- WHERE aj.query_id = %s
- ORDER BY aj.created_at DESC, ar.created_at DESC
- LIMIT 1
- """,
- (query_id,),
- )
- if run is None:
- return {
- "query": query,
- "run": None,
- "jobs": [],
- "items": [],
- "media_assets": [],
- "classifications": [],
- "decode_summaries": [],
- }
- jobs = self._all(
- """
- SELECT aj.* FROM acquisition_jobs aj
- WHERE aj.query_id = %s
- ORDER BY aj.created_at, aj.platform
- """,
- (query_id,),
- )
- items = self._all(
- """
- SELECT
- ci.id,
- ci.query_id,
- ci.job_id,
- ci.platform,
- ci.title,
- LEFT(ci.raw_summary, 700) AS raw_summary,
- ci.status,
- ci.content_mode,
- ci.metadata,
- ci.created_at,
- ci.updated_at
- FROM candidate_items ci
- JOIN acquisition_jobs aj ON aj.id = ci.job_id
- WHERE aj.query_id = %s
- ORDER BY ci.platform, ci.created_at
- """,
- (query_id,),
- )
- item_ids = [row["id"] for row in items]
- media: list[dict[str, Any]] = []
- classifications: list[dict[str, Any]] = []
- decode_summaries: list[dict[str, Any]] = []
- if item_ids:
- media = self._all(
- """
- SELECT DISTINCT ON (item_id)
- id, item_id, media_type, source_url, oss_url, cdn_url, position, status
- FROM media_assets
- WHERE item_id = ANY(%s)
- ORDER BY
- item_id,
- CASE WHEN media_type IN ('cover', 'image', 'frame') THEN 0 ELSE 1 END,
- position,
- created_at
- """,
- (item_ids,),
- )
- classifications = self._all(
- """
- SELECT DISTINCT ON (item_id)
- id, item_id, is_creation_knowledge, label, confidence, status, error_message
- FROM item_classifications
- WHERE item_id = ANY(%s)
- ORDER BY item_id, created_at DESC
- """,
- (item_ids,),
- )
- decode_summaries = self._all(
- """
- SELECT DISTINCT ON (dr.item_id)
- dr.item_id,
- dr.status AS decode_status,
- COUNT(DISTINCT kp.id)::int AS particle_count,
- COUNT(DISTINCT pd.id)::int AS payload_count
- FROM decode_results dr
- LEFT JOIN knowledge_particles kp ON kp.item_id = dr.item_id
- LEFT JOIN payload_drafts pd ON pd.item_id = dr.item_id
- WHERE dr.item_id = ANY(%s)
- GROUP BY dr.item_id, dr.status, dr.created_at
- ORDER BY dr.item_id, dr.created_at DESC
- """,
- (item_ids,),
- )
- return {
- "query": query,
- "run": run,
- "jobs": jobs,
- "items": items,
- "media_assets": media,
- "classifications": classifications,
- "decode_summaries": decode_summaries,
- }
- def get_latest_singleton_overview(self) -> dict[str, Any]:
- batch = self._one_or_none(
- """
- SELECT qb.* FROM query_batches qb
- JOIN acquisition_runs ar ON ar.batch_id = qb.id
- GROUP BY qb.id
- ORDER BY MAX(ar.created_at) DESC
- LIMIT 1
- """,
- (),
- )
- if batch is None:
- return {"batch": None, "run": None, "queries": [], "decoded_items": []}
- run = self._one_or_none(
- """
- SELECT * FROM acquisition_runs
- WHERE batch_id = %s
- ORDER BY created_at DESC
- LIMIT 1
- """,
- (batch["id"],),
- )
- if run is None:
- queries = self._all(
- """
- SELECT
- q.id AS query_id,
- q.query_text,
- q.metadata->>'family_key' AS family_key,
- q.sort_order,
- 0::int AS candidate_count,
- 0::int AS creation_hit_count,
- 0::int AS decoded_count,
- 0::int AS payload_count
- FROM queries q
- WHERE q.batch_id = %s
- ORDER BY q.sort_order, q.created_at
- """,
- (batch["id"],),
- )
- return {"batch": batch, "run": None, "queries": queries, "decoded_items": []}
- latest_queries_cte = """
- WITH query_activity AS (
- SELECT
- q.id AS query_id,
- q.query_text,
- q.created_at AS query_created_at,
- MAX(ar.created_at) AS activity_at
- FROM queries q
- LEFT JOIN acquisition_jobs aj ON aj.query_id = q.id
- LEFT JOIN acquisition_runs ar ON ar.id = aj.run_id
- GROUP BY q.id, q.query_text, q.created_at
- ),
- query_stats AS (
- SELECT
- qa.query_id,
- qa.query_text,
- qa.query_created_at,
- qa.activity_at,
- COUNT(DISTINCT pd.id) FILTER (
- WHERE ic.is_creation_knowledge IS TRUE
- AND dr.status = 'decoded'
- AND kp.id IS NOT NULL
- )::int AS payload_count
- FROM query_activity qa
- LEFT JOIN acquisition_jobs aj ON aj.query_id = qa.query_id
- LEFT JOIN candidate_items ci ON ci.job_id = aj.id
- LEFT JOIN item_classifications ic ON ic.item_id = ci.id
- LEFT JOIN decode_results dr ON dr.item_id = ci.id
- LEFT JOIN knowledge_particles kp ON kp.item_id = ci.id
- LEFT JOIN payload_drafts pd ON pd.item_id = ci.id
- GROUP BY qa.query_id, qa.query_text, qa.query_created_at, qa.activity_at
- ),
- latest_queries AS (
- SELECT DISTINCT ON (query_text)
- query_id
- FROM query_stats
- ORDER BY query_text, (payload_count > 0) DESC, activity_at DESC NULLS LAST, query_created_at DESC
- )
- """
- queries = self._all(
- latest_queries_cte
- + """
- SELECT
- q.id AS query_id,
- q.query_text,
- q.metadata->>'family_key' AS family_key,
- q.sort_order,
- COUNT(DISTINCT ci.id)::int AS candidate_count,
- COUNT(DISTINCT ic.id) FILTER (
- WHERE ic.is_creation_knowledge IS TRUE
- )::int AS creation_hit_count,
- COUNT(DISTINCT ci.id) FILTER (
- WHERE ic.is_creation_knowledge IS TRUE
- AND dr.status = 'decoded'
- AND kp.id IS NOT NULL
- )::int AS decoded_count,
- COUNT(DISTINCT pd.id) FILTER (
- WHERE ic.is_creation_knowledge IS TRUE
- AND dr.status = 'decoded'
- AND kp.id IS NOT NULL
- )::int AS payload_count
- FROM latest_queries lq
- JOIN queries q ON q.id = lq.query_id
- LEFT JOIN acquisition_jobs aj ON aj.query_id = q.id
- LEFT JOIN candidate_items ci ON ci.job_id = aj.id
- LEFT JOIN item_classifications ic ON ic.item_id = ci.id
- LEFT JOIN decode_results dr ON dr.item_id = ci.id
- LEFT JOIN knowledge_particles kp ON kp.item_id = ci.id
- LEFT JOIN payload_drafts pd ON pd.item_id = ci.id
- GROUP BY q.id, q.query_text, q.metadata, q.sort_order
- ORDER BY q.sort_order, q.created_at
- """,
- (),
- )
- decoded_items = self._all(
- latest_queries_cte
- + """
- SELECT
- ci.id AS item_id,
- ci.query_id,
- ci.title,
- ci.platform,
- dr.status AS decode_status,
- COUNT(DISTINCT kp.id)::int AS particle_count,
- COUNT(DISTINCT pd.id)::int AS payload_count
- FROM candidate_items ci
- JOIN latest_queries lq ON lq.query_id = ci.query_id
- JOIN acquisition_jobs aj ON aj.id = ci.job_id
- JOIN decode_results dr ON dr.item_id = ci.id
- LEFT JOIN knowledge_particles kp ON kp.item_id = ci.id
- LEFT JOIN payload_drafts pd ON pd.item_id = ci.id
- GROUP BY ci.id, ci.query_id, ci.title, ci.platform, dr.status, dr.created_at
- ORDER BY dr.created_at DESC
- """,
- (),
- )
- return {
- "batch": batch,
- "run": run,
- "queries": queries,
- "decoded_items": decoded_items,
- }
- def list_creation_candidate_items(
- self,
- *,
- run_id: UUID | None = None,
- limit: int = 100,
- ) -> list[CandidateItem]:
- if run_id is None:
- rows = self._all(
- """
- SELECT ci.* FROM candidate_items ci
- JOIN item_classifications ic ON ic.item_id = ci.id
- WHERE ic.is_creation_knowledge IS TRUE
- AND NOT EXISTS (
- SELECT 1 FROM decode_results dr
- WHERE dr.item_id = ci.id
- AND dr.status IN ('decoded', 'rejected', 'skipped')
- )
- ORDER BY ci.created_at
- LIMIT %s
- """,
- (limit,),
- )
- else:
- rows = self._all(
- """
- SELECT ci.* FROM candidate_items ci
- JOIN acquisition_jobs aj ON aj.id = ci.job_id
- JOIN item_classifications ic ON ic.item_id = ci.id
- WHERE aj.run_id = %s AND ic.is_creation_knowledge IS TRUE
- AND NOT EXISTS (
- SELECT 1 FROM decode_results dr
- WHERE dr.item_id = ci.id
- AND dr.status IN ('decoded', 'rejected', 'skipped')
- )
- ORDER BY ci.created_at
- LIMIT %s
- """,
- (run_id, limit),
- )
- return [CandidateItem.model_validate(row) for row in rows]
- def get_candidate_item(self, item_id: UUID) -> CandidateItem:
- row = self._one("SELECT * FROM candidate_items WHERE id = %s", (item_id,))
- return CandidateItem.model_validate(row)
- def list_media_assets_for_item(self, item_id: UUID) -> list[MediaAsset]:
- rows = self._all(
- """
- SELECT * FROM media_assets
- WHERE item_id = %s
- ORDER BY position
- """,
- (item_id,),
- )
- return [MediaAsset.model_validate(row) for row in rows]
|