service.py 9.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278
  1. """Formal single-item decode orchestration."""
  2. from __future__ import annotations
  3. import re
  4. from dataclasses import dataclass, field
  5. from typing import Any, Callable
  6. from uuid import UUID
  7. from core.config import Settings
  8. from core.llm import chat_json as default_chat_json
  9. from core.models import Post
  10. from decode_content.contracts import SkillContract, load_contract
  11. from decode_content.framing import frame_and_clean
  12. from decode_content.gates import ChatJsonFn, creation_gate
  13. from decode_content.models import DecodeResult, GateResult, PayloadDraft, ReadResult
  14. from decode_content.payloads import build_payloads, validate_ingest_payload
  15. from decode_content.readers.service import read_post
  16. from decode_content.repository import DecodeRepository
  17. from decode_content.scoping import ScopeLinker, apply_scopes, nounify_scopes, scope_candidates
  18. @dataclass
  19. class DecodeWorkflowOutput:
  20. item_id: UUID
  21. read_result: ReadResult
  22. gate_result: GateResult
  23. knowledges: list[dict[str, Any]] = field(default_factory=list)
  24. payloads: list[dict[str, Any]] = field(default_factory=list)
  25. decode_result: DecodeResult | None = None
  26. payload_drafts: list[PayloadDraft] = field(default_factory=list)
  27. status: str = "draft"
  28. def slugify(text: str) -> str:
  29. s = re.sub(r"[^\w一-鿿]+", "-", text or "").strip("-").lower()
  30. return s or "run"
  31. class DecodeService:
  32. """Decode one creation candidate without depending on web files or SQLite."""
  33. def __init__(
  34. self,
  35. *,
  36. settings: Settings | None = None,
  37. repository: DecodeRepository | None = None,
  38. contract: SkillContract | None = None,
  39. scope_linker: ScopeLinker | None = None,
  40. chat_json_fn: ChatJsonFn | None = None,
  41. reader: Callable[[Post], ReadResult] | None = None,
  42. ) -> None:
  43. self.settings = settings
  44. self.repository = repository
  45. self.contract = contract or load_contract()
  46. self.scope_linker = scope_linker
  47. self.chat_json_fn = chat_json_fn or default_chat_json
  48. self.reader = reader
  49. def _read(self, post: Post) -> ReadResult:
  50. if self.reader is not None:
  51. return self.reader(post)
  52. return read_post(post, settings=self.settings)
  53. def decode_post(
  54. self,
  55. *,
  56. item_id: UUID,
  57. post: Post,
  58. read_result: ReadResult | None = None,
  59. ) -> DecodeWorkflowOutput:
  60. repo = self.repository
  61. job = repo.create_decode_job(item_id=item_id, status="running") if repo else None
  62. self._save_contract_snapshots()
  63. read = read_result or self._read(post)
  64. if read.is_empty:
  65. gate = GateResult(passed=False, reason="读懂结果为空", details={"is_empty": True})
  66. decode_result = self._save_decode_result(
  67. item_id=item_id,
  68. job_id=getattr(job, "id", None),
  69. read=read,
  70. gate=gate,
  71. knowledges=[],
  72. status="skipped",
  73. )
  74. return DecodeWorkflowOutput(
  75. item_id=item_id,
  76. read_result=read,
  77. gate_result=gate,
  78. decode_result=decode_result,
  79. status="skipped",
  80. )
  81. gate = creation_gate(read.text, chat_json_fn=self.chat_json_fn)
  82. if not gate.passed:
  83. decode_result = self._save_decode_result(
  84. item_id=item_id,
  85. job_id=getattr(job, "id", None),
  86. read=read,
  87. gate=gate,
  88. knowledges=[],
  89. status="rejected",
  90. )
  91. return DecodeWorkflowOutput(
  92. item_id=item_id,
  93. read_result=read,
  94. gate_result=gate,
  95. decode_result=decode_result,
  96. status="rejected",
  97. )
  98. knowledges = frame_and_clean(
  99. post,
  100. read.text,
  101. contract=self.contract,
  102. chat_json_fn=self.chat_json_fn,
  103. )
  104. scopes = scope_candidates(knowledges, contract=self.contract, chat_json_fn=self.chat_json_fn)
  105. scopes = nounify_scopes(scopes, chat_json_fn=self.chat_json_fn)
  106. apply_scopes(knowledges, scopes, self.scope_linker)
  107. payloads = build_payloads(post, knowledges)
  108. decode_result = self._save_decode_result(
  109. item_id=item_id,
  110. job_id=getattr(job, "id", None),
  111. read=read,
  112. gate=gate,
  113. knowledges=knowledges,
  114. status="decoded",
  115. )
  116. payload_drafts = self._save_particles_and_payloads(
  117. item_id=item_id,
  118. decode_result=decode_result,
  119. knowledges=knowledges,
  120. payloads=payloads,
  121. )
  122. return DecodeWorkflowOutput(
  123. item_id=item_id,
  124. read_result=read,
  125. gate_result=gate,
  126. knowledges=knowledges,
  127. payloads=payloads,
  128. decode_result=decode_result,
  129. payload_drafts=payload_drafts,
  130. status="decoded",
  131. )
  132. def _save_contract_snapshots(self) -> None:
  133. if self.repository is None:
  134. return
  135. for snapshot in self.contract.snapshots():
  136. self.repository.save_contract_snapshot(
  137. contract_name=snapshot.contract_name,
  138. contract_type=snapshot.contract_type,
  139. version_label=snapshot.version_label,
  140. content_hash=snapshot.content_hash,
  141. source_path=snapshot.source_path,
  142. snapshot=snapshot.snapshot,
  143. )
  144. def _save_decode_result(
  145. self,
  146. *,
  147. item_id: UUID,
  148. job_id: UUID | None,
  149. read: ReadResult,
  150. gate: GateResult,
  151. knowledges: list[dict[str, Any]],
  152. status: str,
  153. ) -> DecodeResult:
  154. framing_result = {"knowledges": knowledges, "contract_hash": self.contract.content_hash}
  155. if self.repository is None:
  156. return DecodeResult(
  157. item_id=item_id,
  158. decode_job_id=job_id,
  159. read_result=read.model_dump(mode="json"),
  160. gate_result=gate.model_dump(mode="json"),
  161. framing_result=framing_result,
  162. status=status,
  163. )
  164. return self.repository.save_decode_result(
  165. item_id=item_id,
  166. decode_job_id=job_id,
  167. read_result=read.model_dump(mode="json"),
  168. gate_result=gate.model_dump(mode="json"),
  169. framing_result=framing_result,
  170. status=status,
  171. )
  172. def _save_particles_and_payloads(
  173. self,
  174. *,
  175. item_id: UUID,
  176. decode_result: DecodeResult,
  177. knowledges: list[dict[str, Any]],
  178. payloads: list[dict[str, Any]],
  179. ) -> list[PayloadDraft]:
  180. if self.repository is None:
  181. return [PayloadDraft(item_id=item_id, payload=payload) for payload in payloads]
  182. drafts: list[PayloadDraft] = []
  183. for knowledge, payload in zip(knowledges, payloads):
  184. particle = self.repository.save_knowledge_particle(
  185. item_id=item_id,
  186. decode_result_id=decode_result.id,
  187. particle_type=knowledge.get("type") or "how",
  188. title=knowledge.get("title") or "",
  189. content=knowledge,
  190. status="draft",
  191. )
  192. for scope in payload.get("scopes") or []:
  193. source_scope = self._scope_source(knowledge, scope)
  194. self.repository.save_scope_result(
  195. item_id=item_id,
  196. particle_id=particle.id,
  197. scope_type=scope["scope_type"],
  198. scope_value=scope["value"],
  199. evidence={k: v for k, v in source_scope.items() if k not in {"scope_type", "value"}},
  200. status="draft",
  201. )
  202. drafts.append(
  203. self.repository.save_payload_draft(
  204. item_id=item_id,
  205. particle_id=particle.id,
  206. payload=_validated_payload(payload),
  207. review_status="pending",
  208. ingest_ready=False,
  209. status="draft",
  210. )
  211. )
  212. return drafts
  213. @staticmethod
  214. def _scope_source(knowledge: dict[str, Any], scope: dict[str, Any]) -> dict[str, Any]:
  215. source_scopes = list(knowledge.get("作用域", []))
  216. for step in knowledge.get("steps", []):
  217. source_scopes.extend(step.get("作用域", []))
  218. return next(
  219. (
  220. item
  221. for item in source_scopes
  222. if item.get("scope_type") == scope["scope_type"] and item.get("value") == scope["value"]
  223. ),
  224. scope,
  225. )
  226. def _validated_payload(payload: dict[str, Any]) -> dict[str, Any]:
  227. validate_ingest_payload(payload)
  228. return payload
  229. def decode_post(
  230. *,
  231. item_id: UUID,
  232. post: Post,
  233. settings: Settings | None = None,
  234. repository: DecodeRepository | None = None,
  235. scope_linker: ScopeLinker | None = None,
  236. chat_json_fn: ChatJsonFn | None = None,
  237. reader: Callable[[Post], ReadResult] | None = None,
  238. read_result: ReadResult | None = None,
  239. ) -> DecodeWorkflowOutput:
  240. service = DecodeService(
  241. settings=settings,
  242. repository=repository,
  243. scope_linker=scope_linker,
  244. chat_json_fn=chat_json_fn,
  245. reader=reader,
  246. )
  247. return service.decode_post(item_id=item_id, post=post, read_result=read_result)
  248. def decode_item(
  249. item_id: UUID,
  250. *,
  251. load_post: Callable[[UUID], Post],
  252. service: DecodeService | None = None,
  253. ) -> DecodeWorkflowOutput:
  254. svc = service or DecodeService()
  255. return svc.decode_post(item_id=item_id, post=load_post(item_id))