"""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))