c709b3d9b2
为 agent-cli 增加 stdio JSON-RPC 2.0 服务入口与运行控制协议 补充 Host 实时流监听和流式取消入口 新增 App Server 子进程闭环测试及协议文档 更新 README、架构、TODO 与共享决策记录
55 lines
1.8 KiB
Rust
55 lines
1.8 KiB
Rust
use std::sync::{Arc, Mutex};
|
|
|
|
use agent_host::AgentHost;
|
|
use agent_provider_fake::{FakeProvider, FakeStep};
|
|
use agent_runtime_core::ProviderStreamEvent;
|
|
use agent_runtime_engine::{Cancellation, EngineStreamEvent};
|
|
|
|
#[test]
|
|
fn stream_callback_receives_delta_with_run_id_before_return() {
|
|
let seen = Arc::new(Mutex::new(Vec::<(String, String)>::new()));
|
|
let seen_clone = seen.clone();
|
|
let host = AgentHost::in_memory()
|
|
.unwrap()
|
|
.with_provider(
|
|
Arc::new(FakeProvider::new([FakeStep::stream_text(["你", "好"])])),
|
|
"fake",
|
|
)
|
|
.with_stream_callback(move |run_id: &str, event: &EngineStreamEvent| {
|
|
if let ProviderStreamEvent::TextDelta { delta, .. } = event.event() {
|
|
seen_clone
|
|
.lock()
|
|
.unwrap()
|
|
.push((run_id.to_owned(), delta.clone()));
|
|
}
|
|
});
|
|
|
|
let output = host.run_streaming("问候").unwrap();
|
|
let events = seen.lock().unwrap().clone();
|
|
assert!(!events.is_empty());
|
|
assert!(events.iter().all(|(run_id, _)| run_id == &output.run_id));
|
|
assert_eq!(
|
|
events
|
|
.iter()
|
|
.map(|(_, delta)| delta.as_str())
|
|
.collect::<String>(),
|
|
"你好"
|
|
);
|
|
assert_eq!(output.output.text, "你好");
|
|
}
|
|
|
|
#[test]
|
|
fn cancelled_stream_does_not_start_provider_or_tool() {
|
|
let provider = Arc::new(FakeProvider::new([FakeStep::stream_text(["不会执行"])]));
|
|
let cancellation = Cancellation::new();
|
|
cancellation.cancel();
|
|
let host = AgentHost::in_memory()
|
|
.unwrap()
|
|
.with_provider(provider.clone(), "fake");
|
|
let handle = host.prepare_run("取消").unwrap();
|
|
|
|
let result = host.run_existing_streaming_with_cancellation(&handle.run_id, cancellation);
|
|
assert!(result.is_err());
|
|
assert_eq!(provider.remaining_steps(), 1);
|
|
}
|