1use crate::executor::error::{ExecutorError, ExecutorResult};
2use crate::executor::rehydrate::rehydrate_conversation;
3use crate::executor::request::{ExecutionContext, RequestContext};
4use crate::executor::upstream::fetch_blocking_payload;
5use crate::types::event::MessageStatus;
6use crate::types::io::input::latest_compaction_window;
7use crate::types::io::{
8 CompactionItem, InputContent, InputItem, InputMessage, InputMessageContent, OutputItem, ResponseUsage,
9 ResponsesInput,
10};
11use crate::types::request_response::{CompactRequest, CompactedResponse, RequestPayload, ResponsePayload};
12use crate::utils::common::{serialize_to_string, utcnow_str, uuid7_str};
13
14const COMPACTION_PROMPT: &str = "You are performing a CONTEXT CHECKPOINT COMPACTION. Create a concise handoff summary that preserves current progress, decisions, constraints, unresolved work, and critical references for the next model. Return only the summary.";
15
16fn retained_user_window(items: &[InputItem]) -> Vec<InputItem> {
17 let window = latest_compaction_window(items);
18 items
19 .iter()
20 .enumerate()
21 .filter_map(|(index, item)| {
22 let InputItem::Message(message) = item else {
23 return None;
24 };
25 if message.role != "user"
26 || window.is_some_and(|window| index < window.latest_index() && !window.retains_user_item(index, item))
27 {
28 return None;
29 }
30 let mut retained = message.clone();
31 retained.id = Some(uuid7_str("msg_"));
32 retained.status = Some(MessageStatus::Completed);
33 Some(InputItem::Message(retained))
34 })
35 .collect()
36}
37
38fn finish_compacted_window(mut output: Vec<InputItem>, summary: String) -> Vec<InputItem> {
39 output.push(InputItem::Compaction(CompactionItem {
40 id: Some(uuid7_str("cmp_")),
41 encrypted_content: summary,
42 }));
43 output
44}
45
46#[cfg(test)]
47fn build_compacted_window(items: &[InputItem], summary: String) -> Vec<InputItem> {
48 finish_compacted_window(retained_user_window(items), summary)
49}
50
51fn response_output_text(output: &[OutputItem]) -> Option<String> {
52 let text = output
53 .iter()
54 .filter_map(|item| match item {
55 OutputItem::Message(message) => Some(message),
56 _ => None,
57 })
58 .flat_map(|message| message.content.iter())
59 .map(|content| content.text.trim())
60 .filter(|text| !text.is_empty())
61 .collect::<Vec<_>>()
62 .join("\n");
63 (!text.is_empty()).then_some(text)
64}
65
66fn value_has_content(value: &serde_json::Value) -> bool {
67 match value {
68 serde_json::Value::Null => false,
69 serde_json::Value::String(text) => !text.trim().is_empty(),
70 serde_json::Value::Array(values) => values.iter().any(value_has_content),
71 serde_json::Value::Object(values) => values.values().any(value_has_content),
72 serde_json::Value::Bool(_) | serde_json::Value::Number(_) => true,
73 }
74}
75
76fn item_has_meaningful_context(item: &InputItem) -> bool {
77 match item {
78 InputItem::Message(message) => match &message.content {
79 InputMessageContent::Text(text) => !text.trim().is_empty(),
80 InputMessageContent::Parts(parts) => parts.iter().any(|part| match part {
81 InputContent::InputText(text) | InputContent::OutputText(text) | InputContent::ReasoningText(text) => {
82 !text.text.trim().is_empty()
83 }
84 InputContent::InputImage(image) => image.image_url.as_deref().is_some_and(|url| !url.trim().is_empty()),
85 InputContent::Unknown => false,
86 }),
87 },
88 InputItem::FunctionCall(call) => !call.name.trim().is_empty() || !call.arguments.trim().is_empty(),
89 InputItem::FunctionCallOutput(output) => output.output.has_content(),
90 InputItem::CustomToolCall(call) => !call.name.trim().is_empty() || !call.input.trim().is_empty(),
91 InputItem::CustomToolCallOutput(output) => output.output.has_content(),
92 InputItem::Reasoning(reasoning) => {
93 reasoning.content.iter().any(|content| !content.text.trim().is_empty())
94 || reasoning.summary.iter().any(value_has_content)
95 || reasoning.encrypted_content.as_ref().is_some_and(value_has_content)
96 }
97 InputItem::Compaction(compaction) => !compaction.encrypted_content.trim().is_empty(),
98 InputItem::CompactionTrigger | InputItem::Unknown => false,
99 }
100}
101
102fn completed_summary_text(response: &ResponsePayload) -> ExecutorResult<String> {
103 if response.status != "completed" || response.error.is_some() {
104 let details = response
105 .error
106 .as_ref()
107 .and_then(|error| serialize_to_string(error).ok())
108 .or_else(|| {
109 response
110 .incomplete_details
111 .as_ref()
112 .and_then(|details| details.reason.clone())
113 })
114 .unwrap_or_else(|| "upstream returned no failure details".to_owned());
115 return Err(ExecutorError::CompactionFailed {
116 status: response.status.clone(),
117 details,
118 });
119 }
120 response_output_text(&response.output).ok_or_else(|| ExecutorError::CompactionFailed {
121 status: response.status.clone(),
122 details: "upstream returned no summary text".to_owned(),
123 })
124}
125
126#[must_use]
131pub(crate) fn estimate_input_tokens(input: &ResponsesInput) -> u64 {
132 let serialized = serialize_to_string(&input.model_input()).unwrap_or_default();
133 let bytes = u64::try_from(serialized.len()).unwrap_or(u64::MAX);
134 bytes.saturating_add(3) / 4
135}
136
137fn request_payload(model: String, input: ResponsesInput, instructions: Option<String>) -> RequestPayload {
138 RequestPayload {
139 model,
140 input,
141 instructions,
142 previous_response_id: None,
143 conversation_id: None,
144 tools: None,
145 tool_choice: None,
146 stream: false,
147 store: false,
148 include: None,
149 temperature: None,
150 top_p: None,
151 max_output_tokens: None,
152 truncation: None,
153 metadata: None,
154 parallel_tool_calls: None,
155 cache_salt: None,
156 context_management: None,
157 }
158}
159
160pub(crate) async fn compact_items(
167 model: &str,
168 input: ResponsesInput,
169 instructions: Option<&str>,
170 exec_ctx: &ExecutionContext,
171 auth: Option<&str>,
172) -> ExecutorResult<(Vec<InputItem>, ResponseUsage)> {
173 let original_items = Vec::from(input);
174 if !original_items.iter().any(item_has_meaningful_context) {
175 return Err(ExecutorError::InvalidRequest(
176 "compaction requires non-empty input or previous_response_id context".to_owned(),
177 ));
178 }
179
180 let compacted = retained_user_window(&original_items);
181 let mut summary_items: Vec<InputItem> = original_items
182 .into_iter()
183 .filter(|item| !item.is_compaction_trigger())
184 .collect();
185 summary_items.push(InputItem::Message(InputMessage {
186 id: None,
187 role: "user".to_owned(),
188 status: None,
189 content: InputMessageContent::Text(COMPACTION_PROMPT.to_owned()),
190 }));
191 let instructions = instructions.map(str::to_owned);
192 let original_request = request_payload(
193 model.to_owned(),
194 ResponsesInput::Items(Vec::new()),
195 instructions.clone(),
196 );
197 let enriched_request = request_payload(model.to_owned(), ResponsesInput::Items(summary_items), instructions);
198 let ctx = RequestContext {
199 original_request,
200 enriched_request,
201 new_input_items: Vec::new(),
202 response_id: uuid7_str("resp_"),
203 conversation_id: None,
204 conversation_version: None,
205 };
206 let response = fetch_blocking_payload(&ctx, exec_ctx, auth).await?;
207 let summary = completed_summary_text(&response)?;
208
209 Ok((
210 finish_compacted_window(compacted, summary),
211 response.usage.unwrap_or_default(),
212 ))
213}
214
215pub(crate) async fn maybe_compact_context(
224 ctx: &mut RequestContext,
225 exec_ctx: &ExecutionContext,
226 auth: Option<&str>,
227) -> ExecutorResult<Option<ResponseUsage>> {
228 let threshold = ctx
229 .enriched_request
230 .context_management
231 .as_deref()
232 .unwrap_or_default()
233 .iter()
234 .find(|entry| entry.type_ == "compaction")
235 .and_then(|entry| entry.compact_threshold);
236 let Some(threshold) = threshold else {
237 return Ok(None);
238 };
239 let estimated_tokens = estimate_input_tokens(&ctx.enriched_request.input);
240 if estimated_tokens <= threshold {
241 return Ok(None);
242 }
243
244 tracing::debug!(
245 estimated_tokens,
246 threshold,
247 "automatically compacting resolved response input"
248 );
249 let model = ctx.enriched_request.model.clone();
250 let instructions = ctx.enriched_request.instructions.clone();
251 let input = std::mem::replace(&mut ctx.enriched_request.input, ResponsesInput::Items(Vec::new()));
252 let (compacted, usage) = compact_items(&model, input, instructions.as_deref(), exec_ctx, auth).await?;
253 ctx.enriched_request.input = ResponsesInput::Items(compacted.clone());
254 ctx.new_input_items = compacted;
255 Ok(Some(usage))
256}
257
258pub async fn compact_response(
265 request: CompactRequest,
266 exec_ctx: &ExecutionContext,
267 auth: Option<&str>,
268) -> ExecutorResult<CompactedResponse> {
269 if request.input.is_none() && request.previous_response_id.is_none() {
270 return Err(ExecutorError::InvalidRequest(
271 "compaction requires input or previous_response_id".to_owned(),
272 ));
273 }
274
275 let mut payload = request_payload(
276 request.model,
277 request.input.unwrap_or_else(|| ResponsesInput::Items(Vec::new())),
278 request.instructions,
279 );
280 payload.previous_response_id = request.previous_response_id;
281 let mut ctx = rehydrate_conversation(payload, exec_ctx).await?;
282 let model = ctx.enriched_request.model.clone();
283 let instructions = ctx.enriched_request.instructions.clone();
284 let input = std::mem::replace(&mut ctx.enriched_request.input, ResponsesInput::Items(Vec::new()));
285 let (output, usage) = compact_items(&model, input, instructions.as_deref(), exec_ctx, auth).await?;
286
287 let response_id = ctx.response_id.clone();
288 ctx.new_input_items.clone_from(&output);
289 match exec_ctx.resp_handler.execute_turn(ctx, Vec::new()).await {
290 Ok(()) | Err(ExecutorError::Storage(crate::StorageError::NotConfigured)) => {}
291 Err(error) => return Err(error),
292 }
293
294 Ok(CompactedResponse {
295 id: response_id,
296 object: "response.compaction".to_owned(),
297 created_at: utcnow_str(),
298 output,
299 usage,
300 })
301}
302
303#[cfg(test)]
304mod tests {
305 use std::sync::Arc;
306
307 use axum::Router;
308 use axum::routing::post;
309
310 use crate::executor::modes::{ConversationHandler, ResponseHandler};
311 use crate::executor::request::ExecutionContext;
312 use crate::storage::{ConversationStore, InOutItem, ResponseMetadata, ResponseStore, create_pool_with_schema};
313 use crate::types::event::MessageStatus;
314 use crate::types::io::{
315 CompactionItem, FunctionToolResultMessage, InputItem, InputMessage, InputMessageContent, ResponsesInput,
316 };
317
318 use super::{build_compacted_window, compact_response, completed_summary_text, estimate_input_tokens};
319
320 fn user_message(text: &str) -> InputItem {
321 InputItem::Message(InputMessage {
322 id: None,
323 role: "user".to_owned(),
324 status: None,
325 content: InputMessageContent::Text(text.to_owned()),
326 })
327 }
328
329 async fn mock_execution_context(response_store: ResponseStore) -> (ExecutionContext, tokio::task::JoinHandle<()>) {
330 let app = Router::new().route(
331 "/v1/responses",
332 post(|| async {
333 axum::Json(serde_json::json!({
334 "id": "resp_upstream",
335 "object": "response",
336 "created_at": 0,
337 "model": "test-model",
338 "status": "completed",
339 "output": [{
340 "id": "msg_upstream",
341 "type": "message",
342 "role": "assistant",
343 "status": "completed",
344 "content": [{
345 "type": "output_text",
346 "text": "durable summary",
347 "annotations": []
348 }]
349 }],
350 "usage": {
351 "input_tokens": 12,
352 "output_tokens": 3,
353 "total_tokens": 15
354 },
355 "incomplete_details": null,
356 "error": null,
357 "previous_response_id": null,
358 "conversation_id": null,
359 "instructions": null
360 }))
361 }),
362 );
363 let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
364 .await
365 .expect("bind mock inference server");
366 let address = listener.local_addr().expect("mock server address");
367 let server = tokio::spawn(async move {
368 axum::serve(listener, app).await.ok();
369 });
370 let exec_ctx = ExecutionContext::new(
371 ConversationHandler::new(ConversationStore::disabled()),
372 ResponseHandler::new(response_store),
373 Arc::new(reqwest::Client::new()),
374 format!("http://{address}"),
375 );
376 (exec_ctx, server)
377 }
378
379 fn compact_request() -> crate::CompactRequest {
380 serde_json::from_value(serde_json::json!({
381 "model": "test-model",
382 "input": [{"role": "user", "content": "important context"}]
383 }))
384 .expect("valid compact request")
385 }
386
387 #[test]
388 fn compacted_window_retains_user_messages_and_replaces_old_compaction() {
389 let items = vec![
390 InputItem::Compaction(CompactionItem {
391 id: Some("cmp_old".to_owned()),
392 encrypted_content: "old summary".to_owned(),
393 }),
394 user_message("first"),
395 InputItem::FunctionCallOutput(FunctionToolResultMessage {
396 call_id: "call_1".to_owned(),
397 output: "tool output".into(),
398 }),
399 user_message("second"),
400 ];
401
402 let output = build_compacted_window(&items, "new summary".to_owned());
403
404 assert_eq!(
405 output
406 .iter()
407 .filter(|item| matches!(item, InputItem::Message(_)))
408 .count(),
409 2
410 );
411 assert_eq!(
412 output
413 .iter()
414 .filter(|item| matches!(item, InputItem::Compaction(_)))
415 .count(),
416 1
417 );
418 assert!(matches!(output.last(), Some(InputItem::Compaction(_))));
419 for item in output.iter().filter_map(|item| match item {
420 InputItem::Message(message) => Some(message),
421 _ => None,
422 }) {
423 assert_eq!(item.status, Some(MessageStatus::Completed));
424 assert!(item.id.as_deref().is_some_and(|id| id.starts_with("msg_")));
425 }
426 }
427
428 #[test]
429 fn token_estimate_counts_replayed_content() {
430 let input = ResponsesInput::Items(vec![
431 user_message("hello context"),
432 InputItem::FunctionCallOutput(FunctionToolResultMessage {
433 call_id: "call_1".to_owned(),
434 output: "substantial tool output".into(),
435 }),
436 ]);
437
438 assert!(estimate_input_tokens(&input) > 0);
439 }
440
441 #[test]
442 fn partial_text_from_incomplete_or_failed_summary_is_rejected() {
443 for (status, error, incomplete_details) in [
444 (
445 "incomplete",
446 serde_json::Value::Null,
447 serde_json::json!({"reason": "max_output_tokens"}),
448 ),
449 (
450 "failed",
451 serde_json::json!({"message": "model failed"}),
452 serde_json::Value::Null,
453 ),
454 ] {
455 let response: crate::ResponsePayload = serde_json::from_value(serde_json::json!({
456 "id": "resp_upstream",
457 "object": "response",
458 "created_at": 0,
459 "model": "test-model",
460 "status": status,
461 "output": [{
462 "id": "msg_partial",
463 "type": "message",
464 "role": "assistant",
465 "status": "completed",
466 "content": [{"type": "output_text", "text": "partial summary", "annotations": []}]
467 }],
468 "usage": null,
469 "incomplete_details": incomplete_details,
470 "error": error,
471 "previous_response_id": null,
472 "conversation_id": null,
473 "instructions": null
474 }))
475 .expect("valid upstream response");
476
477 let error = completed_summary_text(&response).expect_err("partial summary must be rejected");
478 assert!(matches!(error, crate::executor::ExecutorError::CompactionFailed { .. }));
479 }
480 }
481
482 #[test]
483 fn completed_response_without_summary_is_an_upstream_failure() {
484 let response: crate::ResponsePayload = serde_json::from_value(serde_json::json!({
485 "id": "resp_upstream",
486 "object": "response",
487 "created_at": 0,
488 "model": "test-model",
489 "status": "completed",
490 "output": [],
491 "usage": null,
492 "incomplete_details": null,
493 "error": null,
494 "previous_response_id": null,
495 "conversation_id": null,
496 "instructions": null
497 }))
498 .expect("valid upstream response");
499
500 let error = completed_summary_text(&response).expect_err("missing summary must be rejected");
501
502 assert!(matches!(
503 error,
504 crate::executor::ExecutorError::CompactionFailed {
505 ref status,
506 ref details
507 } if status == "completed" && details.contains("no summary text")
508 ));
509 assert_eq!(error.http_status(), http::StatusCode::BAD_GATEWAY);
510 assert_eq!(error.error_code(), "upstream_error");
511 }
512
513 #[tokio::test]
514 async fn disabled_storage_does_not_fail_compaction() {
515 let (exec_ctx, server) = mock_execution_context(ResponseStore::disabled()).await;
516
517 let response = compact_response(compact_request(), &exec_ctx, None)
518 .await
519 .expect("disabled persistence should be ignored");
520
521 assert_eq!(response.object, "response.compaction");
522 assert_eq!(response.usage.total_tokens, 15);
523 assert!(response.id.starts_with("resp_"));
524 assert!(
525 matches!(response.output.last(), Some(InputItem::Compaction(item)) if item.encrypted_content == "durable summary")
526 );
527 server.abort();
528 }
529
530 #[tokio::test]
531 async fn empty_context_is_rejected_before_summarization() {
532 let (exec_ctx, server) = mock_execution_context(ResponseStore::disabled()).await;
533 for input in [
534 serde_json::json!(""),
535 serde_json::json!([{"role": "user", "content": " "}]),
536 serde_json::json!([{"type": "future_item"}]),
537 ] {
538 let request = serde_json::from_value(serde_json::json!({
539 "model": "test-model",
540 "input": input
541 }))
542 .expect("structurally valid compact request");
543 let error = compact_response(request, &exec_ctx, None)
544 .await
545 .expect_err("empty context must fail");
546 assert!(matches!(error, crate::executor::ExecutorError::InvalidRequest(_)));
547 }
548 server.abort();
549 }
550
551 #[tokio::test]
552 async fn compaction_persists_a_reusable_response_checkpoint() {
553 let pool = create_pool_with_schema(Some("sqlite::memory:"))
554 .await
555 .expect("create response store");
556 let response_store = ResponseStore::new(pool);
557 let (exec_ctx, server) = mock_execution_context(response_store.clone()).await;
558
559 let response = compact_response(compact_request(), &exec_ctx, None)
560 .await
561 .expect("compaction succeeds");
562 let history = response_store
563 .rehydrate(&response.id)
564 .await
565 .expect("compaction checkpoint can be rehydrated");
566
567 assert_eq!(history.len(), 2);
568 assert!(matches!(history[0], InOutItem::Input(InputItem::Message(_))));
569 assert!(matches!(history[1], InOutItem::Input(InputItem::Compaction(_))));
570 server.abort();
571 }
572
573 #[tokio::test]
574 async fn compaction_resolves_and_replaces_previous_response_history() {
575 let pool = create_pool_with_schema(Some("sqlite::memory:"))
576 .await
577 .expect("create response store");
578 let response_store = ResponseStore::new(pool);
579 response_store
580 .persist(
581 "resp_previous",
582 None,
583 vec![InOutItem::Input(user_message("remember banana"))],
584 &ResponseMetadata {
585 model: "test-model".to_owned(),
586 previous_response_id: None,
587 effective_tools: None,
588 effective_tool_choice: crate::ToolChoice::Auto,
589 effective_instructions: None,
590 },
591 )
592 .await
593 .expect("seed previous response");
594 let (exec_ctx, server) = mock_execution_context(response_store.clone()).await;
595 let request = serde_json::from_value(serde_json::json!({
596 "model": "test-model",
597 "previous_response_id": "resp_previous"
598 }))
599 .expect("valid previous-response compact request");
600
601 let response = compact_response(request, &exec_ctx, None)
602 .await
603 .expect("previous response compacts");
604 let history = response_store
605 .rehydrate(&response.id)
606 .await
607 .expect("compacted continuation rehydrates");
608 let model_input = ResponsesInput::Items(InOutItem::into_input_items(history));
609 let serialized = serde_json::to_value(model_input.model_input()).expect("model input serializes");
610
611 assert_eq!(serialized.as_array().map(Vec::len), Some(2));
612 assert_eq!(serialized[0]["content"], "remember banana");
613 assert_eq!(serialized[1]["role"], "assistant");
614 assert_eq!(serialized[1]["content"][0]["text"], "durable summary");
615 server.abort();
616 }
617}