base.py 5.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193
  1. """Repository contracts for the formal acquisition state store."""
  2. from __future__ import annotations
  3. from typing import Any, Protocol
  4. from uuid import UUID
  5. from acquisition.domain import (
  6. AcquisitionJob,
  7. AcquisitionRun,
  8. CandidateItem,
  9. ItemClassification,
  10. MediaAsset,
  11. Query,
  12. QueryBatch,
  13. )
  14. class AcquisitionRepository(Protocol):
  15. """Persistence boundary used by acquisition runners and API routes."""
  16. def create_query_batch(
  17. self,
  18. *,
  19. name: str,
  20. source_type: str = "manual",
  21. generation_method: str | None = None,
  22. target_platforms: list[str] | None = None,
  23. status: str = "draft",
  24. metadata: dict[str, Any] | None = None,
  25. ) -> QueryBatch:
  26. ...
  27. def add_query(
  28. self,
  29. *,
  30. batch_id: UUID | None,
  31. query_text: str,
  32. axes: dict[str, Any] | None = None,
  33. keep: bool | None = None,
  34. filter_reason: str | None = None,
  35. status: str = "draft",
  36. sort_order: int = 0,
  37. metadata: dict[str, Any] | None = None,
  38. ) -> Query:
  39. ...
  40. def list_queries_for_batch(
  41. self,
  42. batch_id: UUID,
  43. *,
  44. keep: bool | None = None,
  45. ) -> list[Query]:
  46. ...
  47. def get_query_batch(self, batch_id: UUID) -> QueryBatch:
  48. ...
  49. def create_acquisition_run(
  50. self,
  51. *,
  52. batch_id: UUID | None = None,
  53. run_key: str | None = None,
  54. status: str = "pending",
  55. note: str | None = None,
  56. metadata: dict[str, Any] | None = None,
  57. ) -> AcquisitionRun:
  58. ...
  59. def ensure_acquisition_job(
  60. self,
  61. *,
  62. run_id: UUID,
  63. query_id: UUID | None,
  64. platform: str,
  65. search_limit: int | None = None,
  66. display_limit: int | None = None,
  67. status: str = "pending",
  68. metadata: dict[str, Any] | None = None,
  69. ) -> AcquisitionJob:
  70. ...
  71. def update_acquisition_job(
  72. self,
  73. job_id: UUID,
  74. *,
  75. status: str,
  76. attempt_count: int | None = None,
  77. error_message: str | None = None,
  78. metadata: dict[str, Any] | None = None,
  79. ) -> AcquisitionJob:
  80. ...
  81. def update_acquisition_run(
  82. self,
  83. run_id: UUID,
  84. *,
  85. status: str,
  86. error_message: str | None = None,
  87. metadata: dict[str, Any] | None = None,
  88. ) -> AcquisitionRun:
  89. ...
  90. def upsert_candidate_item(
  91. self,
  92. *,
  93. platform: str,
  94. job_id: UUID | None = None,
  95. query_id: UUID | None = None,
  96. platform_item_id: str | None = None,
  97. unique_key: str | None = None,
  98. canonical_url: str | None = None,
  99. content_type: str | None = None,
  100. content_mode: str | None = None,
  101. title: str | None = None,
  102. author_name: str | None = None,
  103. body_text: str | None = None,
  104. raw_summary: str | None = None,
  105. status: str = "candidate",
  106. source_payload: dict[str, Any] | None = None,
  107. metadata: dict[str, Any] | None = None,
  108. error_message: str | None = None,
  109. ) -> CandidateItem:
  110. ...
  111. def get_candidate_item_by_unique_key(self, unique_key: str) -> CandidateItem | None:
  112. ...
  113. def attach_existing_candidate_item(
  114. self,
  115. item_id: UUID,
  116. *,
  117. job_id: UUID,
  118. query_id: UUID,
  119. metadata: dict[str, Any] | None = None,
  120. ) -> CandidateItem:
  121. ...
  122. def add_media_asset(
  123. self,
  124. *,
  125. item_id: UUID,
  126. media_type: str,
  127. source_url: str | None = None,
  128. oss_url: str | None = None,
  129. cdn_url: str | None = None,
  130. position: int = 0,
  131. status: str = "pending",
  132. source_payload: dict[str, Any] | None = None,
  133. metadata: dict[str, Any] | None = None,
  134. ) -> MediaAsset:
  135. ...
  136. def add_item_classification(
  137. self,
  138. *,
  139. item_id: UUID,
  140. is_creation_knowledge: bool | None = None,
  141. label: str | None = None,
  142. confidence: float | None = None,
  143. reason: str | None = None,
  144. model_name: str | None = None,
  145. prompt_version: str | None = None,
  146. result_payload: dict[str, Any] | None = None,
  147. status: str = "pending",
  148. error_message: str | None = None,
  149. ) -> ItemClassification:
  150. ...
  151. def get_run_summary(self, run_id: UUID) -> dict[str, Any]:
  152. ...
  153. def get_query_detail(self, *, run_id: UUID, query_id: UUID) -> dict[str, Any]:
  154. ...
  155. def get_query_detail_for_batch(self, *, batch_id: UUID, query_id: UUID) -> dict[str, Any]:
  156. ...
  157. def get_latest_query_result_list(self, query_id: UUID) -> dict[str, Any]:
  158. ...
  159. def list_creation_candidate_items(
  160. self,
  161. *,
  162. run_id: UUID | None = None,
  163. limit: int = 100,
  164. ) -> list[CandidateItem]:
  165. ...
  166. def get_candidate_item(self, item_id: UUID) -> CandidateItem:
  167. ...
  168. def list_media_assets_for_item(self, item_id: UUID) -> list[MediaAsset]:
  169. ...