feat: add update indicators and stabilize Cursor streams

This commit is contained in:
leookun
2026-08-24 05:24:07 +08:00
parent 0f23a9a9c3
commit 6170778de9
21 changed files with 360 additions and 105 deletions
+48
View File
@@ -146,6 +146,54 @@ async fn registry_shutdown_cancels_runs_and_closes_run_sse_outputs() {
assert_eq!(output.recv().await, None);
}
#[tokio::test]
async fn client_heartbeat_returns_a_server_protocol_heartbeat() {
let (_directory, store) = fixtures::temp_store().await;
let assets = PromptAssets::load(
std::path::Path::new(env!("CARGO_MANIFEST_DIR"))
.join("prompt/cursor")
.as_path(),
)
.unwrap();
let registry = CursorSessionRegistry::new(
store,
Arc::new(fake_provider::FakeProvider::default()),
PromptCompiler::new(assets),
Default::default(),
);
let handle = registry.get_or_create("heartbeat-run").await.unwrap();
let mut output = handle.subscribe();
handle
.command(CursorCommand::Append {
seqno: 0,
message: Box::new(pb::AgentClientMessage {
message: Some(pb::agent_client_message::Message::ClientHeartbeat(
pb::ClientHeartbeat {},
)),
}),
})
.await
.unwrap();
let frame = tokio::time::timeout(std::time::Duration::from_secs(1), output.recv())
.await
.unwrap()
.unwrap();
let (_, payload) = connect::decode_frames(&frame).unwrap().pop().unwrap();
let message = pb::AgentServerMessage::decode(payload).unwrap();
assert!(matches!(
message.message,
Some(pb::agent_server_message::Message::InteractionUpdate(
pb::InteractionUpdate {
message: Some(pb::interaction_update::Message::Heartbeat(_)),
}
))
));
registry.shutdown().await;
}
#[tokio::test]
async fn runtime_user_message_action_aborts_active_exec_before_canceled_end_stream() {
let (_directory, store) = fixtures::temp_store().await;
+31
View File
@@ -329,6 +329,37 @@ async fn openai_responses_raw_stream_does_not_invent_reasoning_effort() {
assert_eq!(replayed, ["opaque-1", "opaque-2"]);
}
#[tokio::test]
async fn openai_responses_streams_openrouter_reasoning_text_events() {
let (base_url, _requests, server) = fixture_server(
"/v1/responses",
concat!(
"data: {\"type\":\"response.reasoning_text.delta\",\"delta\":\"still working\"}\n\n",
"data: {\"type\":\"response.reasoning_text.done\"}\n\n",
"data: {\"type\":\"response.completed\",\"response\":{}}\n\n",
),
)
.await;
let provider = OpenAiResponsesProvider::new(
reqwest::Client::new(),
config(ProviderKind::OpenAiResponses, base_url, None),
);
let events = collect(provider.stream(invocation(), CancellationToken::new())).await;
server.abort();
assert!(events
.iter()
.any(|event| matches!(event, ModelEvent::ThinkingStart)));
assert!(events.iter().any(
|event| matches!(event, ModelEvent::ThinkingDelta(delta) if delta == "still working")
));
assert!(events
.iter()
.any(|event| matches!(event, ModelEvent::ThinkingEnd)));
assert_eq!(events.last(), Some(&ModelEvent::Done(FinishReason::Stop)));
}
#[tokio::test]
async fn openai_responses_reasoning_item_done_closes_an_open_summary() {
let (base_url, _requests, server) = fixture_server(