1use std::collections::HashMap;
2use std::sync::Mutex;
3use std::sync::atomic::{AtomicU64, Ordering};
4use std::time::{SystemTime, UNIX_EPOCH};
5
6use crate::anthropic::sse::parse_sse_events;
7use crate::provider::RequestContext;
8use crate::providers::codex::client::{CodexError, CodexHttpClient};
9
10use super::translate::request::{
11 ResponsesContentPart, ResponsesInputItem, ResponsesRequest, is_compact_message_text,
12};
13
14const RETAINED_MESSAGE_TOKEN_BUDGET: u64 = 20_000;
15const STATE_TTL_MS: u64 = 30 * 60 * 1_000;
16const MAX_STATES: usize = 1_000;
17const MAX_STATE_BYTES: usize = 4 * 1024 * 1024;
18const MAX_TOTAL_STATE_BYTES: usize = 20_000_000;
19const MIN_PORTABLE_SUMMARY_BYTES: usize = 32;
20
21#[derive(Debug)]
22pub enum CompactionError {
23 Upstream(CodexError),
24 InvalidResponse(String),
25}
26
27impl std::fmt::Display for CompactionError {
28 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
29 match self {
30 Self::Upstream(error) => write!(f, "{error}"),
31 Self::InvalidResponse(message) => f.write_str(message),
32 }
33 }
34}
35
36enum CompactionPhase {
37 Preparing,
38 Unconfirmed,
39 Anchored { portable_summary: String },
40}
41
42struct CompactionState {
43 attempt: CompactionAttempt,
44 model: String,
45 native_history: Vec<ResponsesInputItem>,
46 phase: CompactionPhase,
47 updated_at: u64,
48}
49
50#[derive(Debug, Clone, Copy, PartialEq, Eq)]
51pub struct CompactionAttempt(u64);
52
53pub struct CompactionReplay {
54 pub request: ResponsesRequest,
55 pub attempt: CompactionAttempt,
56}
57
58#[derive(Default)]
59struct CompactionRegistry {
60 states: HashMap<String, CompactionState>,
61 total_bytes: usize,
62}
63
64static REGISTRY: Mutex<Option<CompactionRegistry>> = Mutex::new(None);
65static NEXT_ATTEMPT_ID: AtomicU64 = AtomicU64::new(1);
66
67pub async fn request_compaction(
68 client: &CodexHttpClient,
69 request: &ResponsesRequest,
70 ctx: &RequestContext,
71) -> Result<Vec<ResponsesInputItem>, CompactionError> {
72 let (envelope, conversation) = split_input_envelope(&request.input);
73 let conversation = without_compaction_instruction(conversation);
74 let mut compaction_request = request.clone();
75 compaction_request.instructions = None;
76 compaction_request.input = envelope
77 .iter()
78 .filter(|item| matches!(item, ResponsesInputItem::AdditionalTools { .. }))
79 .cloned()
80 .chain(conversation.iter().cloned())
81 .chain(std::iter::once(ResponsesInputItem::CompactionTrigger))
82 .collect();
83 compaction_request.include = Some(vec!["reasoning.encrypted_content".to_string()]);
84
85 let response = client
86 .post_codex_for_owner(&compaction_request, ctx, None)
87 .await
88 .map_err(CompactionError::Upstream)?;
89 let compaction = parse_compaction_response(&response.body)?;
90 Ok(build_compacted_history(&conversation, compaction))
91}
92
93pub fn begin_compaction(session_id: &str, model: &str) -> CompactionAttempt {
94 let attempt = CompactionAttempt(NEXT_ATTEMPT_ID.fetch_add(1, Ordering::Relaxed));
95 let state = CompactionState {
96 attempt,
97 model: model.to_string(),
98 native_history: Vec::new(),
99 phase: CompactionPhase::Preparing,
100 updated_at: now_ms(),
101 };
102 let now = state.updated_at;
103 let mut guard = REGISTRY.lock().unwrap();
104 let registry = guard.get_or_insert_with(CompactionRegistry::default);
105 evict_states(registry, now);
106 registry.states.insert(session_id.to_string(), state);
107 evict_states(registry, now);
108 attempt
109}
110
111pub fn store_compaction(
112 session_id: &str,
113 attempt: CompactionAttempt,
114 native_history: Vec<ResponsesInputItem>,
115) -> bool {
116 let now = now_ms();
117 let mut guard = REGISTRY.lock().unwrap();
118 let Some(registry) = guard.as_mut() else {
119 return false;
120 };
121 evict_states(registry, now);
122 let Some(state) = registry.states.get_mut(session_id) else {
123 return false;
124 };
125 if state.attempt != attempt || !matches!(state.phase, CompactionPhase::Preparing) {
126 return false;
127 }
128 state.native_history = native_history;
129 state.phase = CompactionPhase::Unconfirmed;
130 state.updated_at = now;
131 if state_size(session_id, state) > MAX_STATE_BYTES {
132 registry.states.remove(session_id);
133 update_total_bytes(registry);
134 return false;
135 }
136 evict_states(registry, now);
137 registry
138 .states
139 .get(session_id)
140 .is_some_and(|state| state.attempt == attempt)
141}
142
143pub fn activate_compaction(
144 session_id: Option<&str>,
145 attempt: Option<CompactionAttempt>,
146 model: &str,
147 output: &[ResponsesInputItem],
148) -> bool {
149 let (Some(session_id), Some(attempt)) = (session_id, attempt) else {
150 return false;
151 };
152 let Some(portable_summary) = portable_summary_text(output) else {
153 abort_compaction_attempt(Some(session_id), Some(attempt));
154 return false;
155 };
156
157 let now = now_ms();
158 let mut guard = REGISTRY.lock().unwrap();
159 let Some(registry) = guard.as_mut() else {
160 return false;
161 };
162 evict_states(registry, now);
163 let Some(state) = registry.states.get_mut(session_id) else {
164 return false;
165 };
166 if state.attempt != attempt {
167 return false;
168 }
169 if state.model != model || !matches!(state.phase, CompactionPhase::Unconfirmed) {
170 registry.states.remove(session_id);
171 update_total_bytes(registry);
172 return false;
173 }
174 state.phase = CompactionPhase::Anchored { portable_summary };
175 state.updated_at = now;
176 if state_size(session_id, state) > MAX_STATE_BYTES {
177 registry.states.remove(session_id);
178 update_total_bytes(registry);
179 return false;
180 }
181 evict_states(registry, now);
182 registry.states.contains_key(session_id)
183}
184
185pub fn apply_compaction_replay(
186 session_id: Option<&str>,
187 request: &ResponsesRequest,
188) -> Option<CompactionReplay> {
189 let session_id = session_id?;
190 let now = now_ms();
191 let mut guard = REGISTRY.lock().unwrap();
192 let registry = guard.as_mut()?;
193 evict_states(registry, now);
194 let state = registry.states.get_mut(session_id)?;
195 if !matches!(state.phase, CompactionPhase::Anchored { .. }) {
196 return None;
197 }
198 if state.model != request.model {
199 registry.states.remove(session_id);
200 update_total_bytes(registry);
201 return None;
202 }
203 let CompactionPhase::Anchored { portable_summary } = &state.phase else {
204 return None;
205 };
206
207 let (envelope, conversation) = split_input_envelope(&request.input);
208 let summary_item = conversation.first()?;
209 let Some(text) = message_text(summary_item) else {
210 registry.states.remove(session_id);
211 update_total_bytes(registry);
212 return None;
213 };
214 if text.match_indices(portable_summary).count() != 1 {
215 registry.states.remove(session_id);
216 update_total_bytes(registry);
217 return None;
218 }
219 if conversation.len() == 1 {
220 return None;
221 }
222
223 let mut replay = request.clone();
224 replay.input = envelope
225 .iter()
226 .cloned()
227 .chain(state.native_history.iter().cloned())
228 .chain(conversation[1..].iter().cloned())
229 .collect();
230 if serialized_size(&replay.input) > MAX_STATE_BYTES {
231 registry.states.remove(session_id);
232 update_total_bytes(registry);
233 return None;
234 }
235 state.updated_at = now;
236 Some(CompactionReplay {
237 request: replay,
238 attempt: state.attempt,
239 })
240}
241
242pub fn abort_compaction_attempt(session_id: Option<&str>, attempt: Option<CompactionAttempt>) {
243 let (Some(session_id), Some(attempt)) = (session_id, attempt) else {
244 return;
245 };
246 let mut guard = REGISTRY.lock().unwrap();
247 let Some(registry) = guard.as_mut() else {
248 return;
249 };
250 if registry
251 .states
252 .get(session_id)
253 .is_some_and(|state| state.attempt == attempt)
254 {
255 registry.states.remove(session_id);
256 update_total_bytes(registry);
257 }
258}
259
260pub fn clear_compaction(session_id: &str) {
261 let mut guard = REGISTRY.lock().unwrap();
262 if let Some(registry) = guard.as_mut() {
263 registry.states.remove(session_id);
264 update_total_bytes(registry);
265 }
266}
267
268fn split_input_envelope(
269 input: &[ResponsesInputItem],
270) -> (&[ResponsesInputItem], &[ResponsesInputItem]) {
271 let prefix_len = input
272 .iter()
273 .take_while(|item| is_envelope_item(item))
274 .count();
275 input.split_at(prefix_len)
276}
277
278fn without_compaction_instruction(input: &[ResponsesInputItem]) -> Vec<ResponsesInputItem> {
279 let mut input = input.to_vec();
280 let remove_empty_message = if let Some(ResponsesInputItem::Message { role, content }) =
281 input.last_mut()
282 && role == "user"
283 {
284 content.retain(|part| {
285 !matches!(
286 part,
287 ResponsesContentPart::InputText { text }
288 if is_compact_message_text(text)
289 )
290 });
291 content.is_empty()
292 } else {
293 false
294 };
295 if remove_empty_message {
296 input.pop();
297 }
298 input
299}
300
301fn is_envelope_item(item: &ResponsesInputItem) -> bool {
302 match item {
303 ResponsesInputItem::AdditionalTools { .. } => true,
304 ResponsesInputItem::Message { role, .. } => role == "developer",
305 _ => false,
306 }
307}
308
309fn portable_summary_text(output: &[ResponsesInputItem]) -> Option<String> {
310 let text = output
311 .iter()
312 .filter_map(|item| match item {
313 ResponsesInputItem::Message { role, content } if role == "assistant" => Some(content),
314 _ => None,
315 })
316 .flat_map(|content| content.iter())
317 .filter_map(|part| match part {
318 ResponsesContentPart::InputText { text }
319 | ResponsesContentPart::OutputText { text } => Some(text.as_str()),
320 ResponsesContentPart::InputImage { .. } => None,
321 })
322 .collect::<String>();
323 let trimmed = text.trim();
324 let summary = trimmed
325 .split_once("<summary>")
326 .and_then(|(_, rest)| rest.split_once("</summary>"))
327 .map(|(summary, _)| summary.trim())
328 .filter(|summary| !summary.is_empty())
329 .unwrap_or(trimmed);
330 (summary.len() >= MIN_PORTABLE_SUMMARY_BYTES).then(|| summary.to_string())
331}
332
333fn message_text(item: &ResponsesInputItem) -> Option<String> {
334 let ResponsesInputItem::Message { content, .. } = item else {
335 return None;
336 };
337 Some(
338 content
339 .iter()
340 .filter_map(|part| match part {
341 ResponsesContentPart::InputText { text }
342 | ResponsesContentPart::OutputText { text } => Some(text.as_str()),
343 ResponsesContentPart::InputImage { .. } => None,
344 })
345 .collect(),
346 )
347}
348
349fn parse_compaction_response(body: &[u8]) -> Result<ResponsesInputItem, CompactionError> {
350 let mut completed = false;
351 let mut compacted = Vec::new();
352
353 for event in parse_sse_events(body) {
354 if event.data == "[DONE]" {
355 continue;
356 }
357 let Ok(payload) = serde_json::from_str::<serde_json::Value>(&event.data) else {
358 continue;
359 };
360 match payload.get("type").and_then(serde_json::Value::as_str) {
361 Some("error" | "response.error" | "response.failed") => {
362 let message = payload
363 .pointer("/response/error/message")
364 .or_else(|| payload.pointer("/error/message"))
365 .or_else(|| payload.get("message"))
366 .and_then(serde_json::Value::as_str)
367 .unwrap_or("remote compaction failed");
368 return Err(CompactionError::InvalidResponse(message.to_string()));
369 }
370 Some("response.output_item.done") => {
371 let Some(item) = payload.get("item") else {
372 continue;
373 };
374 if item.get("type").and_then(serde_json::Value::as_str) == Some("compaction")
375 && let Some(encrypted_content) = item
376 .get("encrypted_content")
377 .and_then(serde_json::Value::as_str)
378 {
379 compacted.push(ResponsesInputItem::Compaction {
380 encrypted_content: encrypted_content.to_string(),
381 });
382 }
383 }
384 Some("response.completed") => completed = true,
385 _ => {}
386 }
387 }
388
389 if !completed {
390 return Err(CompactionError::InvalidResponse(
391 "remote compaction stream ended before response.completed".to_string(),
392 ));
393 }
394 if compacted.len() != 1 {
395 return Err(CompactionError::InvalidResponse(format!(
396 "remote compaction expected exactly one compaction item, got {}",
397 compacted.len()
398 )));
399 }
400 Ok(compacted.pop().expect("validated one compaction item"))
401}
402
403fn build_compacted_history(
404 input: &[ResponsesInputItem],
405 compaction: ResponsesInputItem,
406) -> Vec<ResponsesInputItem> {
407 let retained = input
408 .iter()
409 .filter(|item| {
410 matches!(
411 item,
412 ResponsesInputItem::Message { role, .. }
413 if matches!(role.as_str(), "user" | "developer" | "system")
414 )
415 })
416 .cloned()
417 .collect::<Vec<_>>();
418 let mut retained = truncate_retained_messages(retained, RETAINED_MESSAGE_TOKEN_BUDGET);
419 retained.push(compaction);
420 retained
421}
422
423fn truncate_retained_messages(
424 items: Vec<ResponsesInputItem>,
425 max_tokens: u64,
426) -> Vec<ResponsesInputItem> {
427 let mut remaining = max_tokens;
428 let mut retained = Vec::new();
429 for item in items.into_iter().rev() {
430 if remaining == 0 {
431 break;
432 }
433 let tokens = message_tokens(&item).max(1);
434 if tokens <= remaining {
435 retained.push(item);
436 remaining -= tokens;
437 } else if let Some(item) = truncate_message(item, remaining) {
438 retained.push(item);
439 remaining = 0;
440 }
441 }
442 retained.reverse();
443 retained
444}
445
446fn message_tokens(item: &ResponsesInputItem) -> u64 {
447 let ResponsesInputItem::Message { content, .. } = item else {
448 return 0;
449 };
450 content
451 .iter()
452 .map(|part| match part {
453 ResponsesContentPart::InputText { text }
454 | ResponsesContentPart::OutputText { text } => text.len().div_ceil(4) as u64,
455 ResponsesContentPart::InputImage { .. } => 2_000,
456 })
457 .sum()
458}
459
460fn truncate_message(item: ResponsesInputItem, max_tokens: u64) -> Option<ResponsesInputItem> {
461 let ResponsesInputItem::Message { role, content } = item else {
462 return Some(item);
463 };
464 let mut remaining_chars = max_tokens.saturating_mul(4) as usize;
465 let mut truncated = Vec::new();
466 for part in content {
467 match part {
468 ResponsesContentPart::InputImage { .. } => truncated.push(part),
469 ResponsesContentPart::InputText { text } => {
470 let text = truncate_text(text, &mut remaining_chars);
471 if !text.is_empty() {
472 truncated.push(ResponsesContentPart::InputText { text });
473 }
474 }
475 ResponsesContentPart::OutputText { text } => {
476 let text = truncate_text(text, &mut remaining_chars);
477 if !text.is_empty() {
478 truncated.push(ResponsesContentPart::OutputText { text });
479 }
480 }
481 }
482 }
483 (!truncated.is_empty()).then_some(ResponsesInputItem::Message {
484 role,
485 content: truncated,
486 })
487}
488
489fn truncate_text(mut text: String, remaining_chars: &mut usize) -> String {
490 if *remaining_chars == 0 {
491 return String::new();
492 }
493 if text.len() > *remaining_chars {
494 let mut boundary = *remaining_chars;
495 while !text.is_char_boundary(boundary) {
496 boundary -= 1;
497 }
498 text.truncate(boundary);
499 }
500 *remaining_chars -= text.len();
501 text
502}
503
504fn serialized_size(items: &[ResponsesInputItem]) -> usize {
505 serde_json::to_vec(items).map_or(usize::MAX, |value| value.len())
506}
507
508fn state_size(session_id: &str, state: &CompactionState) -> usize {
509 let summary_len = match &state.phase {
510 CompactionPhase::Preparing | CompactionPhase::Unconfirmed => 0,
511 CompactionPhase::Anchored { portable_summary } => portable_summary.len(),
512 };
513 session_id.len() + state.model.len() + summary_len + serialized_size(&state.native_history)
514}
515
516fn now_ms() -> u64 {
517 SystemTime::now()
518 .duration_since(UNIX_EPOCH)
519 .unwrap_or_default()
520 .as_millis() as u64
521}
522
523fn update_total_bytes(registry: &mut CompactionRegistry) {
524 registry.total_bytes = registry
525 .states
526 .iter()
527 .map(|(session_id, state)| state_size(session_id, state))
528 .sum();
529}
530
531fn evict_states(registry: &mut CompactionRegistry, now: u64) {
532 registry
533 .states
534 .retain(|_, state| now.saturating_sub(state.updated_at) <= STATE_TTL_MS);
535 update_total_bytes(registry);
536 while registry.states.len() > MAX_STATES || registry.total_bytes > MAX_TOTAL_STATE_BYTES {
537 let oldest = registry
538 .states
539 .iter()
540 .min_by_key(|(_, state)| state.updated_at)
541 .map(|(session_id, _)| session_id.clone());
542 let Some(oldest) = oldest else {
543 break;
544 };
545 registry.states.remove(&oldest);
546 update_total_bytes(registry);
547 }
548}
549
550pub fn clear_all_compactions_for_tests() {
551 *REGISTRY.lock().unwrap() = None;
552}
553
554#[cfg(test)]
555mod tests {
556 use super::*;
557 use serde_json::json;
558
559 const SUMMARY: &str =
560 "portable summary with enough detail to identify this compacted conversation";
561 const STALE_SUMMARY: &str =
562 "stale portable summary from an older overlapping compaction attempt";
563 static TEST_REGISTRY_LOCK: Mutex<()> = Mutex::new(());
564
565 fn request(input: serde_json::Value) -> ResponsesRequest {
566 serde_json::from_value(json!({
567 "model": "gpt-5.6-sol",
568 "input": input,
569 "store": false,
570 "stream": true,
571 "parallel_tool_calls": false,
572 "client_metadata": {"lite":"true"},
573 "text": {"verbosity":"low"}
574 }))
575 .unwrap()
576 }
577
578 fn output(text: &str) -> Vec<ResponsesInputItem> {
579 serde_json::from_value(json!([{
580 "type":"message",
581 "role":"assistant",
582 "content":[{"type":"output_text","text":text}]
583 }]))
584 .unwrap()
585 }
586
587 fn stored_compaction(
588 session_id: &str,
589 native_history: Vec<ResponsesInputItem>,
590 ) -> CompactionAttempt {
591 let attempt = begin_compaction(session_id, "gpt-5.6-sol");
592 assert!(store_compaction(session_id, attempt, native_history));
593 attempt
594 }
595
596 #[test]
597 fn parses_exactly_one_completed_compaction_item() {
598 let body = b"data: {\"type\":\"response.output_item.done\",\"item\":{\"type\":\"compaction\",\"encrypted_content\":\"opaque\"}}\n\ndata: {\"type\":\"response.completed\",\"response\":{}}\n\n";
599 assert!(matches!(
600 parse_compaction_response(body).unwrap(),
601 ResponsesInputItem::Compaction { encrypted_content } if encrypted_content == "opaque"
602 ));
603 }
604
605 #[test]
606 fn rejects_incomplete_or_ambiguous_compaction_streams() {
607 let incomplete = b"data: {\"type\":\"response.output_item.done\",\"item\":{\"type\":\"compaction\",\"encrypted_content\":\"opaque\"}}\n\n";
608 assert!(
609 parse_compaction_response(incomplete)
610 .unwrap_err()
611 .to_string()
612 .contains("before response.completed")
613 );
614 let missing = b"data: {\"type\":\"response.completed\",\"response\":{}}\n\n";
615 assert!(
616 parse_compaction_response(missing)
617 .unwrap_err()
618 .to_string()
619 .contains("exactly one")
620 );
621 }
622
623 #[test]
624 fn replay_requires_activation_and_wrapped_summary_anchor() {
625 let _guard = TEST_REGISTRY_LOCK.lock().unwrap();
626 clear_all_compactions_for_tests();
627 let attempt = stored_compaction(
628 "session",
629 vec![ResponsesInputItem::Compaction {
630 encrypted_content: "opaque".to_string(),
631 }],
632 );
633 let next = request(json!([
634 {"type":"additional_tools","role":"developer","tools":[]},
635 {"type":"message","role":"developer","content":[{"type":"input_text","text":"instructions"}]},
636 {"type":"message","role":"user","content":[{"type":"input_text","text":format!("<summary>{SUMMARY}</summary>")}]},
637 {"type":"message","role":"user","content":[{"type":"input_text","text":"continue"}]}
638 ]));
639 assert!(apply_compaction_replay(Some("session"), &next).is_none());
640 assert!(activate_compaction(
641 Some("session"),
642 Some(attempt),
643 "gpt-5.6-sol",
644 &output(&format!(
645 "<analysis>summary preparation</analysis>\n<summary>\n{SUMMARY}\n</summary>"
646 ))
647 ));
648
649 let replay = apply_compaction_replay(Some("session"), &next)
650 .unwrap()
651 .request;
652 assert!(matches!(
653 replay.input[0],
654 ResponsesInputItem::AdditionalTools { .. }
655 ));
656 assert!(
657 matches!(replay.input[1], ResponsesInputItem::Message { ref role, .. } if role == "developer")
658 );
659 assert!(matches!(
660 replay.input[2],
661 ResponsesInputItem::Compaction { .. }
662 ));
663 assert_eq!(replay.client_metadata, next.client_metadata);
664 }
665
666 #[test]
667 fn stale_activation_cannot_anchor_newer_native_history() {
668 let _guard = TEST_REGISTRY_LOCK.lock().unwrap();
669 clear_all_compactions_for_tests();
670 let older = stored_compaction(
671 "session",
672 vec![ResponsesInputItem::Compaction {
673 encrypted_content: "older-native-history".to_string(),
674 }],
675 );
676 let newer = stored_compaction(
677 "session",
678 vec![ResponsesInputItem::Compaction {
679 encrypted_content: "newer-native-history".to_string(),
680 }],
681 );
682
683 assert!(!activate_compaction(
684 Some("session"),
685 Some(older),
686 "gpt-5.6-sol",
687 &output(STALE_SUMMARY),
688 ));
689 assert!(activate_compaction(
690 Some("session"),
691 Some(newer),
692 "gpt-5.6-sol",
693 &output(SUMMARY),
694 ));
695 }
696
697 #[test]
698 fn stale_store_cannot_replace_newer_compaction() {
699 let _guard = TEST_REGISTRY_LOCK.lock().unwrap();
700 clear_all_compactions_for_tests();
701 let older = begin_compaction("session", "gpt-5.6-sol");
702 let newer = stored_compaction(
703 "session",
704 vec![ResponsesInputItem::Compaction {
705 encrypted_content: "newer-native-history".to_string(),
706 }],
707 );
708
709 assert!(!store_compaction(
710 "session",
711 older,
712 vec![ResponsesInputItem::Compaction {
713 encrypted_content: "late-older-native-history".to_string(),
714 }],
715 ));
716 assert!(activate_compaction(
717 Some("session"),
718 Some(newer),
719 "gpt-5.6-sol",
720 &output(SUMMARY),
721 ));
722 let next = request(json!([
723 {"type":"message","role":"user","content":[{"type":"input_text","text":SUMMARY}]},
724 {"type":"message","role":"user","content":[{"type":"input_text","text":"continue"}]}
725 ]));
726 let replay = apply_compaction_replay(Some("session"), &next)
727 .unwrap()
728 .request;
729 assert!(replay.input.iter().any(|item| matches!(
730 item,
731 ResponsesInputItem::Compaction { encrypted_content }
732 if encrypted_content == "newer-native-history"
733 )));
734 }
735
736 #[test]
737 fn preparing_compaction_survives_model_mismatched_replay_check() {
738 let _guard = TEST_REGISTRY_LOCK.lock().unwrap();
739 clear_all_compactions_for_tests();
740 let attempt = begin_compaction("session", "gpt-5.6-sol");
741 let mut mismatched = request(json!([
742 {"type":"message","role":"user","content":[{"type":"input_text","text":"continue"}]}
743 ]));
744 mismatched.model = "gpt-5.6-terra".to_string();
745
746 assert!(apply_compaction_replay(Some("session"), &mismatched).is_none());
747 assert!(store_compaction(
748 "session",
749 attempt,
750 vec![ResponsesInputItem::Compaction {
751 encrypted_content: "native-history".to_string(),
752 }],
753 ));
754 }
755
756 #[test]
757 fn stale_abort_cannot_clear_newer_compaction() {
758 let _guard = TEST_REGISTRY_LOCK.lock().unwrap();
759 clear_all_compactions_for_tests();
760 let older = begin_compaction("session", "gpt-5.6-sol");
761 let newer = stored_compaction(
762 "session",
763 vec![ResponsesInputItem::Compaction {
764 encrypted_content: "newer-native-history".to_string(),
765 }],
766 );
767
768 abort_compaction_attempt(Some("session"), Some(older));
769
770 assert!(activate_compaction(
771 Some("session"),
772 Some(newer),
773 "gpt-5.6-sol",
774 &output(SUMMARY),
775 ));
776 }
777
778 #[test]
779 fn stale_replay_abort_cannot_clear_newer_compaction() {
780 let _guard = TEST_REGISTRY_LOCK.lock().unwrap();
781 clear_all_compactions_for_tests();
782 let older = stored_compaction(
783 "session",
784 vec![ResponsesInputItem::Compaction {
785 encrypted_content: "older-native-history".to_string(),
786 }],
787 );
788 assert!(activate_compaction(
789 Some("session"),
790 Some(older),
791 "gpt-5.6-sol",
792 &output(STALE_SUMMARY),
793 ));
794 let older_next = request(json!([
795 {"type":"message","role":"user","content":[{"type":"input_text","text":STALE_SUMMARY}]},
796 {"type":"message","role":"user","content":[{"type":"input_text","text":"continue old"}]}
797 ]));
798 let older_replay = apply_compaction_replay(Some("session"), &older_next).unwrap();
799
800 let newer = stored_compaction(
801 "session",
802 vec![ResponsesInputItem::Compaction {
803 encrypted_content: "newer-native-history".to_string(),
804 }],
805 );
806 assert!(activate_compaction(
807 Some("session"),
808 Some(newer),
809 "gpt-5.6-sol",
810 &output(SUMMARY),
811 ));
812
813 abort_compaction_attempt(Some("session"), Some(older_replay.attempt));
814
815 let newer_next = request(json!([
816 {"type":"message","role":"user","content":[{"type":"input_text","text":SUMMARY}]},
817 {"type":"message","role":"user","content":[{"type":"input_text","text":"continue new"}]}
818 ]));
819 let replay = apply_compaction_replay(Some("session"), &newer_next)
820 .unwrap()
821 .request;
822 assert!(replay.input.iter().any(|item| matches!(
823 item,
824 ResponsesInputItem::Compaction { encrypted_content }
825 if encrypted_content == "newer-native-history"
826 )));
827 }
828
829 #[test]
830 fn invalid_stale_summary_cannot_clear_newer_compaction() {
831 let _guard = TEST_REGISTRY_LOCK.lock().unwrap();
832 clear_all_compactions_for_tests();
833 let older = begin_compaction("session", "gpt-5.6-sol");
834 let newer = stored_compaction(
835 "session",
836 vec![ResponsesInputItem::Compaction {
837 encrypted_content: "newer-native-history".to_string(),
838 }],
839 );
840
841 assert!(!activate_compaction(
842 Some("session"),
843 Some(older),
844 "gpt-5.6-sol",
845 &output("too short"),
846 ));
847 assert!(activate_compaction(
848 Some("session"),
849 Some(newer),
850 "gpt-5.6-sol",
851 &output(SUMMARY),
852 ));
853 }
854
855 #[test]
856 fn replay_clears_on_missing_or_duplicate_anchor() {
857 let _guard = TEST_REGISTRY_LOCK.lock().unwrap();
858 for text in [
859 "different conversation without the expected summary".to_string(),
860 format!("{SUMMARY} and {SUMMARY}"),
861 ] {
862 clear_all_compactions_for_tests();
863 let attempt = stored_compaction(
864 "session",
865 vec![ResponsesInputItem::Compaction {
866 encrypted_content: "opaque".to_string(),
867 }],
868 );
869 activate_compaction(
870 Some("session"),
871 Some(attempt),
872 "gpt-5.6-sol",
873 &output(SUMMARY),
874 );
875 let changed = request(json!([
876 {"type":"message","role":"user","content":[{"type":"input_text","text":text}]},
877 {"type":"message","role":"user","content":[{"type":"input_text","text":"continue"}]}
878 ]));
879 assert!(apply_compaction_replay(Some("session"), &changed).is_none());
880 assert!(apply_compaction_replay(Some("session"), &changed).is_none());
881 }
882 }
883
884 #[test]
885 fn compacted_history_excludes_lite_envelope() {
886 let input: Vec<ResponsesInputItem> = serde_json::from_value(json!([
887 {"type":"additional_tools","role":"developer","tools":[]},
888 {"type":"message","role":"developer","content":[{"type":"input_text","text":"summarize"}]},
889 {"type":"message","role":"user","content":[{"type":"input_text","text":"remember me"}]}
890 ])).unwrap();
891 let (_, conversation) = split_input_envelope(&input);
892 let history = build_compacted_history(
893 conversation,
894 ResponsesInputItem::Compaction {
895 encrypted_content: "opaque".to_string(),
896 },
897 );
898 assert_eq!(history.len(), 2);
899 assert!(
900 matches!(history[0], ResponsesInputItem::Message { ref role, .. } if role == "user")
901 );
902 }
903
904 #[test]
905 fn failed_replay_clears_anchored_state() {
906 let _guard = TEST_REGISTRY_LOCK.lock().unwrap();
907 clear_all_compactions_for_tests();
908 let attempt = stored_compaction(
909 "failed-replay",
910 vec![ResponsesInputItem::Compaction {
911 encrypted_content: "opaque".to_string(),
912 }],
913 );
914 activate_compaction(
915 Some("failed-replay"),
916 Some(attempt),
917 "gpt-5.6-sol",
918 &output(SUMMARY),
919 );
920 let next = request(json!([
921 {"type":"message","role":"user","content":[{"type":"input_text","text":SUMMARY}]},
922 {"type":"message","role":"user","content":[{"type":"input_text","text":"continue"}]}
923 ]));
924 let replay = apply_compaction_replay(Some("failed-replay"), &next).unwrap();
925
926 abort_compaction_attempt(Some("failed-replay"), Some(replay.attempt));
927
928 assert!(apply_compaction_replay(Some("failed-replay"), &next).is_none());
929 }
930
931 #[test]
932 fn replay_clears_on_model_change() {
933 let _guard = TEST_REGISTRY_LOCK.lock().unwrap();
934 clear_all_compactions_for_tests();
935 let attempt = stored_compaction(
936 "session",
937 vec![ResponsesInputItem::Compaction {
938 encrypted_content: "opaque".to_string(),
939 }],
940 );
941 activate_compaction(
942 Some("session"),
943 Some(attempt),
944 "gpt-5.6-sol",
945 &output(SUMMARY),
946 );
947 let mut changed = request(json!([
948 {"type":"message","role":"user","content":[{"type":"input_text","text":SUMMARY}]},
949 {"type":"message","role":"user","content":[{"type":"input_text","text":"continue"}]}
950 ]));
951 changed.model = "gpt-5.4".to_string();
952 assert!(apply_compaction_replay(Some("session"), &changed).is_none());
953 }
954
955 #[test]
956 fn retained_history_obeys_token_budget_at_utf8_boundary() {
957 let text = "é".repeat((RETAINED_MESSAGE_TOKEN_BUDGET as usize + 10) * 4);
958 let input: Vec<ResponsesInputItem> = serde_json::from_value(json!([
959 {"type":"message","role":"user","content":[{"type":"input_text","text":text}]}
960 ]))
961 .unwrap();
962 let history = build_compacted_history(
963 &input,
964 ResponsesInputItem::Compaction {
965 encrypted_content: "opaque".to_string(),
966 },
967 );
968 let ResponsesInputItem::Message { content, .. } = &history[0] else {
969 panic!("expected retained message");
970 };
971 let ResponsesContentPart::InputText { text } = &content[0] else {
972 panic!("expected retained text");
973 };
974 assert!(text.len() <= RETAINED_MESSAGE_TOKEN_BUDGET as usize * 4);
975 assert!(std::str::from_utf8(text.as_bytes()).is_ok());
976 }
977}