diff --git a/server/src/cursor/conversation/output.rs b/server/src/cursor/conversation/output.rs index d66c54d..654937b 100644 --- a/server/src/cursor/conversation/output.rs +++ b/server/src/cursor/conversation/output.rs @@ -399,6 +399,14 @@ impl ConversationOutput { .unwrap_or_else(|_| serde_json::json!({})) }; } + RunEvent::UsageSnapshot(usage) => { + if !self.context.compacting { + if let Some(output_tokens) = usage.output_tokens { + self.handle.emit(&events::token_delta(output_tokens))?; + } + context_tokens = usage.context_input_tokens; + } + } RunEvent::Usage(usage) => { if !self.context.compacting { if let Some(output_tokens) = usage.output_tokens { diff --git a/server/src/cursor/services/blob_sync.rs b/server/src/cursor/services/blob_sync.rs index d9fa6c6..e88947a 100644 --- a/server/src/cursor/services/blob_sync.rs +++ b/server/src/cursor/services/blob_sync.rs @@ -20,6 +20,9 @@ use crate::{ type BlobSetSender = oneshot::Sender>; +const SET_TIMEOUT: Duration = Duration::from_secs(30 * 60); +const GET_TIMEOUT: Duration = Duration::from_secs(10 * 60); + #[derive(Clone)] pub struct BlobSynchronizer { inner: Arc, @@ -130,7 +133,7 @@ impl BlobSynchronizer { let result = tokio::select! { result = receiver => result.map_err(|_| Error::Protocol("KV SET response channel closed".into()))?, _ = cancellation.cancelled() => Err(Error::Cancelled), - _ = tokio::time::sleep(Duration::from_secs(60)) => Err(Error::Protocol(format!("KV SET timed out: {}", blob_id.to_base64()))), + _ = tokio::time::sleep(SET_TIMEOUT) => Err(Error::Protocol(format!("KV SET timed out: {}", blob_id.to_base64()))), }; if result.is_err() { self.inner.set_requests.lock().await.remove(&id); @@ -183,7 +186,7 @@ impl BlobSynchronizer { let result = tokio::select! { result = receiver => result.map_err(|_| Error::Protocol("KV GET response channel closed".into()))?, _ = cancellation.cancelled() => Err(Error::Cancelled), - _ = tokio::time::sleep(Duration::from_secs(60)) => Err(Error::Protocol(format!("KV GET timed out: {}", blob_id.to_base64()))), + _ = tokio::time::sleep(GET_TIMEOUT) => Err(Error::Protocol(format!("KV GET timed out: {}", blob_id.to_base64()))), }; if result.is_err() { self.inner.get_requests.lock().await.remove(&id); @@ -318,3 +321,18 @@ impl BlobSynchronizer { Ok(()) } } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn set_timeout_allows_slow_cursor_acknowledgements() { + assert_eq!(SET_TIMEOUT, Duration::from_secs(30 * 60)); + } + + #[test] + fn get_timeout_allows_slow_cursor_responses() { + assert_eq!(GET_TIMEOUT, Duration::from_secs(10 * 60)); + } +} diff --git a/server/src/run/compaction.rs b/server/src/run/compaction.rs index ec41441..9cb45b6 100644 --- a/server/src/run/compaction.rs +++ b/server/src/run/compaction.rs @@ -40,15 +40,23 @@ pub(super) fn estimated_tokens( .unwrap_or_else(|| estimate_context_tokens(&prepared.prompt, projected_messages)) } +pub(super) fn compaction_estimate( + prepared: &PreparedRun, + projected_messages: &[ProjectedMessage], + anchor: Option, +) -> Option { + let budget = input_budget(prepared)?; + let estimated = estimated_tokens(prepared, projected_messages, anchor); + (estimated > budget).then_some(estimated) +} + +#[cfg(test)] pub(super) fn should_compact( prepared: &PreparedRun, projected_messages: &[ProjectedMessage], anchor: Option, ) -> bool { - let Some(budget) = input_budget(prepared) else { - return false; - }; - estimated_tokens(prepared, projected_messages, anchor) > budget + compaction_estimate(prepared, projected_messages, anchor).is_some() } pub(super) fn validate_compacted( diff --git a/server/src/run/engine.rs b/server/src/run/engine.rs index 42b3808..4694f36 100644 --- a/server/src/run/engine.rs +++ b/server/src/run/engine.rs @@ -183,9 +183,21 @@ impl RunEngine { Ok(history) => history, Err(error) => return (RunOutcome::Failed(error.into()), usage), }; - if prepared.action != RunAction::Compact - && super::compaction::should_compact(prepared, &history, context_usage_anchor) - { + let compaction_estimate = (prepared.action != RunAction::Compact) + .then(|| { + super::compaction::compaction_estimate(prepared, &history, context_usage_anchor) + }) + .flatten(); + if let Some(estimated_tokens) = compaction_estimate { + if emit( + client, + RunEvent::UsageSnapshot(context_usage_snapshot(estimated_tokens)), + ) + .await + .is_err() + { + return (client_failure(), usage); + } match self .auto_compact(prepared, checkpoint, &messages, client, cancellation) .await @@ -703,20 +715,6 @@ impl RunEngine { emit(client, RunEvent::AutoCompactionStarted) .await .map_err(|_| client_failure())?; - emit( - client, - RunEvent::Usage(Usage { - input_tokens: Some(0), - context_input_tokens: Some(0), - output_tokens: Some(0), - total_tokens: Some(0), - cache_read_tokens: Some(0), - cache_write_tokens: Some(0), - reasoning_tokens: Some(0), - }), - ) - .await - .map_err(|_| client_failure())?; let provider_call_index = self .store .begin_provider_call(&prepared.run_id) @@ -862,6 +860,9 @@ impl RunEngine { emit(client, RunEvent::AutoCompactionCompleted) .await .map_err(|_| client_failure())?; + emit(client, RunEvent::UsageSnapshot(context_usage_snapshot(0))) + .await + .map_err(|_| client_failure())?; checkpoint = super::messages::append_batches( &self.store, prepared, @@ -922,6 +923,18 @@ async fn hydrate_tool_images( Ok(()) } +fn context_usage_snapshot(tokens: u64) -> Usage { + Usage { + input_tokens: Some(tokens), + context_input_tokens: Some(tokens), + output_tokens: Some(0), + total_tokens: Some(tokens), + cache_read_tokens: Some(0), + cache_write_tokens: Some(0), + reasoning_tokens: Some(0), + } +} + fn update_context_usage_anchor( anchor: &mut Option, usage: Usage, diff --git a/server/src/run/event.rs b/server/src/run/event.rs index 345be46..01918aa 100644 --- a/server/src/run/event.rs +++ b/server/src/run/event.rs @@ -129,6 +129,7 @@ pub enum RunEvent { ToolCallEnd { index: usize, }, + UsageSnapshot(Usage), Usage(Usage), ExecuteToolRound { round_id: ToolRoundId, diff --git a/server/tests/compaction.rs b/server/tests/compaction.rs index 2781544..4eeb8ee 100644 --- a/server/tests/compaction.rs +++ b/server/tests/compaction.rs @@ -268,9 +268,14 @@ async fn automatic_compaction_preflights_provider_input_and_records_rebuilt_toke assert_eq!(second.summary_started, 1); assert_eq!(second.summary_completed, 1); assert_eq!( - &second.interaction_events[..2], - &["summary_started", "token_delta:0"], - "automatic compaction must immediately reset Cursor usage" + &second.interaction_events[..4], + &[ + "token_delta:0", + "summary_started", + "summary_completed", + "token_delta:0", + ], + "automatic compaction must publish estimated usage before summarizing and zero usage after" ); let compacted_tokens = second .checkpoints @@ -530,7 +535,8 @@ async fn run( output.summary.push_str(&delta.summary) } Some(pb::interaction_update::Message::SummaryCompleted(_)) => { - output.summary_completed += 1 + output.summary_completed += 1; + output.interaction_events.push("summary_completed".into()); } Some(pb::interaction_update::Message::TurnEnded(_)) => output.turn_ended += 1, Some(pb::interaction_update::Message::TokenDelta(delta)) => {