//! streaming_events_demo —— LLM 流式响应事件流消费(含错误路径)。 //! Required features: cargo run --example streaming_events_demo --features "llm,provider-openai" //! //! 演示: //! 1. `MockProvider::chat_stream` 输出标准 `StreamEvent` 流(离线可跑) //! 2. `LlmCycle::submit_stream` 消费流 //! 3. match 各类 `StreamEvent`:MessageStart / ContentBlockStart / TextDelta / //! ContentBlockEnd / CostUpdate / MessageComplete //! 4. 实时累计文本与解析事件计数 //! 5. 从 `MessageComplete` 拿到完整 `MessageResponse` //! 6. **错误路径**:队列耗尽时 `chat_stream` 返回 `Err`,`submit_stream` 同样 //! 返回 `Err`(错误不进流,直接 fail-fast) //! //! 运行:`cargo run --example streaming_events_demo` use std::sync::Arc; use agcore::llm::LlmProvider; use agcore::llm::cycle::{CycleConfig, LlmCycle}; use agcore::llm::mock::MockProvider; use agcore::llm::types::Usage; use agcore::llm::types::message::{ContentBlock, Message}; use agcore::llm::types::response_v2::{MessageResponse, StopReason, StreamEvent}; use futures_util::StreamExt; /// 构造预设的纯文本响应。 fn text_response(id: &str, text: &str) -> MessageResponse { MessageResponse { id: id.to_string(), model: "mock".into(), message: Message::Assistant { content: vec![ContentBlock::Text { text: text.into() }], }, usage: Usage::from_input_output(3, text.chars().count() as u32), stop_reason: StopReason::Stop, extra: Default::default(), } } /// 消费一个流直到终止事件,并打印事件轨迹。 async fn drain_stream(stream: &mut (impl futures_util::Stream + Unpin)) { while let Some(event) = stream.next().await { match event { StreamEvent::MessageStart { id, model } => { println!("[MessageStart] id={id}, model={model}"); } StreamEvent::ContentBlockStart { index, .. } => { println!("[ContentBlockStart] index={index}"); } StreamEvent::TextDelta { text } => { print!("[TextDelta] {text}"); use std::io::Write; std::io::stdout().flush().ok(); } StreamEvent::ContentBlockEnd { index } => { println!("\n[ContentBlockEnd] index={index}"); } StreamEvent::CostUpdate { usage } => { let p = usage.prompt_tokens.unwrap_or(0); let c = usage.completion_tokens.unwrap_or(0); println!("[CostUpdate] prompt={p}, completion={c}"); } StreamEvent::MessageComplete { full_response } => { println!( "[MessageComplete] stop_reason={:?}", full_response.stop_reason ); break; } other => { // 未匹配的事件变体(ThinkingDelta / RefusalDelta / ToolCall* 等) // 在 MockProvider 当前实现下不可达;保留打印以便未来 Provider 扩展时易调试。 eprintln!("[unhandled event] {other:?}"); } } } } #[tokio::main] async fn main() { // 1. 准备 Mock provider(仅含一轮文本响应,第二次调用会触发错误) let provider = Arc::new(MockProvider::new(vec![text_response( "resp-001", "今天天气晴朗,适合户外活动。", )])); let dyn_provider: Arc = provider.clone(); // ===== 阶段 1:正常流式响应 ===== println!("=== 阶段 1:正常流式响应 ==="); let mut cycle = LlmCycle::new_with_arc(dyn_provider.clone(), CycleConfig::default()); let mut stream = cycle .submit_stream("讲个笑话".to_string(), vec![]) .await .expect("阶段 1 submit_stream 应成功"); drain_stream(&mut stream).await; // ===== 阶段 2:错误路径(队列耗尽)===== // 队列中已无响应 → `chat_stream` 返回 `Err(LlmError::Other)` → // `submit_stream` 用 `?` 立即传播(错误不进流,fail-fast)。 // 上层 Agent 通过 `match` 或 `?` 处理 `AgentError::Llm(_)`。 println!("\n=== 阶段 2:错误路径(队列耗尽)==="); let mut cycle = LlmCycle::new_with_arc(dyn_provider, CycleConfig::default()); let result = cycle.submit_stream("第二次提问".to_string(), vec![]).await; match result { Ok(_) => panic!("阶段 2 必须失败(队列耗尽)"), Err(e) => { eprintln!("[expected error] {e}"); assert!( e.to_string().contains("预设响应已用完"), "应包含 MockProvider 的队列耗尽提示" ); } } println!("\n✓ streaming_events_demo 完成"); }