service.py 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399
  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 pipeline.tracing import TraceContext, TraceWriter, hash_prompt, timed_ms
  11. from decode_content.contracts import SkillContract, load_contract
  12. from decode_content.framing import frame_and_clean
  13. from decode_content.gates import ChatJsonFn, creation_gate
  14. from decode_content.models import DecodeResult, GateResult, PayloadDraft, ReadResult
  15. from decode_content.payloads import build_payloads, validate_ingest_payload
  16. from decode_content.readers.service import read_post
  17. from decode_content.repository import DecodeRepository
  18. from decode_content.scoping import ScopeLinker, apply_scopes, nounify_scopes, scope_candidates
  19. import time
  20. @dataclass
  21. class DecodeWorkflowOutput:
  22. item_id: UUID
  23. read_result: ReadResult
  24. gate_result: GateResult
  25. knowledges: list[dict[str, Any]] = field(default_factory=list)
  26. payloads: list[dict[str, Any]] = field(default_factory=list)
  27. decode_result: DecodeResult | None = None
  28. payload_drafts: list[PayloadDraft] = field(default_factory=list)
  29. status: str = "draft"
  30. def slugify(text: str) -> str:
  31. s = re.sub(r"[^\w一-鿿]+", "-", text or "").strip("-").lower()
  32. return s or "run"
  33. class DecodeService:
  34. """Decode one creation candidate without depending on web files or SQLite."""
  35. def __init__(
  36. self,
  37. *,
  38. settings: Settings | None = None,
  39. repository: DecodeRepository | None = None,
  40. contract: SkillContract | None = None,
  41. scope_linker: ScopeLinker | None = None,
  42. chat_json_fn: ChatJsonFn | None = None,
  43. reader: Callable[[Post], ReadResult] | None = None,
  44. trace_writer: TraceWriter | None = None,
  45. trace_context: TraceContext | None = None,
  46. ) -> None:
  47. self.settings = settings
  48. self.repository = repository
  49. self.contract = contract or load_contract()
  50. self.scope_linker = scope_linker
  51. self.chat_json_fn = chat_json_fn or default_chat_json
  52. self.reader = reader
  53. self.trace_writer = trace_writer
  54. self.trace_context = trace_context or TraceContext(stage="decode")
  55. def _read(self, post: Post, context: TraceContext | None = None) -> ReadResult:
  56. if self.reader is not None:
  57. return self.reader(post)
  58. return read_post(
  59. post,
  60. settings=self.settings,
  61. trace_writer=self.trace_writer,
  62. trace_context=context or self.trace_context,
  63. )
  64. def _event(self, context: TraceContext, event_type: str, **kwargs: Any) -> None:
  65. if self.trace_writer is None:
  66. return
  67. self.trace_writer.event(context=context, stage="decode", event_type=event_type, **kwargs)
  68. def _traced_chat(self, context: TraceContext, substage: str) -> ChatJsonFn:
  69. def call(system: str, user: str, **kwargs: Any) -> dict[str, Any]:
  70. started = time.perf_counter()
  71. try:
  72. result = self.chat_json_fn(system, user, **kwargs)
  73. if self.trace_writer is not None:
  74. self.trace_writer.llm_call(
  75. context=context.child(substage=substage),
  76. stage="decode",
  77. substage=substage,
  78. provider="bailian",
  79. model_name=self.settings.llm_model if self.settings else None,
  80. prompt_name=substage,
  81. prompt_hash=hash_prompt(system),
  82. request_payload={
  83. "system": system,
  84. "user": user,
  85. "timeout": kwargs.get("timeout"),
  86. },
  87. parsed_payload=result,
  88. status="done",
  89. latency_ms=timed_ms(started),
  90. )
  91. return result
  92. except Exception as exc:
  93. if self.trace_writer is not None:
  94. self.trace_writer.llm_call(
  95. context=context.child(substage=substage),
  96. stage="decode",
  97. substage=substage,
  98. provider="bailian",
  99. model_name=self.settings.llm_model if self.settings else None,
  100. prompt_name=substage,
  101. prompt_hash=hash_prompt(system),
  102. request_payload={
  103. "system": system,
  104. "user": user,
  105. "timeout": kwargs.get("timeout"),
  106. },
  107. status="failed",
  108. error_message=str(exc),
  109. latency_ms=timed_ms(started),
  110. )
  111. raise
  112. return call
  113. def decode_post(
  114. self,
  115. *,
  116. item_id: UUID,
  117. post: Post,
  118. read_result: ReadResult | None = None,
  119. ) -> DecodeWorkflowOutput:
  120. repo = self.repository
  121. job = repo.create_decode_job(item_id=item_id, status="running") if repo else None
  122. context = self.trace_context.child(
  123. item_id=item_id,
  124. decode_job_id=getattr(job, "id", None),
  125. stage="decode",
  126. )
  127. self._event(
  128. context,
  129. "decode_started",
  130. status="running",
  131. target_table="decode_jobs",
  132. target_id=getattr(job, "id", None),
  133. )
  134. self._save_contract_artifacts()
  135. read = read_result or self._read(post, context=context.child(substage="read_post"))
  136. self._event(
  137. context,
  138. "read_completed",
  139. status="done" if not read.is_empty else "skipped",
  140. payload={"is_empty": read.is_empty, "card_count": len(read.cards or [])},
  141. )
  142. if read.is_empty:
  143. gate = GateResult(passed=False, reason="读懂结果为空", details={"is_empty": True})
  144. decode_result = self._save_decode_result(
  145. item_id=item_id,
  146. job_id=getattr(job, "id", None),
  147. read=read,
  148. gate=gate,
  149. knowledges=[],
  150. status="skipped",
  151. )
  152. return DecodeWorkflowOutput(
  153. item_id=item_id,
  154. read_result=read,
  155. gate_result=gate,
  156. decode_result=decode_result,
  157. status="skipped",
  158. )
  159. gate = creation_gate(read.text, chat_json_fn=self._traced_chat(context, "creation_gate"))
  160. self._event(
  161. context,
  162. "creation_gate_completed",
  163. status="done" if gate.passed else "skipped",
  164. payload=gate.model_dump(mode="json"),
  165. )
  166. if not gate.passed:
  167. decode_result = self._save_decode_result(
  168. item_id=item_id,
  169. job_id=getattr(job, "id", None),
  170. read=read,
  171. gate=gate,
  172. knowledges=[],
  173. status="rejected",
  174. )
  175. return DecodeWorkflowOutput(
  176. item_id=item_id,
  177. read_result=read,
  178. gate_result=gate,
  179. decode_result=decode_result,
  180. status="rejected",
  181. )
  182. knowledges = frame_and_clean(
  183. post,
  184. read.text,
  185. contract=self.contract,
  186. chat_json_fn=self._traced_chat(context, "frame_and_clean"),
  187. )
  188. self._event(
  189. context,
  190. "framing_completed",
  191. status="done",
  192. payload={"knowledge_count": len(knowledges)},
  193. )
  194. scopes = scope_candidates(
  195. knowledges,
  196. contract=self.contract,
  197. chat_json_fn=self._traced_chat(context, "scope_candidates"),
  198. )
  199. scopes = nounify_scopes(scopes, chat_json_fn=self._traced_chat(context, "scope_nounify"))
  200. apply_scopes(knowledges, scopes, self.scope_linker)
  201. self._event(
  202. context,
  203. "scope_completed",
  204. status="done",
  205. payload={"scope_group_count": len(scopes)},
  206. )
  207. payloads = build_payloads(post, knowledges)
  208. self._event(
  209. context,
  210. "payloads_built",
  211. status="done",
  212. payload={"payload_count": len(payloads)},
  213. )
  214. decode_result = self._save_decode_result(
  215. item_id=item_id,
  216. job_id=getattr(job, "id", None),
  217. read=read,
  218. gate=gate,
  219. knowledges=knowledges,
  220. status="decoded",
  221. )
  222. payload_drafts = self._save_particles_and_payloads(
  223. item_id=item_id,
  224. decode_result=decode_result,
  225. knowledges=knowledges,
  226. payloads=payloads,
  227. )
  228. self._event(
  229. context.child(decode_result_id=getattr(decode_result, "id", None)),
  230. "decode_finished",
  231. status="decoded",
  232. target_table="decode_results",
  233. target_id=getattr(decode_result, "id", None),
  234. payload={"knowledge_count": len(knowledges), "payload_draft_count": len(payload_drafts)},
  235. )
  236. return DecodeWorkflowOutput(
  237. item_id=item_id,
  238. read_result=read,
  239. gate_result=gate,
  240. knowledges=knowledges,
  241. payloads=payloads,
  242. decode_result=decode_result,
  243. payload_drafts=payload_drafts,
  244. status="decoded",
  245. )
  246. def _save_contract_artifacts(self) -> None:
  247. if self.repository is None:
  248. return
  249. for snapshot in self.contract.snapshots():
  250. self.repository.save_contract_snapshot(
  251. contract_name=snapshot.contract_name,
  252. contract_type=snapshot.contract_type,
  253. version_label=snapshot.version_label,
  254. content_hash=snapshot.content_hash,
  255. source_path=snapshot.source_path,
  256. snapshot=snapshot.snapshot,
  257. pipeline_run_id=self.trace_context.pipeline_run_id,
  258. )
  259. def _save_decode_result(
  260. self,
  261. *,
  262. item_id: UUID,
  263. job_id: UUID | None,
  264. read: ReadResult,
  265. gate: GateResult,
  266. knowledges: list[dict[str, Any]],
  267. status: str,
  268. ) -> DecodeResult:
  269. framing_result = {"knowledges": knowledges, "contract_hash": self.contract.content_hash}
  270. if self.repository is None:
  271. return DecodeResult(
  272. item_id=item_id,
  273. decode_job_id=job_id,
  274. read_result=read.model_dump(mode="json"),
  275. gate_result=gate.model_dump(mode="json"),
  276. framing_result=framing_result,
  277. status=status,
  278. )
  279. return self.repository.save_decode_result(
  280. item_id=item_id,
  281. decode_job_id=job_id,
  282. read_result=read.model_dump(mode="json"),
  283. gate_result=gate.model_dump(mode="json"),
  284. framing_result=framing_result,
  285. status=status,
  286. )
  287. def _save_particles_and_payloads(
  288. self,
  289. *,
  290. item_id: UUID,
  291. decode_result: DecodeResult,
  292. knowledges: list[dict[str, Any]],
  293. payloads: list[dict[str, Any]],
  294. ) -> list[PayloadDraft]:
  295. if self.repository is None:
  296. return [PayloadDraft(item_id=item_id, payload=payload) for payload in payloads]
  297. drafts: list[PayloadDraft] = []
  298. for knowledge, payload in zip(knowledges, payloads):
  299. particle = self.repository.save_knowledge_particle(
  300. item_id=item_id,
  301. decode_result_id=decode_result.id,
  302. particle_type=knowledge.get("type") or "how",
  303. title=knowledge.get("title") or "",
  304. content=knowledge,
  305. status="draft",
  306. )
  307. for scope in payload.get("scopes") or []:
  308. source_scope = self._scope_source(knowledge, scope)
  309. self.repository.save_scope_result(
  310. item_id=item_id,
  311. particle_id=particle.id,
  312. scope_type=scope["scope_type"],
  313. scope_value=scope["value"],
  314. evidence={k: v for k, v in source_scope.items() if k not in {"scope_type", "value"}},
  315. status="draft",
  316. )
  317. drafts.append(
  318. self.repository.save_payload_draft(
  319. item_id=item_id,
  320. particle_id=particle.id,
  321. payload=_validated_payload(payload),
  322. review_status="pending",
  323. ingest_ready=False,
  324. status="draft",
  325. )
  326. )
  327. return drafts
  328. @staticmethod
  329. def _scope_source(knowledge: dict[str, Any], scope: dict[str, Any]) -> dict[str, Any]:
  330. source_scopes = list(knowledge.get("作用域", []))
  331. for step in knowledge.get("steps", []):
  332. source_scopes.extend(step.get("作用域", []))
  333. return next(
  334. (
  335. item
  336. for item in source_scopes
  337. if item.get("scope_type") == scope["scope_type"] and item.get("value") == scope["value"]
  338. ),
  339. scope,
  340. )
  341. def _validated_payload(payload: dict[str, Any]) -> dict[str, Any]:
  342. validate_ingest_payload(payload)
  343. return payload
  344. def decode_post(
  345. *,
  346. item_id: UUID,
  347. post: Post,
  348. settings: Settings | None = None,
  349. repository: DecodeRepository | None = None,
  350. scope_linker: ScopeLinker | None = None,
  351. chat_json_fn: ChatJsonFn | None = None,
  352. reader: Callable[[Post], ReadResult] | None = None,
  353. read_result: ReadResult | None = None,
  354. trace_writer: TraceWriter | None = None,
  355. trace_context: TraceContext | None = None,
  356. ) -> DecodeWorkflowOutput:
  357. service = DecodeService(
  358. settings=settings,
  359. repository=repository,
  360. scope_linker=scope_linker,
  361. chat_json_fn=chat_json_fn,
  362. reader=reader,
  363. trace_writer=trace_writer,
  364. trace_context=trace_context,
  365. )
  366. return service.decode_post(item_id=item_id, post=post, read_result=read_result)
  367. def decode_item(
  368. item_id: UUID,
  369. *,
  370. load_post: Callable[[UUID], Post],
  371. service: DecodeService | None = None,
  372. ) -> DecodeWorkflowOutput:
  373. svc = service or DecodeService()
  374. return svc.decode_post(item_id=item_id, post=load_post(item_id))