from __future__ import annotations import json from pathlib import Path from typing import Any, Literal from .compat import model_to_dict from .constants import STAGE_IDS, STAGE_LABELS from .models import ( AIReviewCritique, ArtifactStatus, IssueSeverity, IssueStatus, IssueType, LayoutSpec, PipelineState, ProposedChange, ProposedChangeSet, QualityIssue, ReviewModeConfig, ReviewPacket, ReviewQueueItem, ReviewerProductivityMetrics, ReviewRole, RevisionConstraints, SpeakerNotes, SuggestedFix, VersionComparisonSummary, VisualAsset, ) from .quality import ( compute_artifact_diff, compute_deck_health_summary, create_artifact_version, get_current_stage_artifact, grade_stage, now_iso, record_audit, stable_hash, ) REVIEW_MODE_DEFAULTS: dict[ReviewRole, ReviewModeConfig] = { ReviewRole.SME: ReviewModeConfig( role=ReviewRole.SME, label="SME Review", description="Focus on technical correctness, source grounding, and claims.", visible_stage_ids=["text_generation", "technical_review"], focused_issue_types=[ IssueType.UNSUPPORTED_CLAIM, IssueType.TECHNICAL_INACCURACY, ], focused_rubric_dimensions=[ "technical_accuracy", "unsupported_claim_detection", "terminology_consistency", "source_grounding", ], ), ReviewRole.INSTRUCTIONAL_DESIGNER: ReviewModeConfig( role=ReviewRole.INSTRUCTIONAL_DESIGNER, label="Instructional Design Review", description="Focus on objectives, progression, teachability, and cognitive load.", visible_stage_ids=["slide_outline_order", "text_generation", "pedagogical_review"], focused_issue_types=[ IssueType.OBJECTIVE_UNCOVERED, IssueType.OBJECTIVE_WEAKLY_COVERED, IssueType.TEXT_DENSITY_EXCEEDED, IssueType.COGNITIVE_LOAD_HIGH, IssueType.SPEAKER_NOTES_MISSING, ], focused_rubric_dimensions=[ "objective_coverage", "learning_progression", "examples_or_applications", "speaker_notes_teachability", ], ), ReviewRole.VISUAL_DESIGNER: ReviewModeConfig( role=ReviewRole.VISUAL_DESIGNER, label="Visual Design Review", description="Focus on layout, visuals, accessibility, and readability.", visible_stage_ids=[ "image_visual_asset_generation", "aesthetic_ordering_visual_composition", "aesthetic_review", ], focused_issue_types=[ IssueType.LAYOUT_SCHEMA_INVALID, IssueType.LAYOUT_SLOT_VIOLATION, IssueType.VISUAL_ASSET_MISSING, IssueType.VISUAL_ASSET_PURPOSE_MISSING, IssueType.ALT_TEXT_MISSING, IssueType.ACCESSIBILITY_CONTRAST_RISK, ], focused_rubric_dimensions=[ "layout_consistency", "readability", "accessibility", "brand_template_compliance", ], ), ReviewRole.PRODUCER: ReviewModeConfig( role=ReviewRole.PRODUCER, label="Producer Review", description="Focus on blockers, approvals, stale stages, and export readiness.", visible_stage_ids=["final_render_export", "audit_log_version_history"], focused_issue_types=[ IssueType.HUMAN_APPROVAL_MISSING, IssueType.APPROVAL_INVALIDATED, IssueType.STALE_DOWNSTREAM_STAGE, IssueType.EXPORT_PREFLIGHT_FAILED, IssueType.SCORE_BELOW_THRESHOLD, ], focused_rubric_dimensions=[ "preflight_passed", "all_required_approvals_valid", "no_stale_stages", "export_target_safety", ], ), ReviewRole.INSTRUCTOR: ReviewModeConfig( role=ReviewRole.INSTRUCTOR, label="Instructor Review", description="Focus on notes, teachability, timing, confusions, and teaching tips.", visible_stage_ids=["text_generation", "pedagogical_review"], focused_issue_types=[ IssueType.SPEAKER_NOTES_MISSING, IssueType.COGNITIVE_LOAD_HIGH, IssueType.OBJECTIVE_WEAKLY_COVERED, IssueType.TEXT_DENSITY_EXCEEDED, ], focused_rubric_dimensions=[ "speaker_notes_teachability", "examples_or_applications", "learning_progression", ], ), ReviewRole.GENERAL_REVIEWER: ReviewModeConfig( role=ReviewRole.GENERAL_REVIEWER, label="General Review", description="Show all reviewable stages and issue types.", visible_stage_ids=list(STAGE_IDS), focused_issue_types=list(IssueType), focused_rubric_dimensions=[], ), } def _coerce_role(role: ReviewRole | str | None) -> ReviewRole | None: if role in (None, "", "all"): return None if isinstance(role, ReviewRole): return role return ReviewRole(role) def _coerce_severity(severity: IssueSeverity | str | None) -> IssueSeverity | None: if severity in (None, "", "all"): return None if isinstance(severity, IssueSeverity): return severity return IssueSeverity(severity) def _coerce_status(status: IssueStatus | str | None) -> IssueStatus | None: if status in (None, "", "all"): return None if isinstance(status, IssueStatus): return status return IssueStatus(status) def _coerce_issue_type(issue_type: IssueType | str | None) -> IssueType | None: if issue_type in (None, "", "all"): return None if isinstance(issue_type, IssueType): return issue_type return IssueType(issue_type) def _issue_status(issue: QualityIssue) -> IssueStatus: if issue.resolved and issue.status != IssueStatus.RESOLVED: return IssueStatus.RESOLVED return issue.status def _issue_matches_role(issue: QualityIssue, role: ReviewRole | None) -> bool: if role is None or role == ReviewRole.GENERAL_REVIEWER: return True if issue.severity == IssueSeverity.BLOCKER: return True if issue.assigned_role == role: return True config = REVIEW_MODE_DEFAULTS[role] return issue.issue_type in config.focused_issue_types or ( issue.stage_id in config.visible_stage_ids if issue.stage_id else False ) def _stage_sort_key(stage_id: str | None) -> int: if not stage_id or stage_id not in STAGE_IDS: return len(STAGE_IDS) + 1 return STAGE_IDS.index(stage_id) def _slide_sort_key(slide_id: str | None) -> int: if not slide_id: return 10_000 digits = "".join(char for char in slide_id if char.isdigit()) return int(digits) if digits else 10_000 def build_review_queue( state: PipelineState, role: ReviewRole | str | None = None, stage_id: str | None = None, severity: IssueSeverity | str | None = None, status: IssueStatus | str | None = None, issue_type: IssueType | str | None = None, assigned_to: str | None = None, ) -> list[ReviewQueueItem]: review_role = _coerce_role(role) severity_filter = _coerce_severity(severity) status_filter = _coerce_status(status) issue_type_filter = _coerce_issue_type(issue_type) rows: list[ReviewQueueItem] = [] for issue in state.issues.values(): issue_status = _issue_status(issue) if status_filter is None and issue_status == IssueStatus.RESOLVED: continue if status_filter is not None and issue_status != status_filter: continue if stage_id and stage_id != "all" and issue.stage_id != stage_id: continue if severity_filter is not None and issue.severity != severity_filter: continue if issue_type_filter is not None and issue.issue_type != issue_type_filter: continue if assigned_to and issue.assigned_to != assigned_to: continue if not _issue_matches_role(issue, review_role): continue rows.append( ReviewQueueItem( queue_item_id=f"queue:{issue.issue_id}", issue_id=issue.issue_id, issue_type=issue.issue_type, severity=issue.severity, status=issue_status, priority=issue.priority, stage_id=issue.stage_id, slide_id=issue.slide_id, objective_id=issue.objective_id, claim_id=issue.claim_id, artifact_version_id=issue.artifact_version_id, assigned_role=issue.assigned_role, assigned_to=issue.assigned_to, message=issue.message, suggested_fix_summary=issue.suggested_fix, created_at=issue.created_at, ) ) severity_order = { IssueSeverity.BLOCKER: 0, IssueSeverity.MAJOR: 1, IssueSeverity.MINOR: 2, IssueSeverity.INFO: 3, } rows.sort( key=lambda item: ( severity_order[item.severity], _stage_sort_key(item.stage_id), _slide_sort_key(item.slide_id), item.created_at, ) ) record_audit( state, "review_queue_built", metadata={ "review_role": review_role.value if review_role else None, "stage_id": stage_id, "severity": severity_filter.value if severity_filter else None, "status": status_filter.value if status_filter else None, "issue_type": issue_type_filter.value if issue_type_filter else None, "assigned_to": assigned_to, "queue_count": len(rows), }, ) return rows def update_issue_status( issue_id: str, status: IssueStatus | str, state: PipelineState, reviewer_name: str = "human_reviewer", note: str | None = None, ) -> PipelineState: issue = state.issues.get(issue_id) if issue is None: raise ValueError(f"Unknown issue: {issue_id}") new_status = _coerce_status(status) if new_status is None: raise ValueError("Issue status is required.") clean_note = (note or "").strip() if new_status == IssueStatus.WAIVED: return waive_issue(issue_id, clean_note, state, reviewer_name=reviewer_name) if new_status == IssueStatus.WONT_FIX and not clean_note: raise ValueError("Won't fix requires a reason.") if new_status == IssueStatus.RESOLVED and not clean_note: raise ValueError("Resolving an issue requires a resolution note or re-grade evidence.") issue.status = new_status issue.resolved = new_status == IssueStatus.RESOLVED if new_status == IssueStatus.RESOLVED: issue.resolved_by = reviewer_name or "human_reviewer" issue.resolved_at = now_iso() issue.resolution_note = clean_note elif new_status == IssueStatus.WONT_FIX: issue.resolution_note = clean_note else: issue.resolved_at = None record_audit( state, "issue_marked_wont_fix" if new_status == IssueStatus.WONT_FIX else "issue_status_updated", stage_id=issue.stage_id, slide_id=issue.slide_id, issue_id=issue.issue_id, reason=clean_note or new_status.value, metadata={"status": new_status.value, "reviewer_name": reviewer_name}, ) return state def assign_issue( issue_id: str, role: ReviewRole | str | None, assigned_to: str | None, state: PipelineState, ) -> PipelineState: issue = state.issues.get(issue_id) if issue is None: raise ValueError(f"Unknown issue: {issue_id}") review_role = _coerce_role(role) issue.assigned_role = review_role issue.assigned_to = assigned_to or None record_audit( state, "issue_assigned", stage_id=issue.stage_id, slide_id=issue.slide_id, issue_id=issue.issue_id, metadata={ "review_role": review_role.value if review_role else None, "assigned_to": issue.assigned_to, }, ) return state def waive_issue( issue_id: str, reason: str, state: PipelineState, reviewer_name: str = "human_reviewer", ) -> PipelineState: issue = state.issues.get(issue_id) if issue is None: raise ValueError(f"Unknown issue: {issue_id}") clean_reason = (reason or "").strip() if issue.severity == IssueSeverity.BLOCKER: raise ValueError("BLOCKER issues cannot be waived in P1.") if issue.severity in {IssueSeverity.MAJOR, IssueSeverity.MINOR} and not clean_reason: raise ValueError("Waiving a non-blocker issue requires a reason.") issue.status = IssueStatus.WAIVED issue.resolved = False issue.waiver_reason = clean_reason issue.waived_by = reviewer_name or "human_reviewer" issue.waived_at = now_iso() record_audit( state, "issue_waived", stage_id=issue.stage_id, slide_id=issue.slide_id, issue_id=issue.issue_id, reason=clean_reason, metadata={"reviewer_name": reviewer_name}, ) return state def _suggested_fix_fields(issue: QualityIssue) -> tuple[str, str, dict[str, int], str]: mapping: dict[IssueType, tuple[str, str, dict[str, int], str]] = { IssueType.OBJECTIVE_UNCOVERED: ( "improve_objective_mapping", "Map this objective to an existing slide or add a focused slide.", {"objective_coverage": 25}, "medium", ), IssueType.OBJECTIVE_WEAKLY_COVERED: ( "add_example", "Add an example or strengthen the slide's objective alignment.", {"objective_coverage": 15}, "low", ), IssueType.UNSUPPORTED_CLAIM: ( "support_claim", "Attach source support, mark for human review, or remove the unsupported claim.", {"source_grounding": 35}, "high", ), IssueType.TECHNICAL_INACCURACY: ( "mark_for_human_review", "Route this item to SME review for correction.", {"technical_accuracy": 30}, "high", ), IssueType.TEXT_DENSITY_EXCEEDED: ( "reduce_text", "Reduce visible text or move details into speaker notes.", {"text_density": 20}, "low", ), IssueType.COGNITIVE_LOAD_HIGH: ( "rewrite_text", "Simplify the slide to reduce objectives or concepts.", {"learning_progression": 15}, "medium", ), IssueType.SPEAKER_NOTES_MISSING: ( "add_speaker_notes", "Add instructor-facing speaker notes.", {"speaker_notes_teachability": 20}, "low", ), IssueType.LAYOUT_SCHEMA_INVALID: ( "change_layout", "Switch to an approved layout and named slots.", {"layout_consistency": 30}, "medium", ), IssueType.LAYOUT_SLOT_VIOLATION: ( "change_layout", "Fix slot assignments to match the selected template.", {"slot_compliance": 20}, "medium", ), IssueType.VISUAL_ASSET_MISSING: ( "add_visual_asset", "Add a visual asset or mark the visual as not required.", {"asset_completeness": 30}, "medium", ), IssueType.VISUAL_ASSET_PURPOSE_MISSING: ( "add_visual_asset", "Define whether the visual is instructional or decorative.", {"visual_relevance": 15}, "low", ), IssueType.ALT_TEXT_MISSING: ( "add_alt_text", "Add concise alt text for the meaningful visual.", {"accessibility": 20}, "low", ), IssueType.ACCESSIBILITY_CONTRAST_RISK: ( "mark_for_human_review", "Ask visual review to adjust foreground/background contrast.", {"accessibility": 20}, "medium", ), IssueType.COPYRIGHT_OR_LICENSE_RISK: ( "mark_for_human_review", "Review or replace the asset license source.", {"license_or_generation_metadata": 20}, "medium", ), IssueType.STALE_DOWNSTREAM_STAGE: ( "other", "Re-run and re-grade this stale downstream stage.", {"no_stale_stages": 25}, "low", ), IssueType.APPROVAL_INVALIDATED: ( "other", "Re-grade and approve the current artifact.", {"all_required_approvals_valid": 25}, "low", ), } return mapping.get( issue.issue_type, ("other", issue.suggested_fix or "Review this issue and decide the next action.", {}, "medium"), ) def generate_suggested_fix_for_issue(issue_id: str, state: PipelineState) -> SuggestedFix: issue = state.issues.get(issue_id) if issue is None: raise ValueError(f"Unknown issue: {issue_id}") fix_type, description, score_delta, risk = _suggested_fix_fields(issue) fix = SuggestedFix( fix_id=f"fix_{stable_hash([issue_id, issue.message])[:12]}", issue_ids=[issue.issue_id], stage_id=issue.stage_id, slide_ids=[issue.slide_id] if issue.slide_id else [], objective_ids=[issue.objective_id] if issue.objective_id else [], claim_ids=[issue.claim_id] if issue.claim_id else [], fix_type=fix_type, # type: ignore[arg-type] description=description, expected_score_delta=score_delta, risk_level=risk, # type: ignore[arg-type] requires_human_confirmation=True, ) state.suggested_fixes[fix.fix_id] = fix issue.suggested_fix = fix.description record_audit( state, "suggested_fix_created", stage_id=issue.stage_id, slide_id=issue.slide_id, issue_id=issue.issue_id, metadata={"fix_id": fix.fix_id, "fix_type": fix.fix_type}, ) return fix def generate_suggested_fixes(state: PipelineState) -> list[SuggestedFix]: fixes = [ generate_suggested_fix_for_issue(issue.issue_id, state) for issue in state.issues.values() if _issue_status(issue) != IssueStatus.RESOLVED ] return list({fix.fix_id: fix for fix in fixes}.values()) def critique_artifact_for_improvement( stage_id: str, state: PipelineState, constraints: RevisionConstraints | None = None, ) -> AIReviewCritique: current = get_current_stage_artifact(stage_id, state) relevant = [ issue for issue in state.issues.values() if issue.stage_id == stage_id and _issue_status(issue) != IssueStatus.RESOLVED ] blocker_count = sum(issue.severity == IssueSeverity.BLOCKER for issue in relevant) major_count = sum(issue.severity == IssueSeverity.MAJOR for issue in relevant) diagnosis = ( f"{STAGE_LABELS.get(stage_id, stage_id)} has {blocker_count} blocker(s) " f"and {major_count} major issue(s) needing review." ) if not relevant: diagnosis = f"{STAGE_LABELS.get(stage_id, stage_id)} has no open stage-specific issues." risks = [] if any(issue.issue_type == IssueType.UNSUPPORTED_CLAIM for issue in relevant): risks.append("Do not alter technical claim text without SME confirmation.") if constraints and constraints.preserve_objective_mappings: risks.append("Objective mappings must remain unchanged.") critique = AIReviewCritique( critique_id=f"critique_{len(state.critiques) + 1:05d}", stage_id=stage_id, artifact_version_id=current.artifact_version_id if current else None, created_at=now_iso(), issue_ids=[issue.issue_id for issue in relevant], diagnosis=diagnosis, risks=risks, reviewer_questions=[ "Which low-risk changes should be applied now?", "Which items need human subject-matter review?", ], ) state.critiques[critique.critique_id] = critique record_audit( state, "artifact_critiqued", stage_id=stage_id, artifact_version_id=critique.artifact_version_id, metadata={"critique_id": critique.critique_id, "issue_count": len(relevant)}, ) return critique def _change_field(change: ProposedChange) -> str: if change.field_path: return change.field_path.split(".")[-1] return change.target_type def _constraint_violation(change: ProposedChange, constraints: RevisionConstraints) -> str | None: field = _change_field(change) fields = {field, change.field_path or "", change.target_type} if constraints.preserve_slide_order and (field in {"slide_order", "slide_number"}): return "preserve_slide_order blocks slide order changes" if constraints.preserve_slide_titles and change.target_type == "slide_title": return "preserve_slide_titles blocks title changes" if constraints.preserve_technical_claims and change.target_type == "claim" and field == "claim_text": return "preserve_technical_claims blocks claim text changes" if constraints.preserve_objective_mappings and change.target_type == "objective_mapping": return "preserve_objective_mappings blocks objective mapping changes" if constraints.preserve_visual_assets and change.target_type == "visual_asset": return "preserve_visual_assets blocks visual asset changes" if constraints.preserve_layout_ids and change.target_type == "layout" and field == "layout_id": return "preserve_layout_ids blocks layout ID changes" if constraints.selected_slide_ids and change.target_id not in constraints.selected_slide_ids: return "selected_slide_ids limits changes to selected slides" if constraints.blocked_fields and fields & set(constraints.blocked_fields): return "blocked_fields contains this change field" if constraints.allowed_fields and not fields & set(constraints.allowed_fields): return "allowed_fields does not include this change field" return None def validate_change_set_against_constraints( change_set: ProposedChangeSet, constraints: RevisionConstraints, ) -> list[QualityIssue]: issues: list[QualityIssue] = [] for change in change_set.changes: reason = _constraint_violation(change, constraints) if not reason: continue issues.append( QualityIssue( issue_id=f"issue:revision_constraint:{change.change_id}", issue_type=IssueType.RENDER_SAFETY_VIOLATION, severity=IssueSeverity.BLOCKER, message=f"Change {change.change_id} violates revision constraints: {reason}.", stage_id=change_set.stage_id, slide_id=change.target_id if change.target_id.startswith("slide_") else None, artifact_version_id=change_set.artifact_version_id, suggested_fix="Remove this proposed change or relax the constraint.", priority="critical", created_at=now_iso(), ) ) return issues def _build_change_for_issue( change_id: str, issue: QualityIssue, state: PipelineState, ) -> ProposedChange: slide = state.slides.get(issue.slide_id or "") if issue.issue_type in {IssueType.TEXT_DENSITY_EXCEEDED, IssueType.COGNITIVE_LOAD_HIGH} and slide: after = " ".join(slide.visible_text.split()[:30]) or slide.title or "Focused slide text." return ProposedChange( change_id=change_id, target_type="slide_text", target_id=slide.slide_id, field_path="visible_text", before=slide.visible_text, after=f"{after}\n\nDetails moved to speaker notes.", rationale="Reduce visible text while preserving the instructional point.", issue_ids=[issue.issue_id], risk_level="low", ) if issue.issue_type == IssueType.SPEAKER_NOTES_MISSING and slide: return ProposedChange( change_id=change_id, target_type="speaker_notes", target_id=slide.slide_id, field_path="notes_text", before=slide.speaker_notes.notes_text if slide.speaker_notes else None, after=f"Teach {slide.title or slide.slide_id} with one example and one check for understanding.", rationale="Add instructor-facing notes required for teachability.", issue_ids=[issue.issue_id], risk_level="low", ) if issue.issue_type in {IssueType.LAYOUT_SCHEMA_INVALID, IssueType.LAYOUT_SLOT_VIOLATION} and slide: layout = state.layout_specs.get(slide.slide_id) return ProposedChange( change_id=change_id, target_type="layout", target_id=slide.slide_id, field_path="layout_id", before=layout.layout_id if layout else None, after="title_body", rationale="Switch to an approved layout with named slots.", issue_ids=[issue.issue_id], risk_level="medium", ) if issue.issue_type == IssueType.ALT_TEXT_MISSING and slide: asset = next((asset for asset in state.visual_assets.values() if asset.slide_id == slide.slide_id), None) return ProposedChange( change_id=change_id, target_type="visual_asset", target_id=slide.slide_id, field_path="alt_text", before=asset.alt_text if asset else None, after=f"Instructional visual supporting {slide.title or slide.slide_id}.", rationale="Add concise alt text for accessibility.", issue_ids=[issue.issue_id], risk_level="low", ) if issue.issue_type == IssueType.VISUAL_ASSET_MISSING and slide: return ProposedChange( change_id=change_id, target_type="visual_asset", target_id=slide.slide_id, field_path="prompt", before=None, after=f"Create a simple instructional diagram for {slide.title or slide.slide_id}.", rationale="Provide a visual brief without generating external assets.", issue_ids=[issue.issue_id], risk_level="medium", ) if issue.issue_type == IssueType.UNSUPPORTED_CLAIM and issue.claim_id: claim = state.claims.get(issue.claim_id) return ProposedChange( change_id=change_id, target_type="claim", target_id=issue.claim_id, field_path="review_status", before=claim.review_status if claim else None, after="needs_human_review", rationale="Preserve claim text while blocking export until support is reviewed.", issue_ids=[issue.issue_id], risk_level="high", ) if issue.issue_type in {IssueType.OBJECTIVE_UNCOVERED, IssueType.OBJECTIVE_WEAKLY_COVERED}: target_slide_id = issue.slide_id or next(iter(state.slides), "slide_1") return ProposedChange( change_id=change_id, target_type="objective_mapping", target_id=target_slide_id, field_path="objective_ids", before=state.slides[target_slide_id].objective_ids if target_slide_id in state.slides else [], after=[issue.objective_id] if issue.objective_id else [], rationale="Improve objective coverage by mapping an appropriate slide.", issue_ids=[issue.issue_id], risk_level="medium", ) return ProposedChange( change_id=change_id, target_type="metadata", target_id=issue.stage_id or "deck", field_path="review_instruction", before=None, after=issue.suggested_fix or issue.message, rationale="Convert the issue into explicit reviewer guidance.", issue_ids=[issue.issue_id], risk_level="medium", ) def create_proposed_change_set( stage_id: str, state: PipelineState, critique: AIReviewCritique | None = None, constraints: RevisionConstraints | None = None, ) -> ProposedChangeSet: current = get_current_stage_artifact(stage_id, state) active_constraints = constraints or RevisionConstraints() critique = critique or critique_artifact_for_improvement(stage_id, state, active_constraints) change_set_id = f"changes_{len(state.proposed_change_sets) + 1:05d}" changes: list[ProposedChange] = [] for index, issue_id in enumerate(critique.issue_ids, start=1): issue = state.issues.get(issue_id) if issue is None: continue change = _build_change_for_issue(f"{change_set_id}_chg_{index}", issue, state) if _constraint_violation(change, active_constraints): record_audit( state, "revision_constraints_applied", stage_id=stage_id, slide_id=change.target_id if change.target_id.startswith("slide_") else None, issue_id=issue.issue_id, metadata={"change_id": change.change_id, "excluded": True}, ) continue changes.append(change) change_set = ProposedChangeSet( change_set_id=change_set_id, stage_id=stage_id, artifact_version_id=current.artifact_version_id if current else None, created_at=now_iso(), constraints=active_constraints, changes=changes, summary=f"{len(changes)} proposed change(s) for {STAGE_LABELS.get(stage_id, stage_id)}.", ) state.proposed_change_sets[change_set.change_set_id] = change_set critique.proposed_change_set_id = change_set.change_set_id record_audit( state, "proposed_change_set_created", stage_id=stage_id, artifact_version_id=change_set.artifact_version_id, metadata={"change_set_id": change_set.change_set_id, "change_count": len(changes)}, ) return change_set def _find_or_create_asset(state: PipelineState, slide_id: str) -> VisualAsset: for asset in state.visual_assets.values(): if asset.slide_id == slide_id: return asset asset_id = f"asset_{len(state.visual_assets) + 1}" asset = VisualAsset( asset_id=asset_id, slide_id=slide_id, asset_type="placeholder", source="mock", license_status="generated", ) state.visual_assets[asset_id] = asset return asset def _apply_change_to_state(state: PipelineState, change: ProposedChange) -> None: if change.target_type == "slide_title" and change.target_id in state.slides: state.slides[change.target_id].title = str(change.after or "") return if change.target_type == "slide_text" and change.target_id in state.slides: state.slides[change.target_id].visible_text = str(change.after or "") return if change.target_type == "speaker_notes" and change.target_id in state.slides: slide = state.slides[change.target_id] if slide.speaker_notes is None: slide.speaker_notes = SpeakerNotes(slide_id=slide.slide_id) slide.speaker_notes.notes_text = str(change.after or "") return if change.target_type == "layout" and change.target_id in state.slides: slide = state.slides[change.target_id] layout = state.layout_specs.get(slide.slide_id) if layout is None: layout = LayoutSpec(slide_id=slide.slide_id, layout_id="title_body") state.layout_specs[slide.slide_id] = layout if change.field_path == "layout_id": layout.layout_id = str(change.after or "title_body") layout.slot_assignments = { "title": slide.title or "", "body": slide.visible_text, } return if change.target_type == "visual_asset" and change.target_id in state.slides: asset = _find_or_create_asset(state, change.target_id) if change.field_path == "alt_text": asset.alt_text = str(change.after or "") elif change.field_path == "prompt": asset.prompt = str(change.after or "") asset.purpose = "instructional" asset.approved_for_export = False return if change.target_type == "claim" and change.target_id in state.claims: claim = state.claims[change.target_id] if change.field_path == "review_status": claim.review_status = str(change.after or "needs_human_review") # type: ignore[assignment] elif change.field_path == "claim_text": claim.claim_text = str(change.after or "") return if change.target_type == "objective_mapping" and change.target_id in state.slides: after = change.after if isinstance(change.after, list) else [] state.slides[change.target_id].objective_ids = [str(item) for item in after] def apply_proposed_change_set( change_set_id: str, state: PipelineState, selected_change_ids: list[str] | None = None, reviewer_name: str = "human_reviewer", ) -> PipelineState: change_set = state.proposed_change_sets.get(change_set_id) if change_set is None: raise ValueError(f"Unknown proposed change set: {change_set_id}") selected = set(selected_change_ids or [change.change_id for change in change_set.changes]) applied: list[ProposedChange] = [] for change in change_set.changes: if change.change_id not in selected: continue reason = _constraint_violation(change, change_set.constraints) if reason: record_audit( state, "revision_constraints_applied", stage_id=change_set.stage_id, slide_id=change.target_id if change.target_id.startswith("slide_") else None, reason=reason, metadata={"change_set_id": change_set_id, "change_id": change.change_id}, ) continue _apply_change_to_state(state, change) applied.append(change) if not applied: change_set.status = "rejected" record_audit( state, "proposed_change_set_rejected", stage_id=change_set.stage_id, artifact_version_id=change_set.artifact_version_id, reason="No changes passed revision constraints.", metadata={"change_set_id": change_set_id}, ) return state content = { "applied_change_set_id": change_set_id, "applied_by": reviewer_name, "changes": [model_to_dict(change) for change in applied], "slides": [model_to_dict(slide) for slide in state.slides.values()], "claims": [model_to_dict(claim) for claim in state.claims.values()], "visual_assets": [model_to_dict(asset) for asset in state.visual_assets.values()], "layout_specs": [model_to_dict(layout) for layout in state.layout_specs.values()], } artifact = create_artifact_version( state, change_set.stage_id, content, created_by="ai", status=ArtifactStatus.CANDIDATE, mark_downstream_stale=True, ) change_set.status = "applied" if len(applied) == len(change_set.changes) else "partially_applied" record_audit( state, "proposed_change_set_applied", stage_id=change_set.stage_id, artifact_version_id=artifact.artifact_version_id, metadata={ "change_set_id": change_set_id, "applied_change_ids": [change.change_id for change in applied], "reviewer_name": reviewer_name, }, ) return state def reject_proposed_change_set( change_set_id: str, state: PipelineState, reviewer_name: str = "human_reviewer", reason: str | None = None, ) -> PipelineState: change_set = state.proposed_change_sets.get(change_set_id) if change_set is None: raise ValueError(f"Unknown proposed change set: {change_set_id}") change_set.status = "rejected" record_audit( state, "proposed_change_set_rejected", stage_id=change_set.stage_id, artifact_version_id=change_set.artifact_version_id, reason=reason, metadata={"change_set_id": change_set_id, "reviewer_name": reviewer_name}, ) return state def improve_selected_slides( stage_id: str, slide_ids: list[str], state: PipelineState, constraints: RevisionConstraints | None = None, action: Literal[ "rewrite_title", "rewrite_visible_text", "rewrite_speaker_notes", "improve_visual_brief", "suggest_layout_switch", "add_alt_text", ] = "rewrite_visible_text", ) -> ProposedChangeSet: active_constraints = constraints or RevisionConstraints(selected_slide_ids=slide_ids) active_constraints.selected_slide_ids = slide_ids change_set_id = f"changes_{len(state.proposed_change_sets) + 1:05d}" changes: list[ProposedChange] = [] for index, slide_id in enumerate(slide_ids, start=1): slide = state.slides.get(slide_id) if slide is None: continue before: Any after: Any target_type: str field_path: str if action == "rewrite_title": target_type, field_path = "slide_title", "title" before, after = slide.title, f"{slide.title or slide.slide_id} (review draft)" elif action == "rewrite_speaker_notes": target_type, field_path = "speaker_notes", "notes_text" before = slide.speaker_notes.notes_text if slide.speaker_notes else None after = f"Teach {slide.title or slide_id} with a concise explanation and a learner check." elif action == "improve_visual_brief": target_type, field_path = "visual_asset", "prompt" before = next((asset.prompt for asset in state.visual_assets.values() if asset.slide_id == slide_id), None) after = f"Instructional visual brief for {slide.title or slide_id}." elif action == "suggest_layout_switch": target_type, field_path = "layout", "layout_id" before = state.layout_specs[slide_id].layout_id if slide_id in state.layout_specs else None after = "title_body" elif action == "add_alt_text": target_type, field_path = "visual_asset", "alt_text" before = next((asset.alt_text for asset in state.visual_assets.values() if asset.slide_id == slide_id), None) after = f"Visual explaining {slide.title or slide_id}." else: target_type, field_path = "slide_text", "visible_text" before, after = slide.visible_text, f"{slide.visible_text}\n\nReview draft: clearer and shorter." change = ProposedChange( change_id=f"{change_set_id}_chg_{index}", target_type=target_type, # type: ignore[arg-type] target_id=slide_id, field_path=field_path, before=before, after=after, rationale=f"Targeted {action.replace('_', ' ')} for selected slide.", risk_level="low" if action != "suggest_layout_switch" else "medium", ) if not _constraint_violation(change, active_constraints): changes.append(change) current = get_current_stage_artifact(stage_id, state) change_set = ProposedChangeSet( change_set_id=change_set_id, stage_id=stage_id, artifact_version_id=current.artifact_version_id if current else None, created_at=now_iso(), constraints=active_constraints, changes=changes, summary=f"Targeted {action.replace('_', ' ')} for {len(changes)} slide(s).", ) state.proposed_change_sets[change_set.change_set_id] = change_set record_audit( state, "targeted_slide_improvement_requested", stage_id=stage_id, metadata={ "change_set_id": change_set.change_set_id, "slide_ids": slide_ids, "action": action, }, ) record_audit( state, "proposed_change_set_created", stage_id=stage_id, metadata={"change_set_id": change_set.change_set_id, "change_count": len(changes)}, ) return change_set def _items_by_id(items: Any, id_field: str) -> dict[str, dict[str, Any]]: if isinstance(items, dict): values = list(items.values()) elif isinstance(items, list): values = items else: values = [] result: dict[str, dict[str, Any]] = {} for item in values: if not isinstance(item, dict): continue item_id = item.get(id_field) if isinstance(item_id, str): result[item_id] = item return result def _artifact_content(state: PipelineState, artifact_version_id: str) -> dict[str, Any]: artifact = state.artifacts.get(artifact_version_id) if artifact is None: raise ValueError(f"Unknown artifact version: {artifact_version_id}") content = artifact.metadata.get("content", {}) return content if isinstance(content, dict) else {"value": content} def compare_artifact_versions_semantically( previous_artifact_version_id: str, current_artifact_version_id: str, state: PipelineState, ) -> VersionComparisonSummary: previous_content = _artifact_content(state, previous_artifact_version_id) current_content = _artifact_content(state, current_artifact_version_id) previous_slides = _items_by_id(previous_content.get("slides"), "slide_id") current_slides = _items_by_id(current_content.get("slides"), "slide_id") previous_claims = _items_by_id(previous_content.get("claims"), "claim_id") current_claims = _items_by_id(current_content.get("claims"), "claim_id") previous_layouts = _items_by_id(previous_content.get("layout_specs"), "slide_id") current_layouts = _items_by_id(current_content.get("layout_specs"), "slide_id") previous_assets = _items_by_id(previous_content.get("visual_assets"), "asset_id") current_assets = _items_by_id(current_content.get("visual_assets"), "asset_id") added = sorted(set(current_slides) - set(previous_slides)) removed = sorted(set(previous_slides) - set(current_slides)) changed_titles: list[str] = [] changed_slide_ids: set[str] = set(added + removed) changed_notes: list[str] = [] changed_objectives: list[str] = [] for slide_id in sorted(set(previous_slides) & set(current_slides)): previous = previous_slides[slide_id] current = current_slides[slide_id] if previous.get("title") != current.get("title"): changed_titles.append(slide_id) changed_slide_ids.add(slide_id) if previous.get("visible_text") != current.get("visible_text"): changed_slide_ids.add(slide_id) if previous.get("speaker_notes") != current.get("speaker_notes"): changed_notes.append(slide_id) changed_slide_ids.add(slide_id) if previous.get("objective_ids") != current.get("objective_ids"): changed_objectives.append(slide_id) changed_slide_ids.add(slide_id) changed_claims = sorted( claim_id for claim_id in set(previous_claims) | set(current_claims) if previous_claims.get(claim_id) != current_claims.get(claim_id) ) changed_layouts = sorted( slide_id for slide_id in set(previous_layouts) | set(current_layouts) if previous_layouts.get(slide_id) != current_layouts.get(slide_id) ) changed_assets = sorted( asset_id for asset_id in set(previous_assets) | set(current_assets) if previous_assets.get(asset_id) != current_assets.get(asset_id) ) summary = VersionComparisonSummary( comparison_id=f"compare_{len(state.version_comparisons) + 1:05d}", previous_artifact_version_id=previous_artifact_version_id, current_artifact_version_id=current_artifact_version_id, created_at=now_iso(), changed_slide_ids=sorted(changed_slide_ids), added_slide_ids=added, removed_slide_ids=removed, changed_titles=changed_titles, changed_objective_mappings=changed_objectives, changed_claims=changed_claims, changed_layouts=changed_layouts, changed_visual_assets=changed_assets, changed_speaker_notes=changed_notes, summary=( f"{len(changed_slide_ids)} changed slide(s), {len(changed_claims)} changed claim(s), " f"{len(changed_layouts)} changed layout(s)." ), diff_text=compute_artifact_diff(previous_content, current_content), ) state.version_comparisons[summary.comparison_id] = summary record_audit( state, "semantic_diff_created", artifact_version_id=current_artifact_version_id, metadata={ "comparison_id": summary.comparison_id, "previous_artifact_version_id": previous_artifact_version_id, "changed_slide_ids": summary.changed_slide_ids, }, ) return summary def _packet_issues( state: PipelineState, role: ReviewRole | None, stage_ids: list[str], slide_ids: list[str], ) -> list[QualityIssue]: issues = [] for issue in state.issues.values(): if stage_ids and issue.stage_id not in stage_ids: continue if slide_ids and issue.slide_id not in slide_ids: continue if not _issue_matches_role(issue, role): continue issues.append(issue) return issues def _render_markdown_packet( state: PipelineState, role: ReviewRole | None, stage_ids: list[str], slide_ids: list[str], issues: list[QualityIssue], ) -> str: health = compute_deck_health_summary(state) config = REVIEW_MODE_DEFAULTS.get(role) if role else None blockers = [issue for issue in issues if issue.severity == IssueSeverity.BLOCKER] majors = [issue for issue in issues if issue.severity == IssueSeverity.MAJOR] waived = [issue for issue in issues if _issue_status(issue) == IssueStatus.WAIVED] fixes = generate_suggested_fixes(state) latest_comparison = next(reversed(state.version_comparisons.values()), None) lines = [ f"# Review Packet: {state.deck_title or state.job_id}", "", "## Job Summary", f"- Job ID: {state.job_id}", f"- Deck title: {state.deck_title or ''}", f"- Role: {config.label if config else 'All reviewers'}", "", "## Deck Health", f"- Can export: {health.can_export}", f"- Unresolved blockers: {health.unresolved_blocker_count}", f"- Unresolved major issues: {health.unresolved_major_issue_count}", f"- Stale stages: {health.stale_stage_count}", "", "## Role-Specific Review Focus", config.description if config else "Full review packet across all roles.", "", "## Relevant Stages", ", ".join(stage_ids or STAGE_IDS), "", "## Relevant Slides", ", ".join(slide_ids or sorted(state.slides)) or "No slides yet.", "", "## Open Blockers", *[f"- {issue.issue_id}: {issue.message}" for issue in blockers], "", "## Major Issues", *[f"- {issue.issue_id}: {issue.message}" for issue in majors], "", "## Waived Issues", *[f"- {issue.issue_id}: {issue.waiver_reason or issue.message}" for issue in waived], "", "## Suggested Fixes", *[f"- {fix.fix_id}: {fix.description}" for fix in fixes if set(fix.issue_ids) & {issue.issue_id for issue in issues}], "", "## Objective Coverage Summary", *[ f"- {trace.objective_id}: {trace.coverage_status} ({trace.coverage_score})" for trace in state.objective_traces.values() ], "", "## Claim Support Summary", *[ f"- {claim.claim_id}: {claim.review_status}" for claim in state.claims.values() if not slide_ids or claim.slide_id in slide_ids ], "", "## Version Comparison Summary", latest_comparison.summary if latest_comparison else "No semantic comparison available.", "", "## Reviewer Notes", *[f"- {note.reviewer_name}: {note.summary or ''}" for note in state.reviewer_notes], "", ] return "\n".join(lines) def create_review_packet( state: PipelineState, role: ReviewRole | str | None = None, stage_ids: list[str] | None = None, slide_ids: list[str] | None = None, format: Literal["markdown", "json"] = "markdown", ) -> ReviewPacket: review_role = _coerce_role(role) selected_stage_ids = stage_ids or [] selected_slide_ids = slide_ids or [] issues = _packet_issues(state, review_role, selected_stage_ids, selected_slide_ids) packet_id = f"packet_{len(state.review_packets) + 1:05d}" packet_dir = Path("artifacts") / state.job_id packet_dir.mkdir(parents=True, exist_ok=True) suffix = "md" if format == "markdown" else "json" path = packet_dir / f"{packet_id}.{suffix}" if format == "markdown": payload = _render_markdown_packet(state, review_role, selected_stage_ids, selected_slide_ids, issues) path.write_text(payload, encoding="utf-8") else: payload_data = { "job": {"job_id": state.job_id, "deck_title": state.deck_title}, "role": review_role.value if review_role else None, "stage_ids": selected_stage_ids, "slide_ids": selected_slide_ids, "issues": [model_to_dict(issue) for issue in issues], "suggested_fixes": [model_to_dict(fix) for fix in generate_suggested_fixes(state)], "objective_traces": [model_to_dict(trace) for trace in state.objective_traces.values()], "claims": [model_to_dict(claim) for claim in state.claims.values()], "version_comparisons": [model_to_dict(item) for item in state.version_comparisons.values()], "reviewer_notes": [model_to_dict(note) for note in state.reviewer_notes], } path.write_text(json.dumps(payload_data, indent=2, sort_keys=True), encoding="utf-8") packet = ReviewPacket( packet_id=packet_id, job_id=state.job_id, role=review_role, stage_ids=selected_stage_ids, slide_ids=selected_slide_ids, issue_ids=[issue.issue_id for issue in issues], objective_ids=list(state.objectives), artifact_version_ids=[ artifact.artifact_version_id for artifact in state.artifacts.values() if artifact.is_current ], created_at=now_iso(), format=format, path=str(path), summary=f"{format.upper()} review packet with {len(issues)} issue(s).", ) state.review_packets[packet.packet_id] = packet record_audit( state, "review_packet_created", metadata={ "packet_id": packet.packet_id, "review_role": review_role.value if review_role else None, "format": format, "path": packet.path, }, ) return packet def list_prompt_runs_for_stage(stage_id: str, state: PipelineState): runs = [ run for run in state.prompt_runs.values() if stage_id in {"", "all"} or run.stage_id == stage_id ] runs.sort(key=lambda run: run.created_at) record_audit( state, "prompt_run_viewed", stage_id=None if stage_id in {"", "all"} else stage_id, metadata={"prompt_run_count": len(runs)}, ) return runs def restore_artifact_version_as_candidate( artifact_version_id: str, state: PipelineState, reviewer_name: str = "human_reviewer", ) -> PipelineState: artifact = state.artifacts.get(artifact_version_id) if artifact is None: raise ValueError(f"Unknown artifact version: {artifact_version_id}") if artifact.status in {ArtifactStatus.STALE, ArtifactStatus.INVALIDATED, ArtifactStatus.EXPORTED}: raise ValueError(f"Cannot restore {artifact.status.value} artifacts as candidates.") content = artifact.metadata.get("content", {}) restored = create_artifact_version( state, artifact.stage_id, { "restored_from_artifact_version_id": artifact.artifact_version_id, "content": content, }, created_by="human", status=ArtifactStatus.CANDIDATE, mark_downstream_stale=True, ) restored.parent_artifact_version_ids = [artifact.artifact_version_id] record_audit( state, "artifact_version_restored_as_candidate", stage_id=artifact.stage_id, artifact_version_id=restored.artifact_version_id, metadata={ "source_artifact_version_id": artifact_version_id, "reviewer_name": reviewer_name, }, ) return state def compute_reviewer_productivity_metrics(state: PipelineState) -> ReviewerProductivityMetrics: issues = list(state.issues.values()) issues_by_stage: dict[str, int] = {} issues_by_role: dict[str, int] = {} for issue in issues: if issue.stage_id: issues_by_stage[issue.stage_id] = issues_by_stage.get(issue.stage_id, 0) + 1 role = issue.assigned_role.value if issue.assigned_role else "unassigned" issues_by_role[role] = issues_by_role.get(role, 0) + 1 average_score_by_stage = { stage_id: grade.score for stage_id, stage in state.stages.items() if (grade := stage.grade_result) is not None } metrics = ReviewerProductivityMetrics( job_id=state.job_id, total_issues=len(issues), open_issues=sum(_issue_status(issue) in {IssueStatus.OPEN, IssueStatus.ACKNOWLEDGED, IssueStatus.IN_PROGRESS} for issue in issues), resolved_issues=sum(_issue_status(issue) == IssueStatus.RESOLVED for issue in issues), waived_issues=sum(_issue_status(issue) == IssueStatus.WAIVED for issue in issues), blocker_count=sum(issue.severity == IssueSeverity.BLOCKER for issue in issues), major_count=sum(issue.severity == IssueSeverity.MAJOR for issue in issues), issues_by_stage=issues_by_stage, issues_by_role=issues_by_role, average_score_by_stage=average_score_by_stage, candidate_versions_created=sum( artifact.status == ArtifactStatus.CANDIDATE for artifact in state.artifacts.values() ), approvals_invalidated=sum( approval.approval_status == "invalidated" for approval in state.approvals ), stale_events=sum(event.event_type == "stage_marked_stale" for event in state.audit_events), generated_at=now_iso(), ) record_audit( state, "reviewer_metrics_computed", metadata={"total_issues": metrics.total_issues, "open_issues": metrics.open_issues}, ) return metrics def role_review_mode_summary(role: ReviewRole | str | None, state: PipelineState) -> str: review_role = _coerce_role(role) or ReviewRole.GENERAL_REVIEWER config = REVIEW_MODE_DEFAULTS[review_role] record_audit( state, "role_review_mode_selected", metadata={"review_role": review_role.value}, ) return ( f"{config.label}: {config.description or ''}\n" f"Stages: {', '.join(config.visible_stage_ids)}\n" f"Rubrics: {', '.join(config.focused_rubric_dimensions) or 'All'}" ) def grade_and_refresh_review_state(stage_id: str, state: PipelineState) -> None: grade_stage(stage_id, state) generate_suggested_fixes(state) compute_deck_health_summary(state)