Skip to main content

agentic_core/executor/
compaction.rs

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/// Estimate the current model-facing context size without requiring a model-specific tokenizer.
127///
128/// The approximation deliberately includes JSON structure and rounds up at four UTF-8 bytes per
129/// token. It is deterministic, inexpensive, and errs slightly toward compacting early.
130#[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
160/// Summarize an already-resolved item history and return its canonical compacted window.
161///
162/// # Errors
163///
164/// Returns an invalid-request error for empty input, an upstream error for an unusable model
165/// summary, and propagates inference and serialization failures.
166pub(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
215/// Apply the first configured compaction threshold to the resolved request history.
216///
217/// Returns the summarization usage when compaction ran. The configuration remains on the
218/// client-facing request context but is never part of [`crate::types::request_response::UpstreamRequest`].
219///
220/// # Errors
221///
222/// Propagates inference errors from the summarization request.
223pub(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
258/// Compact direct input or a stored previous-response chain into a reusable item window.
259///
260/// # Errors
261///
262/// Returns an invalid-request error when neither input nor a previous response ID is supplied,
263/// and propagates history, inference, and persistence failures.
264pub 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}