doubleo7/deep_research/src/core.rs

202 lines
8.7 KiB
Rust
Raw Normal View History

use crate::progress::Spinner;
use crate::review::{self, Review};
use crate::tools::{FetchPage, SearchWeb};
use clap::Parser;
use futures::StreamExt;
use rig::agent::MultiTurnStreamItem;
use rig::client::{AgentClientExt, Nothing};
use rig::message::Text;
use rig::providers::ollama;
use rig::streaming::{StreamedAssistantContent, StreamingPrompt};
use std::io::Write;
/// The tool-calling research loop needs to reliably decide what to search
/// for, when a page is worth fetching, and when it has enough evidence —
/// that's a reasoning-heavy job best given to the largest local Gemma
/// variant. Turning the gathered notes into prose afterwards is comparatively
/// mechanical, so the smaller/faster variant handles that pass instead.
const RESEARCHER_MODEL: &str = "gemma4:26b";
const WRITER_MODEL: &str = "gemma4-e4b:latest";
const MAX_RESEARCH_TURNS: usize = 12;
/// Research/review rounds before giving up and writing the report from
/// whatever the last pass produced, rather than looping forever on a topic
/// the reviewer can never be satisfied with.
const MAX_RESEARCH_ROUNDS: usize = 3;
/// Deep research agentic loop over local Gemma models: a tool-calling agent
/// gathers and cross-checks web evidence, a reviewer agent gates it, and a
/// writer agent turns approved findings into a structured report.
#[derive(Parser)]
#[command(name = "doubleo7-research", version, about)]
pub(crate) struct Cli {
/// Research topic to investigate
pub(crate) topic: Option<String>,
/// Emit logs at this level (off by default; passing this also enables a
/// progress spinner to switch off, since the logs already show progress)
#[arg(short = 'l', long, value_name = "LEVEL")]
pub(crate) log_level: Option<tracing::Level>,
}
/// Only initializes a subscriber (and thus produces any log output at all)
/// when the caller opted in via `--log-level` — otherwise tracing's macros
/// are no-ops, leaving the terminal clean for the spinner.
pub(crate) fn initialize_observability(log_level: tracing::Level) {
tracing_subscriber::fmt()
.with_env_filter(
tracing_subscriber::EnvFilter::try_from_default_env()
.unwrap_or_else(|_| tracing_subscriber::EnvFilter::new(log_level.to_string())),
)
.with_span_events(tracing_subscriber::fmt::format::FmtSpan::CLOSE)
.with_writer(std::io::stderr)
.init();
}
/// The least-agentic shape that fits: a plain Rust loop putting *this code*,
/// not a model, in charge of when to stop — re-running research with the
/// reviewer's feedback folded in until it approves or the round budget runs
/// out, then writing the report from whatever the last pass produced.
pub(crate) async fn research(topic: &str, show_progress: bool) -> anyhow::Result<String> {
let client = ollama::Client::new(Nothing)?;
let mut findings = String::new();
let mut feedback: Option<Review> = None;
for round in 1..=MAX_RESEARCH_ROUNDS {
findings = gather_findings(&client, topic, feedback.as_ref(), round, show_progress).await?;
let review = review::review_findings(&client, topic, &findings, show_progress).await?;
let approved = review.approved;
tracing::info!(round, approved, "review verdict");
if approved || round == MAX_RESEARCH_ROUNDS {
break;
}
feedback = Some(review);
}
write_report(&client, topic, &findings, show_progress).await
}
/// Wraps the tool-calling research loop in its own span so it's visible as a
/// single unit in traces, distinct from the writing and review phases and
/// nesting rig's own per-turn `chat`/`execute_tool` spans underneath it.
#[tracing::instrument(skip(client, feedback), fields(gen_ai.agent.name = "researcher"))]
async fn gather_findings(
client: &ollama::Client,
topic: &str,
feedback: Option<&Review>,
round: usize,
show_progress: bool,
) -> anyhow::Result<String> {
let current_date = chrono::offset::Local::now().to_string();
let researcher = client
.agent(RESEARCHER_MODEL)
.name("researcher")
.preamble(format!(
"You are a meticulous research assistant. Use the search_web and fetch_page tools to \
investigate the user's topic: run several searches with varied phrasing, fetch the \
most promising pages, and cross-check claims across at least two sources before \
trusting them. Ensure sources are up to date: the current date is {current_date}. \
Once you are confident you have enough evidence, stop calling tools and reply with a \
plain-text dump of every fact you gathered, using footnote-style citations: write each \
fact followed by a bracketed number like [1], then at the end of your reply list a \
'Sources' section mapping each number to the exact URL it came from, one per line, e.g. \
'[1] https://example.com/page'. Reuse the same number when multiple facts come from the \
same URL do not give one URL two different numbers. Also call out any open questions \
or contradictions between sources, citing the footnotes involved. This is raw research \
material for a writer, not a final report, so favor completeness over polish.")
.as_str(),
)
.tool(SearchWeb)
.tool(FetchPage)
.build();
let task = match feedback {
None => topic.to_string(),
Some(review) => format!(
"Topic: {topic}\n\n\
You already ran a research pass on this topic. A reviewer checked it against its \
cited sources and found it insufficient. Do more research to address the reviewer's \
feedback, then produce an updated findings dump: carry forward what's solid, and \
add, correct, or better-source whatever the gaps call for. Gaps include out-of-date,\
irrelevant, or clearly wrong information. The date is {current_date}.\n\n\
Solid findings from the last pass keep and build on these:\n{}\n\n\
Gaps the reviewer found conclusions not actually backed by their source, sources \
that don't line up with the conclusion drawn from them, or parts of the topic still \
uncovered:\n{}",
review.solid_findings, review.gaps
),
};
let spinner = Spinner::start(show_progress, "Researching...");
let findings = researcher
.runner(task)
.max_turns(MAX_RESEARCH_TURNS)
.run()
.await?
.output;
drop(spinner);
tracing::info!(round, findings = %findings, "research phase complete");
Ok(findings)
}
#[tracing::instrument(skip(client, findings), fields(gen_ai.agent.name = "writer"))]
async fn write_report(
client: &ollama::Client,
topic: &str,
findings: &str,
show_progress: bool,
) -> anyhow::Result<String> {
let writer = client
.agent(WRITER_MODEL)
.name("writer")
.preamble(
"You turn raw research notes into a clear, well-organized report for the reader. The \
notes use footnote-style citations a bracketed number like [1] after a fact, with a \
Sources section mapping numbers to URLs. Preserve this scheme in your report: keep the \
same [n] markers next to the claims they support (renumbering only if you drop unused \
sources), and end the report with a 'Sources' section listing every footnote number \
still in use next to its exact URL. Structure the body with headings, and call out any \
open questions or contradictions the research turned up. Do not invent facts beyond \
what the notes provide.",
)
.build();
// Drop the spinner before streaming starts: report text is about to print
// to the same terminal line, so the two must not race over stdout.
let spinner = Spinner::start(show_progress, "Writing report...");
let mut response_stream = writer
.stream_prompt(format!("Topic: {topic}\n\nResearch notes:\n{findings}"))
.await;
drop(spinner);
let mut report = String::new();
while let Some(chunk) = response_stream.next().await {
match chunk? {
MultiTurnStreamItem::StreamAssistantItem(StreamedAssistantContent::Text(Text {
text,
..
})) => {
report.push_str(&text);
// Terminal stdout is line-buffered, so a flush is needed here —
// otherwise a chunk without a trailing newline sits in the
// buffer instead of appearing as it streams in.
print!("{text}");
std::io::stdout().flush()?;
}
_ => continue,
}
}
println!();
Ok(report)
}