Skip to main content

claude_codex/providers/codex/
compaction.rs

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}