2d0d5c1592
新增 7 个示例覆盖全 Phase 公共 API(不依赖 API key 即可运行): - agent_session_demo:AgentBuilder → AgentSession → SessionMemory - custom_tool:BaseTool 注册 + invoke/invoke_all + 权限检查 - prompt_composer:PromptTemplate + PromptComposer + validate_messages - task_agent_demo:JsonPlanParser + Step 状态机 + 错误路径 - conversation_memory_demo:滑动窗口 + 多角色 + 隔离 - knowledge_search_demo:KnowledgeStore + 关键词检索 + 停用词过滤 - streaming_events_demo:submit_stream 事件消费 + 队列耗尽错误路径
117 lines
4.7 KiB
Rust
117 lines
4.7 KiB
Rust
//! streaming_events_demo —— LLM 流式响应事件流消费(含错误路径)。
|
||
//!
|
||
//! 演示:
|
||
//! 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::cycle::{CycleConfig, LlmCycle};
|
||
use agcore::llm::mock::MockProvider;
|
||
use agcore::llm::provider::LlmProvider;
|
||
use agcore::llm::types::message::{ContentBlock, Message};
|
||
use agcore::llm::types::response_v2::{MessageResponse, StopReason, StreamEvent};
|
||
use agcore::llm::types::Usage;
|
||
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<Item = StreamEvent> + 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<dyn LlmProvider> = 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 完成");
|
||
} |