| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399 |
- """Formal single-item decode orchestration."""
- from __future__ import annotations
- import re
- from dataclasses import dataclass, field
- from typing import Any, Callable
- from uuid import UUID
- from core.config import Settings
- from core.llm import chat_json as default_chat_json
- from core.models import Post
- from pipeline.tracing import TraceContext, TraceWriter, hash_prompt, timed_ms
- from decode_content.contracts import SkillContract, load_contract
- from decode_content.framing import frame_and_clean
- from decode_content.gates import ChatJsonFn, creation_gate
- from decode_content.models import DecodeResult, GateResult, PayloadDraft, ReadResult
- from decode_content.payloads import build_payloads, validate_ingest_payload
- from decode_content.readers.service import read_post
- from decode_content.repository import DecodeRepository
- from decode_content.scoping import ScopeLinker, apply_scopes, nounify_scopes, scope_candidates
- import time
- @dataclass
- class DecodeWorkflowOutput:
- item_id: UUID
- read_result: ReadResult
- gate_result: GateResult
- knowledges: list[dict[str, Any]] = field(default_factory=list)
- payloads: list[dict[str, Any]] = field(default_factory=list)
- decode_result: DecodeResult | None = None
- payload_drafts: list[PayloadDraft] = field(default_factory=list)
- status: str = "draft"
- def slugify(text: str) -> str:
- s = re.sub(r"[^\w一-鿿]+", "-", text or "").strip("-").lower()
- return s or "run"
- class DecodeService:
- """Decode one creation candidate without depending on web files or SQLite."""
- def __init__(
- self,
- *,
- settings: Settings | None = None,
- repository: DecodeRepository | None = None,
- contract: SkillContract | None = None,
- scope_linker: ScopeLinker | None = None,
- chat_json_fn: ChatJsonFn | None = None,
- reader: Callable[[Post], ReadResult] | None = None,
- trace_writer: TraceWriter | None = None,
- trace_context: TraceContext | None = None,
- ) -> None:
- self.settings = settings
- self.repository = repository
- self.contract = contract or load_contract()
- self.scope_linker = scope_linker
- self.chat_json_fn = chat_json_fn or default_chat_json
- self.reader = reader
- self.trace_writer = trace_writer
- self.trace_context = trace_context or TraceContext(stage="decode")
- def _read(self, post: Post, context: TraceContext | None = None) -> ReadResult:
- if self.reader is not None:
- return self.reader(post)
- return read_post(
- post,
- settings=self.settings,
- trace_writer=self.trace_writer,
- trace_context=context or self.trace_context,
- )
- def _event(self, context: TraceContext, event_type: str, **kwargs: Any) -> None:
- if self.trace_writer is None:
- return
- self.trace_writer.event(context=context, stage="decode", event_type=event_type, **kwargs)
- def _traced_chat(self, context: TraceContext, substage: str) -> ChatJsonFn:
- def call(system: str, user: str, **kwargs: Any) -> dict[str, Any]:
- started = time.perf_counter()
- try:
- result = self.chat_json_fn(system, user, **kwargs)
- if self.trace_writer is not None:
- self.trace_writer.llm_call(
- context=context.child(substage=substage),
- stage="decode",
- substage=substage,
- provider="bailian",
- model_name=self.settings.llm_model if self.settings else None,
- prompt_name=substage,
- prompt_hash=hash_prompt(system),
- request_payload={
- "system": system,
- "user": user,
- "timeout": kwargs.get("timeout"),
- },
- parsed_payload=result,
- status="done",
- latency_ms=timed_ms(started),
- )
- return result
- except Exception as exc:
- if self.trace_writer is not None:
- self.trace_writer.llm_call(
- context=context.child(substage=substage),
- stage="decode",
- substage=substage,
- provider="bailian",
- model_name=self.settings.llm_model if self.settings else None,
- prompt_name=substage,
- prompt_hash=hash_prompt(system),
- request_payload={
- "system": system,
- "user": user,
- "timeout": kwargs.get("timeout"),
- },
- status="failed",
- error_message=str(exc),
- latency_ms=timed_ms(started),
- )
- raise
- return call
- def decode_post(
- self,
- *,
- item_id: UUID,
- post: Post,
- read_result: ReadResult | None = None,
- ) -> DecodeWorkflowOutput:
- repo = self.repository
- job = repo.create_decode_job(item_id=item_id, status="running") if repo else None
- context = self.trace_context.child(
- item_id=item_id,
- decode_job_id=getattr(job, "id", None),
- stage="decode",
- )
- self._event(
- context,
- "decode_started",
- status="running",
- target_table="decode_jobs",
- target_id=getattr(job, "id", None),
- )
- self._save_contract_artifacts()
- read = read_result or self._read(post, context=context.child(substage="read_post"))
- self._event(
- context,
- "read_completed",
- status="done" if not read.is_empty else "skipped",
- payload={"is_empty": read.is_empty, "card_count": len(read.cards or [])},
- )
- if read.is_empty:
- gate = GateResult(passed=False, reason="读懂结果为空", details={"is_empty": True})
- decode_result = self._save_decode_result(
- item_id=item_id,
- job_id=getattr(job, "id", None),
- read=read,
- gate=gate,
- knowledges=[],
- status="skipped",
- )
- return DecodeWorkflowOutput(
- item_id=item_id,
- read_result=read,
- gate_result=gate,
- decode_result=decode_result,
- status="skipped",
- )
- gate = creation_gate(read.text, chat_json_fn=self._traced_chat(context, "creation_gate"))
- self._event(
- context,
- "creation_gate_completed",
- status="done" if gate.passed else "skipped",
- payload=gate.model_dump(mode="json"),
- )
- if not gate.passed:
- decode_result = self._save_decode_result(
- item_id=item_id,
- job_id=getattr(job, "id", None),
- read=read,
- gate=gate,
- knowledges=[],
- status="rejected",
- )
- return DecodeWorkflowOutput(
- item_id=item_id,
- read_result=read,
- gate_result=gate,
- decode_result=decode_result,
- status="rejected",
- )
- knowledges = frame_and_clean(
- post,
- read.text,
- contract=self.contract,
- chat_json_fn=self._traced_chat(context, "frame_and_clean"),
- )
- self._event(
- context,
- "framing_completed",
- status="done",
- payload={"knowledge_count": len(knowledges)},
- )
- scopes = scope_candidates(
- knowledges,
- contract=self.contract,
- chat_json_fn=self._traced_chat(context, "scope_candidates"),
- )
- scopes = nounify_scopes(scopes, chat_json_fn=self._traced_chat(context, "scope_nounify"))
- apply_scopes(knowledges, scopes, self.scope_linker)
- self._event(
- context,
- "scope_completed",
- status="done",
- payload={"scope_group_count": len(scopes)},
- )
- payloads = build_payloads(post, knowledges)
- self._event(
- context,
- "payloads_built",
- status="done",
- payload={"payload_count": len(payloads)},
- )
- decode_result = self._save_decode_result(
- item_id=item_id,
- job_id=getattr(job, "id", None),
- read=read,
- gate=gate,
- knowledges=knowledges,
- status="decoded",
- )
- payload_drafts = self._save_particles_and_payloads(
- item_id=item_id,
- decode_result=decode_result,
- knowledges=knowledges,
- payloads=payloads,
- )
- self._event(
- context.child(decode_result_id=getattr(decode_result, "id", None)),
- "decode_finished",
- status="decoded",
- target_table="decode_results",
- target_id=getattr(decode_result, "id", None),
- payload={"knowledge_count": len(knowledges), "payload_draft_count": len(payload_drafts)},
- )
- return DecodeWorkflowOutput(
- item_id=item_id,
- read_result=read,
- gate_result=gate,
- knowledges=knowledges,
- payloads=payloads,
- decode_result=decode_result,
- payload_drafts=payload_drafts,
- status="decoded",
- )
- def _save_contract_artifacts(self) -> None:
- if self.repository is None:
- return
- for snapshot in self.contract.snapshots():
- self.repository.save_contract_snapshot(
- contract_name=snapshot.contract_name,
- contract_type=snapshot.contract_type,
- version_label=snapshot.version_label,
- content_hash=snapshot.content_hash,
- source_path=snapshot.source_path,
- snapshot=snapshot.snapshot,
- pipeline_run_id=self.trace_context.pipeline_run_id,
- )
- def _save_decode_result(
- self,
- *,
- item_id: UUID,
- job_id: UUID | None,
- read: ReadResult,
- gate: GateResult,
- knowledges: list[dict[str, Any]],
- status: str,
- ) -> DecodeResult:
- framing_result = {"knowledges": knowledges, "contract_hash": self.contract.content_hash}
- if self.repository is None:
- return DecodeResult(
- item_id=item_id,
- decode_job_id=job_id,
- read_result=read.model_dump(mode="json"),
- gate_result=gate.model_dump(mode="json"),
- framing_result=framing_result,
- status=status,
- )
- return self.repository.save_decode_result(
- item_id=item_id,
- decode_job_id=job_id,
- read_result=read.model_dump(mode="json"),
- gate_result=gate.model_dump(mode="json"),
- framing_result=framing_result,
- status=status,
- )
- def _save_particles_and_payloads(
- self,
- *,
- item_id: UUID,
- decode_result: DecodeResult,
- knowledges: list[dict[str, Any]],
- payloads: list[dict[str, Any]],
- ) -> list[PayloadDraft]:
- if self.repository is None:
- return [PayloadDraft(item_id=item_id, payload=payload) for payload in payloads]
- drafts: list[PayloadDraft] = []
- for knowledge, payload in zip(knowledges, payloads):
- particle = self.repository.save_knowledge_particle(
- item_id=item_id,
- decode_result_id=decode_result.id,
- particle_type=knowledge.get("type") or "how",
- title=knowledge.get("title") or "",
- content=knowledge,
- status="draft",
- )
- for scope in payload.get("scopes") or []:
- source_scope = self._scope_source(knowledge, scope)
- self.repository.save_scope_result(
- item_id=item_id,
- particle_id=particle.id,
- scope_type=scope["scope_type"],
- scope_value=scope["value"],
- evidence={k: v for k, v in source_scope.items() if k not in {"scope_type", "value"}},
- status="draft",
- )
- drafts.append(
- self.repository.save_payload_draft(
- item_id=item_id,
- particle_id=particle.id,
- payload=_validated_payload(payload),
- review_status="pending",
- ingest_ready=False,
- status="draft",
- )
- )
- return drafts
- @staticmethod
- def _scope_source(knowledge: dict[str, Any], scope: dict[str, Any]) -> dict[str, Any]:
- source_scopes = list(knowledge.get("作用域", []))
- for step in knowledge.get("steps", []):
- source_scopes.extend(step.get("作用域", []))
- return next(
- (
- item
- for item in source_scopes
- if item.get("scope_type") == scope["scope_type"] and item.get("value") == scope["value"]
- ),
- scope,
- )
- def _validated_payload(payload: dict[str, Any]) -> dict[str, Any]:
- validate_ingest_payload(payload)
- return payload
- def decode_post(
- *,
- item_id: UUID,
- post: Post,
- settings: Settings | None = None,
- repository: DecodeRepository | None = None,
- scope_linker: ScopeLinker | None = None,
- chat_json_fn: ChatJsonFn | None = None,
- reader: Callable[[Post], ReadResult] | None = None,
- read_result: ReadResult | None = None,
- trace_writer: TraceWriter | None = None,
- trace_context: TraceContext | None = None,
- ) -> DecodeWorkflowOutput:
- service = DecodeService(
- settings=settings,
- repository=repository,
- scope_linker=scope_linker,
- chat_json_fn=chat_json_fn,
- reader=reader,
- trace_writer=trace_writer,
- trace_context=trace_context,
- )
- return service.decode_post(item_id=item_id, post=post, read_result=read_result)
- def decode_item(
- item_id: UUID,
- *,
- load_post: Callable[[UUID], Post],
- service: DecodeService | None = None,
- ) -> DecodeWorkflowOutput:
- svc = service or DecodeService()
- return svc.decode_post(item_id=item_id, post=load_post(item_id))
|