use std::sync::Arc;
use langchainrust::agents::{
serve_agent_sse_on, AgentStreamFactory, StreamingFunctionCallingAgent,
};
use langchainrust::{OpenAIChat, OpenAIConfig};
use tokio::net::TcpListener;
fn real_llm() -> OpenAIChat {
OpenAIChat::new(OpenAIConfig {
api_key: std::env::var("AGENT_SSE_API_KEY")
.unwrap_or_else(|_| "sk-6eb65fcf5d17491ca10b984efe1f43e7".to_string()),
base_url: std::env::var("AGENT_SSE_BASE_URL").unwrap_or_else(|_| {
"https://llm-8xo1b7o30z27y2xc.cn-beijing.maas.aliyuncs.com/compatible-mode/v1"
.to_string()
}),
model: std::env::var("AGENT_SSE_MODEL")
.unwrap_or_else(|_| "qwen3.7-max-2026-06-08".to_string()),
streaming: true,
temperature: Some(0.3),
max_tokens: Some(1024),
..Default::default()
})
}
#[tokio::main]
async fn main() {
let host = std::env::var("AGENT_SSE_HOST").unwrap_or_else(|_| "127.0.0.1".to_string());
let port: u16 = std::env::var("AGENT_SSE_PORT")
.ok()
.and_then(|v| v.parse().ok())
.unwrap_or(8090);
let agent = Arc::new(StreamingFunctionCallingAgent::new(real_llm()));
let factory: Arc<dyn AgentStreamFactory> = Arc::new(move |input: String| {
let agent = agent.clone();
async move { agent.invoke_stream(input).await }
});
let listener = TcpListener::bind((host.as_str(), port))
.await
.unwrap_or_else(|e| {
eprintln!("failed to bind {host}:{port}: {e}");
std::process::exit(1);
});
let bound = listener.local_addr().unwrap();
println!("Streaming agent SSE server started ✅");
println!(" POST http://{bound}/agent/stream body: {{\"input\":\"...\"}}");
println!(" GET http://{bound}/agent/stream?input=... (native EventSource)");
println!("events: text* → final_answer | error, then done; 20s keep-alive");
println!("press Ctrl+C to stop.");
serve_agent_sse_on(factory, listener)
.await
.unwrap_or_else(|e| {
eprintln!("SSE service exited with error: {e}");
std::process::exit(1);
});
}