From 064b823020c338308439a3179b808de31df57983 Mon Sep 17 00:00:00 2001 From: Yiorgis Gozadinos Date: Wed, 24 Sep 2025 15:42:18 +0300 Subject: [PATCH] Break Evaluation node into AnalyzeInsights & DecisionNode. Use a list of updateable insights & gaps --- docs/agents.md | 12 +- src/haiku/rag/research/common.py | 67 +++++++- src/haiku/rag/research/dependencies.py | 186 +++++++++++++++++++++-- src/haiku/rag/research/graph.py | 8 +- src/haiku/rag/research/models.py | 132 +++++++++++++++- src/haiku/rag/research/nodes/analysis.py | 181 ++++++++++++++++++++++ src/haiku/rag/research/nodes/evaluate.py | 91 ----------- src/haiku/rag/research/nodes/search.py | 4 +- src/haiku/rag/research/prompts.py | 89 ++++++++--- src/haiku/rag/research/state.py | 3 +- src/haiku/rag/research/stream.py | 10 +- tests/test_research_graph_integration.py | 59 +++++-- 12 files changed, 689 insertions(+), 153 deletions(-) create mode 100644 src/haiku/rag/research/nodes/analysis.py delete mode 100644 src/haiku/rag/research/nodes/evaluate.py diff --git a/docs/agents.md b/docs/agents.md index 5e6374ec..59c9f8f8 100644 --- a/docs/agents.md +++ b/docs/agents.md @@ -47,9 +47,10 @@ title: Research graph --- stateDiagram-v2 PlanNode --> SearchDispatchNode - SearchDispatchNode --> EvaluateNode - EvaluateNode --> SearchDispatchNode - EvaluateNode --> SynthesizeNode + SearchDispatchNode --> AnalyzeInsightsNode + AnalyzeInsightsNode --> DecisionNode + DecisionNode --> SearchDispatchNode + DecisionNode --> SynthesizeNode SynthesizeNode --> [*] ``` @@ -57,12 +58,15 @@ Key nodes: - Plan: builds up to 3 standalone sub‑questions (uses an internal presearch tool) - Search (batched): answers sub‑questions using the KB with minimal, verbatim context -- Evaluate: extracts insights, proposes new questions, and checks sufficiency/confidence +- Analyze: aggregates fresh insights, updates gaps, and suggests new sub-questions +- Decision: checks sufficiency/confidence thresholds and chooses whether to iterate - Synthesize: generates a final structured report Primary models: - `SearchAnswer` — one per sub‑question (query, answer, context, sources) +- `InsightRecord` / `GapRecord` — structured tracking of findings and open issues +- `InsightAnalysis` — output of the analysis stage (insights, gaps, commentary) - `EvaluationResult` — insights, new questions, sufficiency, confidence - `ResearchReport` — final report (title, executive summary, findings, conclusions, …) diff --git a/src/haiku/rag/research/common.py b/src/haiku/rag/research/common.py index ca877a31..fddddb99 100644 --- a/src/haiku/rag/research/common.py +++ b/src/haiku/rag/research/common.py @@ -7,6 +7,7 @@ from pydantic_ai.providers.openai import OpenAIProvider from haiku.rag.config import Config from haiku.rag.research.dependencies import ResearchContext +from haiku.rag.research.models import InsightAnalysis if TYPE_CHECKING: # pragma: no cover from haiku.rag.research.state import ResearchDeps, ResearchState @@ -49,7 +50,69 @@ def format_context_for_prompt(context: ResearchContext) -> str: } for qa in context.qa_responses ], - "insights": context.insights, - "gaps": context.gaps, + "insights": [ + { + "id": insight.id, + "summary": insight.summary, + "status": insight.status.value, + "supporting_sources": insight.supporting_sources, + "originating_questions": insight.originating_questions, + "notes": insight.notes, + } + for insight in context.insights + ], + "gaps": [ + { + "id": gap.id, + "description": gap.description, + "severity": gap.severity.value, + "blocking": gap.blocking, + "resolved": gap.resolved, + "resolved_by": gap.resolved_by, + "supporting_sources": gap.supporting_sources, + "notes": gap.notes, + } + for gap in context.gaps + ], } return format_as_xml(context_data, root_tag="research_context") + + +def format_analysis_for_prompt( + analysis: InsightAnalysis | None, +) -> str: + """Format the latest insight analysis as XML for prompts.""" + + if analysis is None: + return "" + + data = { + "commentary": analysis.commentary, + "highlights": [ + { + "id": insight.id, + "summary": insight.summary, + "status": insight.status.value, + "supporting_sources": insight.supporting_sources, + "originating_questions": insight.originating_questions, + "notes": insight.notes, + } + for insight in analysis.highlights + ], + "gap_assessments": [ + { + "id": gap.id, + "description": gap.description, + "severity": gap.severity.value, + "blocking": gap.blocking, + "resolved": gap.resolved, + "resolved_by": gap.resolved_by, + "supporting_sources": gap.supporting_sources, + "notes": gap.notes, + } + for gap in analysis.gap_assessments + ], + "resolved_gaps": analysis.resolved_gaps, + "new_questions": analysis.new_questions, + } + return format_as_xml(data, root_tag="latest_analysis") diff --git a/src/haiku/rag/research/dependencies.py b/src/haiku/rag/research/dependencies.py index 729264e6..2c698fa9 100644 --- a/src/haiku/rag/research/dependencies.py +++ b/src/haiku/rag/research/dependencies.py @@ -1,8 +1,15 @@ +from collections.abc import Iterable + from pydantic import BaseModel, Field from rich.console import Console from haiku.rag.client import HaikuRAG -from haiku.rag.research.models import SearchAnswer +from haiku.rag.research.models import ( + GapRecord, + InsightAnalysis, + InsightRecord, + SearchAnswer, +) from haiku.rag.research.stream import ResearchStream @@ -16,10 +23,10 @@ class ResearchContext(BaseModel): qa_responses: list[SearchAnswer] = Field( default_factory=list, description="Structured QA pairs used during research" ) - insights: list[str] = Field( + insights: list[InsightRecord] = Field( default_factory=list, description="Key insights discovered" ) - gaps: list[str] = Field( + gaps: list[GapRecord] = Field( default_factory=list, description="Identified information gaps" ) @@ -27,15 +34,147 @@ class ResearchContext(BaseModel): """Add a structured QA response (minimal context already included).""" self.qa_responses.append(qa) - def add_insight(self, insight: str) -> None: - """Add a key insight.""" - if insight not in self.insights: - self.insights.append(insight) + def upsert_insights(self, records: Iterable[InsightRecord]) -> list[InsightRecord]: + """Merge one or more insights into the shared context with deduplication.""" - def add_gap(self, gap: str) -> None: - """Identify an information gap.""" - if gap not in self.gaps: - self.gaps.append(gap) + merged: list[InsightRecord] = [] + for record in records: + candidate = InsightRecord.model_validate(record) + existing = next( + (ins for ins in self.insights if ins.id == candidate.id), None + ) + if not existing: + existing = next( + (ins for ins in self.insights if ins.summary == candidate.summary), + None, + ) + + if existing: + existing.summary = candidate.summary + existing.status = candidate.status + if candidate.notes: + existing.notes = candidate.notes + existing.supporting_sources = _merge_unique( + existing.supporting_sources, candidate.supporting_sources + ) + existing.originating_questions = _merge_unique( + existing.originating_questions, candidate.originating_questions + ) + merged.append(existing) + else: + candidate = candidate.model_copy(deep=True) + if candidate.id is None: # pragma: no cover - defensive + raise ValueError( + "InsightRecord.id must be populated after validation" + ) + candidate_id: str = candidate.id + candidate.id = self._allocate_insight_id(candidate_id) + self.insights.append(candidate) + merged.append(candidate) + + return merged + + def upsert_gaps(self, records: Iterable[GapRecord]) -> list[GapRecord]: + """Merge one or more gap records into the shared context with deduplication.""" + + merged: list[GapRecord] = [] + for record in records: + candidate = GapRecord.model_validate(record) + existing = next((gap for gap in self.gaps if gap.id == candidate.id), None) + if not existing: + existing = next( + ( + gap + for gap in self.gaps + if gap.description == candidate.description + ), + None, + ) + + if existing: + existing.description = candidate.description + existing.severity = candidate.severity + existing.blocking = candidate.blocking + existing.resolved = candidate.resolved + if candidate.notes: + existing.notes = candidate.notes + existing.supporting_sources = _merge_unique( + existing.supporting_sources, candidate.supporting_sources + ) + existing.resolved_by = _merge_unique( + existing.resolved_by, candidate.resolved_by + ) + merged.append(existing) + else: + candidate = candidate.model_copy(deep=True) + if candidate.id is None: # pragma: no cover - defensive + raise ValueError("GapRecord.id must be populated after validation") + candidate_id: str = candidate.id + candidate.id = self._allocate_gap_id(candidate_id) + self.gaps.append(candidate) + merged.append(candidate) + + return merged + + def mark_gap_resolved( + self, identifier: str, resolved_by: Iterable[str] | None = None + ) -> GapRecord | None: + """Mark a gap as resolved by identifier (id or description).""" + + gap = self._find_gap(identifier) + if gap is None: + return None + + gap.resolved = True + gap.blocking = False + if resolved_by: + gap.resolved_by = _merge_unique(gap.resolved_by, list(resolved_by)) + return gap + + def integrate_analysis(self, analysis: InsightAnalysis) -> None: + """Apply an analysis result to the shared context.""" + + merged_insights: list[InsightRecord] = [] + if analysis.highlights: + merged_insights = self.upsert_insights(analysis.highlights) + analysis.highlights = merged_insights + if analysis.gap_assessments: + merged_gaps = self.upsert_gaps(analysis.gap_assessments) + analysis.gap_assessments = merged_gaps + if analysis.resolved_gaps: + resolved_by_list = ( + [ins.id for ins in merged_insights if ins.id is not None] + if merged_insights + else None + ) + for resolved in analysis.resolved_gaps: + self.mark_gap_resolved(resolved, resolved_by=resolved_by_list) + for question in analysis.new_questions: + if question not in self.sub_questions: + self.sub_questions.append(question) + + def _allocate_insight_id(self, candidate_id: str) -> str: + taken: set[str] = set() + for ins in self.insights: + if ins.id is not None: + taken.add(ins.id) + return _allocate_sequential_id(candidate_id, taken) + + def _allocate_gap_id(self, candidate_id: str) -> str: + taken: set[str] = set() + for gap in self.gaps: + if gap.id is not None: + taken.add(gap.id) + return _allocate_sequential_id(candidate_id, taken) + + def _find_gap(self, identifier: str) -> GapRecord | None: + normalized = identifier.lower().strip() + for gap in self.gaps: + if gap.id is not None and gap.id == normalized: + return gap + if gap.description.lower().strip() == normalized: + return gap + return None class ResearchDependencies(BaseModel): @@ -49,3 +188,28 @@ class ResearchDependencies(BaseModel): stream: ResearchStream | None = Field( default=None, description="Optional research event stream" ) + + +def _merge_unique(existing: list[str], incoming: Iterable[str]) -> list[str]: + """Merge two iterables preserving order while removing duplicates.""" + + merged = list(existing) + seen = {item for item in existing if item} + for item in incoming: + if item and item not in seen: + merged.append(item) + seen.add(item) + return merged + + +def _allocate_sequential_id(candidate: str, taken: set[str]) -> str: + slug = candidate + if slug not in taken: + return slug + base = slug + counter = 2 + while True: + slug = f"{base}-{counter}" + if slug not in taken: + return slug + counter += 1 diff --git a/src/haiku/rag/research/graph.py b/src/haiku/rag/research/graph.py index 9c8282e9..7f70895f 100644 --- a/src/haiku/rag/research/graph.py +++ b/src/haiku/rag/research/graph.py @@ -1,7 +1,7 @@ from pydantic_graph import Graph from haiku.rag.research.models import ResearchReport -from haiku.rag.research.nodes.evaluate import EvaluateNode +from haiku.rag.research.nodes.analysis import AnalyzeInsightsNode, DecisionNode from haiku.rag.research.nodes.plan import PlanNode from haiku.rag.research.nodes.search import SearchDispatchNode from haiku.rag.research.nodes.synthesize import SynthesizeNode @@ -10,7 +10,8 @@ from haiku.rag.research.state import ResearchDeps, ResearchState __all__ = [ "PlanNode", "SearchDispatchNode", - "EvaluateNode", + "AnalyzeInsightsNode", + "DecisionNode", "SynthesizeNode", "ResearchState", "ResearchDeps", @@ -23,7 +24,8 @@ def build_research_graph() -> Graph[ResearchState, ResearchDeps, ResearchReport] nodes=[ PlanNode, SearchDispatchNode, - EvaluateNode, + AnalyzeInsightsNode, + DecisionNode, SynthesizeNode, ] ) diff --git a/src/haiku/rag/research/models.py b/src/haiku/rag/research/models.py index ab16601d..eedcb835 100644 --- a/src/haiku/rag/research/models.py +++ b/src/haiku/rag/research/models.py @@ -1,4 +1,134 @@ -from pydantic import BaseModel, Field +import re +from enum import Enum + +from pydantic import BaseModel, Field, model_validator + +_SLUG_RE = re.compile(r"[^a-z0-9]+") + + +def _make_slug(text: str, prefix: str) -> str: + """Generate a lowercase slug with the given prefix as fallback.""" + + base = _SLUG_RE.sub("-", text.lower()).strip("-") + if not base: + base = prefix + # Trim overly long slugs but keep enough entropy for readability + return base[:48] + + +class InsightStatus(str, Enum): + OPEN = "open" + VALIDATED = "validated" + TENTATIVE = "tentative" + + +class GapSeverity(str, Enum): + LOW = "low" + MEDIUM = "medium" + HIGH = "high" + + +class InsightRecord(BaseModel): + """Structured insight with provenance and lifecycle metadata.""" + + id: str | None = Field( + default=None, + description="Stable slug identifier for the insight (auto-generated if omitted)", + ) + summary: str = Field(description="Concise description of the insight") + status: InsightStatus = Field( + default=InsightStatus.OPEN, + description="Lifecycle status for the insight", + ) + supporting_sources: list[str] = Field( + default_factory=list, + description="Source identifiers backing the insight", + ) + originating_questions: list[str] = Field( + default_factory=list, + description="Research sub-questions that produced this insight", + ) + notes: str | None = Field( + default=None, + description="Optional elaboration or caveats for the insight", + ) + + @model_validator(mode="after") + def _set_defaults(self) -> "InsightRecord": + if not self.id: + self.id = _make_slug(self.summary, "insight") + self.id = self.id.lower() + self.supporting_sources = list(dict.fromkeys(self.supporting_sources)) + self.originating_questions = list(dict.fromkeys(self.originating_questions)) + return self + + +class GapRecord(BaseModel): + """Structured representation of an identified research gap.""" + + id: str | None = Field( + default=None, + description="Stable slug identifier for the gap (auto-generated if omitted)", + ) + description: str = Field(description="Concrete statement of what is missing") + severity: GapSeverity = Field( + default=GapSeverity.MEDIUM, + description="Severity of the gap for answering the main question", + ) + blocking: bool = Field( + default=True, + description="Whether this gap blocks a confident answer", + ) + resolved: bool = Field( + default=False, + description="Flag indicating if the gap has been resolved", + ) + resolved_by: list[str] = Field( + default_factory=list, + description="Insight IDs or notes explaining how the gap was closed", + ) + supporting_sources: list[str] = Field( + default_factory=list, + description="Sources confirming the gap status (e.g., evidence of absence)", + ) + notes: str | None = Field( + default=None, + description="Optional clarification about the gap or follow-up actions", + ) + + @model_validator(mode="after") + def _set_defaults(self) -> "GapRecord": + if not self.id: + self.id = _make_slug(self.description, "gap") + self.id = self.id.lower() + self.resolved_by = list(dict.fromkeys(self.resolved_by)) + self.supporting_sources = list(dict.fromkeys(self.supporting_sources)) + return self + + +class InsightAnalysis(BaseModel): + """Output of the insight aggregation agent.""" + + highlights: list[InsightRecord] = Field( + default_factory=list, + description="New or updated insights discovered this iteration", + ) + gap_assessments: list[GapRecord] = Field( + default_factory=list, + description="New or updated gap records based on current evidence", + ) + resolved_gaps: list[str] = Field( + default_factory=list, + description="Gap identifiers or descriptions considered resolved", + ) + new_questions: list[str] = Field( + default_factory=list, + max_length=3, + description="Up to three follow-up sub-questions to pursue next", + ) + commentary: str = Field( + description="Short narrative summary of the incremental findings", + ) class ResearchPlan(BaseModel): diff --git a/src/haiku/rag/research/nodes/analysis.py b/src/haiku/rag/research/nodes/analysis.py new file mode 100644 index 00000000..18057089 --- /dev/null +++ b/src/haiku/rag/research/nodes/analysis.py @@ -0,0 +1,181 @@ +from dataclasses import dataclass + +from pydantic_ai import Agent +from pydantic_graph import BaseNode, GraphRunContext + +from haiku.rag.research.common import ( + format_analysis_for_prompt, + format_context_for_prompt, + get_model, + log, +) +from haiku.rag.research.dependencies import ResearchDependencies +from haiku.rag.research.models import EvaluationResult, InsightAnalysis, ResearchReport +from haiku.rag.research.nodes.synthesize import SynthesizeNode +from haiku.rag.research.prompts import DECISION_AGENT_PROMPT, INSIGHT_AGENT_PROMPT +from haiku.rag.research.state import ResearchDeps, ResearchState + + +@dataclass +class AnalyzeInsightsNode(BaseNode[ResearchState, ResearchDeps, ResearchReport]): + provider: str + model: str + + async def run( + self, ctx: GraphRunContext[ResearchState, ResearchDeps] + ) -> BaseNode[ResearchState, ResearchDeps, ResearchReport]: + state = ctx.state + deps = ctx.deps + + log( + deps, + state, + "\n[bold cyan]🧭 Synthesizing new insights and gap status...[/bold cyan]", + ) + + agent = Agent( + model=get_model(self.provider, self.model), + output_type=InsightAnalysis, + instructions=INSIGHT_AGENT_PROMPT, + retries=3, + deps_type=ResearchDependencies, + ) + + context_xml = format_context_for_prompt(state.context) + prompt = ( + "Review the latest research context and update the shared ledger of insights, gaps," + " and follow-up questions.\n\n" + f"{context_xml}" + ) + agent_deps = ResearchDependencies( + client=deps.client, + context=state.context, + console=deps.console, + stream=deps.stream, + ) + result = await agent.run(prompt, deps=agent_deps) + analysis: InsightAnalysis = result.output + + state.context.integrate_analysis(analysis) + state.last_analysis = analysis + + if analysis.commentary: + log(deps, state, f" Summary: {analysis.commentary}") + if analysis.highlights: + log(deps, state, " [bold]Updated insights:[/bold]") + for insight in analysis.highlights: + label = insight.status.value + log( + deps, + state, + f" • ({label}) {insight.summary}", + ) + if analysis.gap_assessments: + log(deps, state, " [bold yellow]Gap updates:[/bold yellow]") + for gap in analysis.gap_assessments: + status = "resolved" if gap.resolved else "open" + severity = gap.severity.value + log( + deps, + state, + f" • ({severity}/{status}) {gap.description}", + ) + if analysis.resolved_gaps: + log(deps, state, " [green]Resolved gaps:[/green]") + for resolved in analysis.resolved_gaps: + log(deps, state, f" • {resolved}") + if analysis.new_questions: + log(deps, state, " [cyan]Proposed follow-ups:[/cyan]") + for question in analysis.new_questions: + log(deps, state, f" • {question}") + + return DecisionNode(self.provider, self.model) + + +@dataclass +class DecisionNode(BaseNode[ResearchState, ResearchDeps, ResearchReport]): + provider: str + model: str + + async def run( + self, ctx: GraphRunContext[ResearchState, ResearchDeps] + ) -> BaseNode[ResearchState, ResearchDeps, ResearchReport]: + state = ctx.state + deps = ctx.deps + + log( + deps, + state, + "\n[bold cyan]📊 Evaluating research sufficiency...[/bold cyan]", + ) + + agent = Agent( + model=get_model(self.provider, self.model), + output_type=EvaluationResult, + instructions=DECISION_AGENT_PROMPT, + retries=3, + deps_type=ResearchDependencies, + ) + + context_xml = format_context_for_prompt(state.context) + analysis_xml = format_analysis_for_prompt(state.last_analysis) + prompt_parts = [ + "Assess whether the research now answers the original question with adequate confidence.", + context_xml, + analysis_xml, + ] + if state.last_eval is not None: + prev = state.last_eval + prompt_parts.append( + "" + f"{prev.confidence_score:.2f}" + f"{str(prev.is_sufficient).lower()}" + f"{prev.reasoning}" + "" + ) + prompt = "\n\n".join(part for part in prompt_parts if part) + + agent_deps = ResearchDependencies( + client=deps.client, + context=state.context, + console=deps.console, + stream=deps.stream, + ) + decision_result = await agent.run(prompt, deps=agent_deps) + output = decision_result.output + + state.last_eval = output + state.iterations += 1 + + for new_q in output.new_questions: + if new_q not in state.context.sub_questions: + state.context.sub_questions.append(new_q) + + if output.key_insights: + log(deps, state, " [bold]Key insights:[/bold]") + for insight in output.key_insights: + log(deps, state, f" • {insight}") + + if output.gaps: + log(deps, state, " [bold yellow]Remaining gaps:[/bold yellow]") + for gap in output.gaps: + log(deps, state, f" • {gap}") + + log( + deps, + state, + f" Confidence: [yellow]{output.confidence_score:.1%}[/yellow]", + ) + status = "[green]Yes[/green]" if output.is_sufficient else "[red]No[/red]" + log(deps, state, f" Sufficient: {status}") + + from haiku.rag.research.nodes.search import SearchDispatchNode + + if ( + output.is_sufficient + and output.confidence_score >= state.confidence_threshold + ) or state.iterations >= state.max_iterations: + log(deps, state, "\n[bold green]✅ Stopping research.[/bold green]") + return SynthesizeNode(self.provider, self.model) + + return SearchDispatchNode(self.provider, self.model) diff --git a/src/haiku/rag/research/nodes/evaluate.py b/src/haiku/rag/research/nodes/evaluate.py deleted file mode 100644 index 6289934e..00000000 --- a/src/haiku/rag/research/nodes/evaluate.py +++ /dev/null @@ -1,91 +0,0 @@ -from dataclasses import dataclass - -from pydantic_ai import Agent -from pydantic_graph import BaseNode, GraphRunContext - -from haiku.rag.research.common import format_context_for_prompt, get_model, log -from haiku.rag.research.dependencies import ( - ResearchDependencies, -) -from haiku.rag.research.models import EvaluationResult, ResearchReport -from haiku.rag.research.nodes.synthesize import SynthesizeNode -from haiku.rag.research.prompts import EVALUATION_AGENT_PROMPT -from haiku.rag.research.state import ResearchDeps, ResearchState - - -@dataclass -class EvaluateNode(BaseNode[ResearchState, ResearchDeps, ResearchReport]): - provider: str - model: str - - async def run( - self, ctx: GraphRunContext[ResearchState, ResearchDeps] - ) -> BaseNode[ResearchState, ResearchDeps, ResearchReport]: - state = ctx.state - deps = ctx.deps - - log( - deps, - state, - "\n[bold cyan]📊 Analyzing and evaluating research progress...[/bold cyan]", - ) - - agent = Agent( - model=get_model(self.provider, self.model), - output_type=EvaluationResult, - instructions=EVALUATION_AGENT_PROMPT, - retries=3, - deps_type=ResearchDependencies, - ) - - context_xml = format_context_for_prompt(state.context) - prompt = ( - "Analyze gathered information and evaluate completeness for the original question.\n\n" - f"{context_xml}" - ) - agent_deps = ResearchDependencies( - client=deps.client, - context=state.context, - console=deps.console, - stream=deps.stream, - ) - eval_result = await agent.run(prompt, deps=agent_deps) - output = eval_result.output - - for insight in output.key_insights: - state.context.add_insight(insight) - for new_q in output.new_questions: - if new_q not in state.context.sub_questions: - state.context.sub_questions.append(new_q) - for gap in output.gaps: - state.context.add_gap(gap) - - state.last_eval = output - state.iterations += 1 - - if output.key_insights: - log(deps, state, " [bold]Key insights:[/bold]") - for ins in output.key_insights: - log(deps, state, f" • {ins}") - if output.gaps: - log(deps, state, " [bold yellow]Remaining gaps:[/bold yellow]") - for gap in output.gaps: - log(deps, state, f" • {gap}") - log( - deps, - state, - f" Confidence: [yellow]{output.confidence_score:.1%}[/yellow]", - ) - status = "[green]Yes[/green]" if output.is_sufficient else "[red]No[/red]" - log(deps, state, f" Sufficient: {status}") - - from haiku.rag.research.nodes.search import SearchDispatchNode - - if ( - output.is_sufficient - and output.confidence_score >= state.confidence_threshold - ) or state.iterations >= state.max_iterations: - log(deps, state, "\n[bold green]✅ Stopping research.[/bold green]") - return SynthesizeNode(self.provider, self.model) - - return SearchDispatchNode(self.provider, self.model) diff --git a/src/haiku/rag/research/nodes/search.py b/src/haiku/rag/research/nodes/search.py index c7e471a9..664a01cc 100644 --- a/src/haiku/rag/research/nodes/search.py +++ b/src/haiku/rag/research/nodes/search.py @@ -25,9 +25,9 @@ class SearchDispatchNode(BaseNode[ResearchState, ResearchDeps, ResearchReport]): state = ctx.state deps = ctx.deps if not state.context.sub_questions: - from haiku.rag.research.nodes.evaluate import EvaluateNode + from haiku.rag.research.nodes.analysis import AnalyzeInsightsNode - return EvaluateNode(self.provider, self.model) + return AnalyzeInsightsNode(self.provider, self.model) # Take up to max_concurrency questions and answer them concurrently take = max(1, state.max_concurrency) diff --git a/src/haiku/rag/research/prompts.py b/src/haiku/rag/research/prompts.py index 29e10424..4f32e5fe 100644 --- a/src/haiku/rag/research/prompts.py +++ b/src/haiku/rag/research/prompts.py @@ -44,38 +44,77 @@ Answering rules: - Prefer concise phrasing; avoid copying long passages. - When evidence is partial, state the limits explicitly in the answer.""" -EVALUATION_AGENT_PROMPT = """You are an analysis and evaluation specialist for -the research workflow. +INSIGHT_AGENT_PROMPT = """You are the insight aggregation specialist for the +research workflow. Inputs available: -- Original research question -- Question–answer pairs produced by search -- Raw search results and source metadata -- Previously identified insights +- Original research question and sub-questions +- Question–answer pairs with supporting snippets and sources +- Existing insights and gaps (with status metadata) -ANALYSIS: -1. Extract the most important, non‑obvious insights from the collected evidence. -2. Identify patterns, agreements, and disagreements across sources. -3. Note material uncertainties and assumptions. +Tasks: +1. Extract new or refined insights that advance understanding of the question. +2. Update gap status, creating new gap entries when necessary and marking + resolved ones explicitly. +3. Suggest up to 3 high-impact follow-up sub_questions that would close the + most important remaining gaps. -EVALUATION: -1. Decide if we have sufficient information to answer the original question. -2. Provide a confidence_score in [0,1] considering: - - Coverage of the main question’s aspects - - Quality, consistency, and diversity of sources - - Depth and specificity of evidence -3. List concrete gaps that still need investigation. -4. Propose up to 3 new sub_questions that would close the highest‑value gaps. +Output format (map directly to fields): +- highlights: list of insights with fields {summary, status, supporting_sources, + originating_questions, notes}. Use status one of {validated, open, tentative}. +- gap_assessments: list of gaps with fields {description, severity, blocking, + resolved, resolved_by, supporting_sources, notes}. Severity must be one of + {low, medium, high}. resolved_by may reference related insight summaries if no + stable identifier yet. +- resolved_gaps: list of identifiers or descriptions for gaps now closed. +- new_questions: up to 3 standalone, specific sub-questions (no duplicates with + existing ones). +- commentary: 1–3 sentences summarizing what changed this round. + +Guidance: +- Be concise and avoid repeating previously recorded information unless it + changed materially. +- Tie supporting_sources to the evidence used; omit if unavailable. +- Only propose new sub_questions that directly address remaining gaps. +- When marking a gap as resolved, ensure the rationale is clear via + resolved_by or notes.""" + +DECISION_AGENT_PROMPT = """You are the research governor responsible for making +stop/go decisions. + +Inputs available: +- Original research question and current plan +- Full insight ledger with status metadata +- Up-to-date gap tracker, including resolved indicators +- Latest insight analysis summary (highlights, gap changes, new questions) +- Previous evaluation decision (if any) + +Tasks: +1. Determine whether the collected evidence now answers the original question. +2. Provide a confidence_score in [0,1] that reflects coverage, evidence quality, + and agreement across sources. +3. List the highest-priority gaps that still block a confident answer. Reference + existing gap descriptions rather than inventing new ones. +4. Optionally propose up to 3 new sub_questions only if they are not already in + the current backlog. Strictness: -- Only mark research as sufficient when all major aspects are addressed with - consistent, reliable evidence and no critical gaps remain. +- Only mark research as sufficient when every critical aspect of the main + question is addressed with reliable, corroborated evidence. +- Treat unresolved high-severity or blocking gaps as a hard stop. -New sub_questions must: -- Be genuinely new (not answered or duplicative; check qa_responses). -- Be standalone and specific (entities, scope, timeframe/region if relevant). -- Be actionable and scoped to the knowledge base (narrow if necessary). -- Be ordered by expected impact (most valuable first).""" +Output fields must line up with EvaluationResult: +- key_insights: concise bullet-ready statements of the most decision-relevant + insights (cite status if helpful). +- new_questions: follow-up sub-questions (max 3) meeting the specificity rules. +- gaps: list remaining blockers; reuse wording from the tracked gaps when + possible to aid downstream reconciliation. +- confidence_score: numeric in [0,1]. +- is_sufficient: true only when no blocking gaps remain. +- reasoning: short narrative tying the decision to evidence coverage. + +Remember: prefer maintaining continuity with the structured context over +introducing new terminology.""" SYNTHESIS_AGENT_PROMPT = """You are a synthesis specialist producing the final research report. diff --git a/src/haiku/rag/research/state.py b/src/haiku/rag/research/state.py index 238accf6..e6df8c41 100644 --- a/src/haiku/rag/research/state.py +++ b/src/haiku/rag/research/state.py @@ -4,7 +4,7 @@ from rich.console import Console from haiku.rag.client import HaikuRAG from haiku.rag.research.dependencies import ResearchContext -from haiku.rag.research.models import EvaluationResult +from haiku.rag.research.models import EvaluationResult, InsightAnalysis from haiku.rag.research.stream import ResearchStream @@ -29,3 +29,4 @@ class ResearchState: max_concurrency: int = 1 confidence_threshold: float = 0.8 last_eval: EvaluationResult | None = None + last_analysis: InsightAnalysis | None = None diff --git a/src/haiku/rag/research/stream.py b/src/haiku/rag/research/stream.py index 9ecf2d80..3c1d56e7 100644 --- a/src/haiku/rag/research/stream.py +++ b/src/haiku/rag/research/stream.py @@ -42,8 +42,14 @@ class ResearchStateSnapshot: confidence_threshold=state.confidence_threshold, pending_sub_questions=len(context.sub_questions), answered_questions=len(context.qa_responses), - insights=list(context.insights), - gaps=list(context.gaps), + insights=[ + f"{insight.status.value}:{insight.summary}" + for insight in context.insights + ], + gaps=[ + f"{gap.severity.value}/{'resolved' if gap.resolved else 'open'}:{gap.description}" + for gap in context.gaps + ], last_confidence=last_confidence, last_sufficient=last_sufficient, ) diff --git a/tests/test_research_graph_integration.py b/tests/test_research_graph_integration.py index d869af7d..3f3aef78 100644 --- a/tests/test_research_graph_integration.py +++ b/tests/test_research_graph_integration.py @@ -4,7 +4,8 @@ import pytest from haiku.rag.research.dependencies import ResearchContext from haiku.rag.research.graph import ( - EvaluateNode, + AnalyzeInsightsNode, + DecisionNode, PlanNode, ResearchDeps, ResearchState, @@ -12,7 +13,16 @@ from haiku.rag.research.graph import ( SynthesizeNode, build_research_graph, ) -from haiku.rag.research.models import EvaluationResult, ResearchReport, SearchAnswer +from haiku.rag.research.models import ( + EvaluationResult, + GapRecord, + GapSeverity, + InsightAnalysis, + InsightRecord, + InsightStatus, + ResearchReport, + SearchAnswer, +) from haiku.rag.research.stream import stream_research_graph @@ -39,7 +49,7 @@ async def test_graph_end_to_end_with_patched_nodes(monkeypatch): return SearchDispatchNode(self.provider, self.model) async def fake_search_dispatch_run(self, ctx) -> Any: - # Answer all pending questions deterministically, then move to evaluation + # Answer all pending questions deterministically, then move to analysis while ctx.state.context.sub_questions: q = ctx.state.context.sub_questions.pop(0) # pydantic BaseModel kwargs not fully typed for pyright @@ -47,20 +57,46 @@ async def test_graph_end_to_end_with_patched_nodes(monkeypatch): SearchAnswer(query=q, answer="A", context=["x"], sources=["s"]) # pyright: ignore[reportCallIssue] ) ctx.deps.emit_log(f"answered:{q}", ctx.state) - return EvaluateNode(self.provider, self.model) + return AnalyzeInsightsNode(self.provider, self.model) - async def fake_evaluate_run(self, ctx) -> Any: - ctx.state.last_eval = EvaluationResult( - key_insights=["ok"], + async def fake_analyze_run(self, ctx) -> Any: + analysis = InsightAnalysis( + highlights=[ + InsightRecord( + summary="haiku.rag orchestrates research stages", + status=InsightStatus.VALIDATED, + supporting_sources=["s"], + originating_questions=["Describe haiku.rag in one sentence"], + ) + ], + gap_assessments=[ + GapRecord( + description="Need a final summary", + severity=GapSeverity.LOW, + blocking=False, + resolved=False, + ) + ], + resolved_gaps=[], new_questions=[], - gaps=["gap"], + commentary="Insights captured for synthesis", + ) + ctx.state.context.integrate_analysis(analysis) + ctx.state.last_analysis = analysis + ctx.deps.emit_log("analysis", ctx.state) + return DecisionNode(self.provider, self.model) + + async def fake_decision_run(self, ctx) -> Any: + ctx.state.last_eval = EvaluationResult( + key_insights=["haiku.rag coordinates planning, search, and synthesis"], + new_questions=[], + gaps=["Need a final summary"], confidence_score=1.0, is_sufficient=True, reasoning="done", ) ctx.state.iterations += 1 - ctx.state.context.add_gap("gap") - ctx.deps.emit_log("evaluated", ctx.state) + ctx.deps.emit_log("decision", ctx.state) return SynthesizeNode(self.provider, self.model) async def fake_synthesize_run(self, ctx) -> Any: @@ -81,7 +117,8 @@ async def test_graph_end_to_end_with_patched_nodes(monkeypatch): monkeypatch.setattr( SearchDispatchNode, "run", fake_search_dispatch_run, raising=False ) - monkeypatch.setattr(EvaluateNode, "run", fake_evaluate_run, raising=False) + monkeypatch.setattr(AnalyzeInsightsNode, "run", fake_analyze_run, raising=False) + monkeypatch.setattr(DecisionNode, "run", fake_decision_run, raising=False) monkeypatch.setattr(SynthesizeNode, "run", fake_synthesize_run, raising=False) start = PlanNode(provider="test", model="test")