from __future__ import annotations import json import re import time from dataclasses import asdict from pathlib import Path from typing import Any from claridoc.lint import lint_document, render_lint_markdown from claridoc.models import ( Brief, LintIssue, LintReport, ModelReview, Outline, PipelineConfig, ReviewIssue, RoundResult, RunResult, Severity, SourcePack, ValidationError, ) from claridoc.prompts import drafting_prompt, planning_prompt, review_prompt, revision_prompt from claridoc.providers import ProviderError, ProviderRequest, create_provider from claridoc.provenance import build_evidence_map, render_provenance from claridoc.report import render_run_report from claridoc.structures import create_outline, reconcile_outline from claridoc.utils import atomic_write_text, extract_json_object, sha256_file, utc_now_iso, write_json class PipelineExecutionError(RuntimeError): """Raised when a required stage cannot complete.""" def run_pipeline( brief: Brief, sources: SourcePack, config: PipelineConfig, output_dir: str | Path, ) -> RunResult: output = Path(output_dir).resolve() output.mkdir(parents=True, exist_ok=True) for directory in ("inputs", "stages", "rounds", "final"): (output / directory).mkdir(parents=True, exist_ok=True) warnings: list[str] = [] events: list[dict[str, Any]] = [] provider_warning = _mock_provider_warning(config) if provider_warning: warnings.append(provider_warning) write_json(output / "inputs" / "brief.normalized.json", brief.to_dict()) write_json(output / "inputs" / "sources.normalized.json", sources.to_dict()) write_json(output / "inputs" / "pipeline.normalized.json", config.to_dict()) base_outline = create_outline(brief, sources) outline = base_outline planner = create_provider(config.planner) plan_prompt = planning_prompt(brief, base_outline, sources) try: response = _invoke(planner, ProviderRequest("plan", plan_prompt, output, {"document_type": brief.document_type.value}), events) atomic_write_text(output / "stages" / "01-planner.raw.txt", response.text + "\n") candidate = Outline.from_dict(extract_json_object(response.text)) outline = reconcile_outline(base_outline, candidate, sources) except (ProviderError, ValidationError) as exc: warning = f"Planner fallback: {exc}. The deterministic document-type outline was used." warnings.append(warning) atomic_write_text(output / "stages" / "01-planner.error.txt", warning + "\n") write_json(output / "stages" / "02-outline.json", outline.to_dict()) atomic_write_text(output / "stages" / "02-outline.md", _render_outline(outline)) writer = create_provider(config.writer) try: response = _invoke(writer, ProviderRequest("draft", drafting_prompt(brief, outline, sources), output), events) except ProviderError as exc: _write_events(output, events) raise PipelineExecutionError(f"writer stage failed: {exc}") from exc atomic_write_text(output / "stages" / "03-writer.raw.txt", response.text + "\n") draft = _clean_markdown_response(response.text) if not draft: raise PipelineExecutionError("writer stage returned no Markdown") rounds: list[RoundResult] = [] for revision_index in range(config.quality_gate.max_revisions + 1): round_number = revision_index + 1 round_dir = output / "rounds" / f"round-{round_number:02d}" round_dir.mkdir(parents=True, exist_ok=True) draft_path = atomic_write_text(round_dir / "draft.md", draft.rstrip() + "\n") lint_report = lint_document(draft, brief, outline, sources) write_json(round_dir / "lint.json", lint_report.to_dict()) atomic_write_text(round_dir / "lint.md", render_lint_markdown(lint_report)) reviews: list[ModelReview] = [] for reviewer_index, reviewer_spec in enumerate(config.reviewers, start=1): provider = create_provider(reviewer_spec.provider) role_slug = _artifact_slug(reviewer_spec.role) prompt = review_prompt(brief, outline, sources, draft, lint_report, reviewer_spec.role) try: review_response = _invoke( provider, ProviderRequest("review", prompt, output, {"role": reviewer_spec.role}), events, ) raw_path = round_dir / f"review-{reviewer_index:02d}-{role_slug}.raw.txt" atomic_write_text(raw_path, review_response.text + "\n") review = ModelReview.from_dict( extract_json_object(review_response.text), role=reviewer_spec.role, provider=review_response.provider, raw_response=review_response.text, ) except (ProviderError, ValidationError) as exc: if config.fail_on_reviewer_error: _write_events(output, events) raise PipelineExecutionError( f"reviewer stage failed ({reviewer_spec.role}/{reviewer_spec.provider.provider}): {exc}" ) from exc warning = f"Reviewer unavailable ({reviewer_spec.role}/{reviewer_spec.provider.provider}): {exc}" warnings.append(warning) review = _failed_review(reviewer_spec.role, reviewer_spec.provider.provider, warning) reviews.append(review) write_json(round_dir / f"review-{reviewer_index:02d}-{role_slug}.json", review.to_dict()) model_mean = sum(review.score for review in reviews) / len(reviews) if reviews else lint_report.score composite = round( lint_report.score * config.quality_gate.deterministic_weight + model_mean * config.quality_gate.model_weight, 1, ) blockers = lint_report.count(Severity.BLOCKER) + sum(review.blocker_count for review in reviews) errors = lint_report.count(Severity.ERROR) + sum( sum(issue.severity == "error" for issue in review.issues) for review in reviews ) passed = ( composite >= config.quality_gate.minimum_score and blockers <= config.quality_gate.max_blockers and errors <= config.quality_gate.max_errors ) round_result = RoundResult( round_number=round_number, draft_path=draft_path, lint_report=lint_report, reviews=reviews, composite_score=composite, blocker_count=blockers, error_count=errors, passed=passed, ) rounds.append(round_result) write_json( round_dir / "quality-gate.json", { "round": round_number, "deterministic_score": lint_report.score, "model_mean_score": round(model_mean, 1), "composite_score": composite, "blockers": blockers, "errors": errors, "passed": passed, }, ) if passed or revision_index >= config.quality_gate.max_revisions: break reviser = create_provider(config.reviser) try: revision_response = _invoke( reviser, ProviderRequest( "revise", revision_prompt(brief, outline, sources, draft, lint_report, reviews), output, {"round": round_number}, ), events, ) except ProviderError as exc: _write_events(output, events) raise PipelineExecutionError(f"revision stage failed after round {round_number}: {exc}") from exc atomic_write_text(round_dir / "revision.raw.txt", revision_response.text + "\n") revised = _clean_markdown_response(revision_response.text) if not revised or revised.strip() == draft.strip(): warnings.append(f"Revision after round {round_number} produced no material change.") draft = revised or draft if not rounds: raise PipelineExecutionError("pipeline produced no quality-gate round") final_round = rounds[-1] final_path = atomic_write_text(output / "final" / "document.md", draft.rstrip() + "\n") report_path = atomic_write_text( output / "final" / "quality-report.md", render_run_report(brief, config, rounds, warnings), ) provenance_path = atomic_write_text( output / "final" / "provenance.md", render_provenance(brief, outline, sources), ) evidence_map_path = write_json( output / "final" / "evidence-map.json", build_evidence_map(brief, outline, sources), ) _write_events(output, events) run_data = { "schema_version": 1, "created_at": utc_now_iso(), "document": brief.title, "document_type": brief.document_type.value, "passed": final_round.passed, "final_score": final_round.composite_score, "rounds": [ { "round": item.round_number, "draft": str(item.draft_path.relative_to(output)), "deterministic_score": item.lint_report.score, "review_scores": {review.role: review.score for review in item.reviews}, "composite_score": item.composite_score, "blockers": item.blocker_count, "errors": item.error_count, "passed": item.passed, } for item in rounds ], "warnings": warnings, "artifacts": { "document": str(final_path.relative_to(output)), "quality_report": str(report_path.relative_to(output)), "provenance": str(provenance_path.relative_to(output)), "evidence_map": str(evidence_map_path.relative_to(output)), "outline": "stages/02-outline.json", "events": "provider-events.jsonl", }, } write_json(output / "run.json", run_data) manifest_path = _write_manifest(output) return RunResult( output_dir=output, final_path=final_path, report_path=report_path, manifest_path=manifest_path, passed=final_round.passed, final_score=final_round.composite_score, rounds=rounds, warnings=warnings, ) def _configured_provider_names(config: PipelineConfig) -> list[str]: specs = [ config.planner, config.writer, config.reviser, *[reviewer.provider for reviewer in config.reviewers], ] return [spec.provider.casefold().strip() for spec in specs if spec.provider.strip()] def _mock_provider_warning(config: PipelineConfig) -> str: provider_names = _configured_provider_names(config) if not provider_names or "mock" not in provider_names: return "" if set(provider_names) == {"mock"}: return ( "All providers are deterministic mocks. This run validates pipeline mechanics only; " "model-review scores are synthetic and must not be used as evidence of document quality." ) return ( "This pipeline mixes external providers with deterministic mocks. Any mock-authored stage " "or mock review score is synthetic; the composite score is not an all-model quality signal." ) def _artifact_slug(value: str) -> str: slug = re.sub(r"[^A-Za-z0-9_-]+", "-", value).strip("-_") return (slug or "reviewer")[:48] def _invoke(provider: Any, request: ProviderRequest, events: list[dict[str, Any]]) -> Any: started = time.perf_counter() event = { "at": utc_now_iso(), "stage": request.stage, "provider": provider.name, "model": provider.spec.model, "metadata": request.metadata, "status": "started", } events.append(event) try: response = provider.generate(request) except Exception as exc: events.append({ **event, "at": utc_now_iso(), "status": "failed", "duration_ms": round((time.perf_counter() - started) * 1000, 1), "error": str(exc), }) raise events.append({ **event, "at": utc_now_iso(), "status": "completed", "duration_ms": round((time.perf_counter() - started) * 1000, 1), "response_characters": len(response.text), "command": response.command, }) return response def _failed_review(role: str, provider: str, message: str) -> ModelReview: return ModelReview( role=role, provider=provider, score=0, dimension_scores={}, issues=[ReviewIssue("document", message, "The independent review did not complete.", "Restore the provider and rerun.", "blocker")], strengths=[], questions=[], raw_response="", ) def _clean_markdown_response(text: str) -> str: stripped = text.strip() full_fence = re.fullmatch(r"```(?:markdown|md)?\s*\n(.*?)\n```", stripped, flags=re.DOTALL | re.IGNORECASE) if full_fence: stripped = full_fence.group(1).strip() return stripped def _render_outline(outline: Outline) -> str: lines = [f"# Outline contract: {outline.title}", ""] for section in outline.sections: lines.extend([ f"## {section.title}", "", f"- Intent: `{section.intent}`", f"- Reader question: {section.reader_question}", f"- Purpose: {section.purpose}", f"- Must include: {', '.join(section.must_include) if section.must_include else '—'}", f"- Evidence IDs: {', '.join(section.evidence_ids) if section.evidence_ids else '—'}", f"- Decision requirements: {', '.join(section.decision_requirements) if section.decision_requirements else '—'}", f"- Transition: {section.transition_to_next or '—'}", "", ]) return "\n".join(lines) def _write_events(output: Path, events: list[dict[str, Any]]) -> None: content = "".join(json.dumps(event, ensure_ascii=False) + "\n" for event in events) atomic_write_text(output / "provider-events.jsonl", content) def _write_manifest(output: Path) -> Path: entries = [] for path in sorted(output.rglob("*")): if not path.is_file() or path.name == "manifest.json": continue entries.append({ "path": str(path.relative_to(output)), "bytes": path.stat().st_size, "sha256": sha256_file(path), }) return write_json( output / "manifest.json", {"schema_version": 1, "created_at": utc_now_iso(), "files": entries}, )