mirror of
https://wget.la/https://github.com/leookun/cursor-byok
synced 2026-10-04 02:40:50 +08:00
Merge pull request #385 from kevin9327/fix/runtime-user-message-injection-leak
fix(conversation): clear pending runtime user-message injections
This commit is contained in:
@@ -431,21 +431,28 @@ impl ConversationOutput {
|
||||
streams.clear();
|
||||
}
|
||||
if let CommitCause::RuntimeEvent { event_id } = &state.cause {
|
||||
if let Some(injection_id) = event_id.strip_prefix("inject-context:") {
|
||||
if let Some(pending) = self.pending_injections.remove(injection_id)
|
||||
{
|
||||
let delivered_at_ms = crate::cursor::tools::runtime::now_ms()
|
||||
.min(i64::MAX as u64)
|
||||
as i64;
|
||||
self.handle.emit(&events::context_injection_delivered(
|
||||
injection_id.to_owned(),
|
||||
pending.delivery_batch_id.clone(),
|
||||
delivered_at_ms,
|
||||
))?;
|
||||
if let Some(user_message) = pending.user_message {
|
||||
self.handle
|
||||
.emit(&events::user_message_appended(user_message))?;
|
||||
}
|
||||
// Injections key `pending_injections` by their raw
|
||||
// injection id and commit under `inject-context:{id}`,
|
||||
// while runtime user messages key it by (and commit
|
||||
// under) the full `user-message:{id}` event id. Strip
|
||||
// the injection prefix when present and otherwise use
|
||||
// the event id verbatim so both are cleared and emit
|
||||
// their delivered/appended events.
|
||||
let injection_id = event_id
|
||||
.strip_prefix("inject-context:")
|
||||
.unwrap_or(event_id.as_str());
|
||||
if let Some(pending) = self.pending_injections.remove(injection_id) {
|
||||
let delivered_at_ms = crate::cursor::tools::runtime::now_ms()
|
||||
.min(i64::MAX as u64)
|
||||
as i64;
|
||||
self.handle.emit(&events::context_injection_delivered(
|
||||
injection_id.to_owned(),
|
||||
pending.delivery_batch_id.clone(),
|
||||
delivered_at_ms,
|
||||
))?;
|
||||
if let Some(user_message) = pending.user_message {
|
||||
self.handle
|
||||
.emit(&events::user_message_appended(user_message))?;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -506,6 +506,122 @@ async fn runtime_user_message_action_interrupts_and_continues_with_new_message()
|
||||
assert!(history.contains("queued follow-up"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn runtime_user_message_reports_delivered_and_appended() {
|
||||
// A runtime user message is queued into `pending_injections` under a
|
||||
// `user-message:{id}` key, but the commit correlation only handled the
|
||||
// `inject-context:` prefix, so the entry was never cleared: the client
|
||||
// never saw Delivered/UserMessageAppended and every later tool round was
|
||||
// detached (a hang). This asserts the full delivery sequence.
|
||||
let (_directory, store) = fixtures::temp_store().await;
|
||||
let provider = fake_provider::FakeProvider::default();
|
||||
provider.push_pending();
|
||||
provider.push(text_response("continued after user interruption"));
|
||||
let assets = PromptAssets::load(
|
||||
std::path::Path::new(env!("CARGO_MANIFEST_DIR"))
|
||||
.join("prompt/cursor")
|
||||
.as_path(),
|
||||
)
|
||||
.unwrap();
|
||||
let registry = TransportRegistry::new(
|
||||
store,
|
||||
Arc::new(provider.clone()),
|
||||
PromptCompiler::new(assets),
|
||||
);
|
||||
let handle = registry
|
||||
.get_or_create("user-message-events-request")
|
||||
.await
|
||||
.unwrap();
|
||||
let mut output = handle.subscribe();
|
||||
handle
|
||||
.command(TransportCommand::Append {
|
||||
seqno: 0,
|
||||
message: Box::new(client_run_for(
|
||||
"user-message-events-request",
|
||||
"user-message-events-conversation",
|
||||
)),
|
||||
})
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let mut append_seqno = 1;
|
||||
let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(5);
|
||||
while provider.requests().is_empty() {
|
||||
assert!(
|
||||
tokio::time::Instant::now() < deadline,
|
||||
"provider did not start"
|
||||
);
|
||||
if let Ok(Some(frame)) =
|
||||
tokio::time::timeout(std::time::Duration::from_millis(20), output.recv()).await
|
||||
{
|
||||
acknowledge_kv(&handle, &mut append_seqno, &frame).await;
|
||||
}
|
||||
}
|
||||
handle
|
||||
.command(TransportCommand::Append {
|
||||
seqno: append_seqno,
|
||||
message: Box::new(runtime_user_message()),
|
||||
})
|
||||
.await
|
||||
.unwrap();
|
||||
append_seqno += 1;
|
||||
|
||||
let mut protocol_events = Vec::new();
|
||||
loop {
|
||||
let frame = tokio::time::timeout(std::time::Duration::from_secs(5), output.recv())
|
||||
.await
|
||||
.unwrap()
|
||||
.expect("RunSSE closed before successful EndStream");
|
||||
let (flags, payload) = connect::decode_frames(&frame).unwrap().pop().unwrap();
|
||||
if flags & connect::END_STREAM_FLAG != 0 {
|
||||
assert_eq!(payload.as_ref(), b"{}");
|
||||
break;
|
||||
}
|
||||
let server = pb::AgentServerMessage::decode(payload).unwrap();
|
||||
if let Some(pb::agent_server_message::Message::InteractionUpdate(update)) = server.message {
|
||||
match update.message {
|
||||
Some(pb::interaction_update::Message::ContextInjectionState(update)) => {
|
||||
assert_eq!(update.injection_id, "user-message:queued-user");
|
||||
match update.state.and_then(|state| state.state) {
|
||||
Some(pb::context_injection_state::State::Queued(_)) => {
|
||||
protocol_events.push("queued")
|
||||
}
|
||||
Some(pb::context_injection_state::State::Delivered(delivered)) => {
|
||||
assert!(!delivered.delivery_batch_id.is_empty());
|
||||
assert!(delivered.delivered_at_ms > 0);
|
||||
protocol_events.push("delivered");
|
||||
}
|
||||
_ => {}
|
||||
}
|
||||
}
|
||||
Some(pb::interaction_update::Message::UserMessageAppended(update)) => {
|
||||
let user = update.user_message.expect("appended user message");
|
||||
assert_eq!(user.message_id, "queued-user");
|
||||
assert_eq!(user.text, "queued follow-up");
|
||||
protocol_events.push("user_message_appended");
|
||||
}
|
||||
Some(pb::interaction_update::Message::TextDelta(update))
|
||||
if update.text.contains("continued after user interruption") =>
|
||||
{
|
||||
protocol_events.push("continued_output");
|
||||
}
|
||||
_ => {}
|
||||
}
|
||||
}
|
||||
acknowledge_kv(&handle, &mut append_seqno, &frame).await;
|
||||
}
|
||||
|
||||
assert_eq!(
|
||||
protocol_events,
|
||||
[
|
||||
"queued",
|
||||
"delivered",
|
||||
"user_message_appended",
|
||||
"continued_output"
|
||||
]
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn tool_call_with_empty_arguments_does_not_fail_the_run() {
|
||||
// A tool call that carries no arguments streams no argument text. Parsing it
|
||||
|
||||
Reference in New Issue
Block a user