#![cfg(feature = "ai")]
use std::pin::Pin;
use async_trait::async_trait;
use futures_core::Stream;
use futures_util::StreamExt as _;
use rtb_app::app::App;
use rtb_chat::{AiClient, ChatRequest, ChatStreamEvent, Config, Message, Provider};
use secrecy::SecretString;
use crate::error::{DocsError, Result};
pub type AnswerStream = Pin<Box<dyn Stream<Item = String> + Send>>;
#[async_trait]
pub trait AiAnswerStream: Send + Sync + 'static {
async fn ask(&self, context: &str, question: &str) -> Result<AnswerStream>;
}
pub struct AiClientStream {
client: AiClient,
}
impl AiClientStream {
#[must_use]
pub const fn new(client: AiClient) -> Self {
Self { client }
}
}
#[async_trait]
impl AiAnswerStream for AiClientStream {
async fn ask(&self, context: &str, question: &str) -> Result<AnswerStream> {
let req = ChatRequest {
system: Some(format!(
"You are a documentation assistant. Answer using ONLY the doc-tree \
context below. Quote relevant page paths in parentheses.\n\n\
----- DOC TREE -----\n{context}\n----- END DOC TREE -----"
)),
messages: vec![Message::user(question)],
cache_control: true,
..Default::default()
};
let stream = self
.client
.chat_stream(req)
.await
.map_err(|e| DocsError::Assets(format!("ai stream: {e}")))?;
let token_stream = stream.filter_map(|event| async move {
match event {
ChatStreamEvent::Token(t) => Some(t),
_ => None,
}
});
Ok(Box::pin(token_stream))
}
}
pub fn default_answer_stream(_app: &App) -> Result<AiClientStream> {
let api_key = std::env::var("ANTHROPIC_API_KEY").map_err(|_| {
DocsError::Assets(
"docs ask: no `ANTHROPIC_API_KEY` env var set. Set it or wire a \
custom `AiAnswerStream` impl on your tool."
.into(),
)
})?;
let config = Config {
provider: Provider::Anthropic,
model: "claude-opus-4-7".into(),
base_url: None,
api_key: SecretString::from(api_key),
timeout: std::time::Duration::from_secs(60),
allow_insecure_base_url: false,
};
let client = AiClient::new(config).map_err(|e| DocsError::Assets(e.to_string()))?;
Ok(AiClientStream::new(client))
}