Skip to main content

claude_codex/providers/codex/
continuation.rs

1use std::collections::HashMap;
2use std::sync::Mutex;
3use std::sync::atomic::{AtomicU64, Ordering};
4
5use crate::request_identity::ConversationIdentity;
6
7use super::translate::request::{ResponsesInputItem, ResponsesRequest};
8
9const TTL_MS: u64 = 30 * 60 * 1000;
10const MAX_STATES: usize = 10_000;
11const MAX_OWNER_TRANSCRIPT_BYTES: u64 = 2_000_000;
12const MAX_TOTAL_TRANSCRIPT_BYTES: u64 = 20_000_000;
13
14#[derive(Clone)]
15struct ContinuationState {
16    response_id: String,
17    socket_id: u64,
18    prompt_signature: String,
19    transcript: Vec<ResponsesInputItem>,
20    transcript_bytes: u64,
21    updated_at: u64,
22}
23
24struct OwnerState {
25    current_turn: u64,
26    continuation: Option<ContinuationState>,
27    updated_at: u64,
28}
29
30#[derive(Default)]
31struct ContinuationRegistry {
32    owners: HashMap<ConversationIdentity, OwnerState>,
33    total_transcript_bytes: u64,
34}
35
36static REGISTRY: Mutex<Option<ContinuationRegistry>> = Mutex::new(None);
37static NEXT_TURN_ID: AtomicU64 = AtomicU64::new(1);
38
39#[cfg(test)]
40static TEST_REGISTRY_LOCK: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(());
41
42#[cfg(test)]
43pub(crate) fn lock_continuation_registry_for_tests() -> tokio::sync::MutexGuard<'static, ()> {
44    TEST_REGISTRY_LOCK.blocking_lock()
45}
46
47#[cfg(test)]
48pub(crate) async fn lock_continuation_registry_for_async_tests()
49-> tokio::sync::MutexGuard<'static, ()> {
50    TEST_REGISTRY_LOCK.lock().await
51}
52
53#[derive(Clone)]
54pub struct ContinuationCandidate {
55    pub turn_id: Option<u64>,
56    pub previous_response_id: Option<String>,
57    pub input_delta: Option<Vec<ResponsesInputItem>>,
58    pub input_delta_count: usize,
59    pub disabled_reason: Option<String>,
60}
61
62#[derive(Clone)]
63pub(crate) struct ContinuationReservation {
64    candidate: ContinuationCandidate,
65    owner: Option<ConversationIdentity>,
66    origin_socket_id: Option<u64>,
67}
68
69impl ContinuationReservation {
70    pub(crate) fn new(
71        candidate: ContinuationCandidate,
72        owner: Option<ConversationIdentity>,
73        origin_socket_id: Option<u64>,
74    ) -> Self {
75        Self {
76            candidate,
77            owner,
78            origin_socket_id,
79        }
80    }
81
82    pub(crate) fn from_public_candidate(candidate: &ContinuationCandidate) -> Self {
83        Self::new(candidate.clone(), None, None)
84    }
85
86    pub(crate) fn for_owner_turn(
87        owner: Option<&ConversationIdentity>,
88        turn_id: Option<u64>,
89    ) -> Self {
90        Self::new(
91            ContinuationCandidate {
92                turn_id,
93                previous_response_id: None,
94                input_delta: None,
95                input_delta_count: 0,
96                disabled_reason: None,
97            },
98            owner.cloned(),
99            None,
100        )
101    }
102
103    pub(crate) fn candidate(&self) -> &ContinuationCandidate {
104        &self.candidate
105    }
106
107    pub(crate) fn owner(&self) -> Option<&ConversationIdentity> {
108        self.owner.as_ref()
109    }
110
111    pub(crate) fn turn_id(&self) -> Option<u64> {
112        self.candidate.turn_id
113    }
114
115    pub(crate) fn origin_socket_id(&self) -> Option<u64> {
116        self.origin_socket_id
117    }
118
119    pub(crate) fn into_candidate(self) -> ContinuationCandidate {
120        self.candidate
121    }
122
123    pub(crate) fn full_context_retry(&self) -> Self {
124        let mut candidate = self.candidate.clone();
125        candidate.previous_response_id = None;
126        candidate.input_delta = None;
127        candidate.disabled_reason = Some("full_context_retry".to_string());
128        Self::new(candidate, self.owner.clone(), None)
129    }
130}
131
132fn now_ms() -> u64 {
133    std::time::SystemTime::now()
134        .duration_since(std::time::UNIX_EPOCH)
135        .unwrap_or_default()
136        .as_millis() as u64
137}
138
139#[deprecated(note = "use the owner-aware provider flow for typed conversation ownership")]
140pub fn continuation_candidate(
141    session_id: Option<&str>,
142    body: &ResponsesRequest,
143    enabled: bool,
144) -> ContinuationCandidate {
145    let owner = session_id.map(|session_id| ConversationIdentity::Main(session_id.to_owned()));
146    continuation_candidate_inner(owner.as_ref(), body, enabled, "missing_session").into_candidate()
147}
148
149pub(crate) fn continuation_candidate_for_owner(
150    owner: Option<&ConversationIdentity>,
151    body: &ResponsesRequest,
152    enabled: bool,
153) -> ContinuationReservation {
154    continuation_candidate_inner(owner, body, enabled, "missing_identity")
155}
156
157fn continuation_candidate_inner(
158    owner: Option<&ConversationIdentity>,
159    body: &ResponsesRequest,
160    enabled: bool,
161    missing_owner_reason: &str,
162) -> ContinuationReservation {
163    if !enabled {
164        return ContinuationReservation::new(
165            ContinuationCandidate {
166                turn_id: None,
167                previous_response_id: None,
168                input_delta: None,
169                input_delta_count: body.input.len(),
170                disabled_reason: Some("disabled".to_string()),
171            },
172            owner.cloned(),
173            None,
174        );
175    }
176
177    let Some(owner) = owner else {
178        return ContinuationReservation::new(
179            ContinuationCandidate {
180                turn_id: None,
181                previous_response_id: None,
182                input_delta: None,
183                input_delta_count: body.input.len(),
184                disabled_reason: Some(missing_owner_reason.to_string()),
185            },
186            None,
187            None,
188        );
189    };
190
191    let turn_id = NEXT_TURN_ID.fetch_add(1, Ordering::Relaxed);
192    let now = now_ms();
193    let (state, superseded_turn) = {
194        let mut guard = REGISTRY.lock().unwrap();
195        let registry = guard.get_or_insert_with(ContinuationRegistry::default);
196        let existing = registry.owners.remove(owner);
197        let superseded_turn = existing.is_some();
198        let state = existing.and_then(|owner| owner.continuation);
199        if let Some(state) = &state {
200            registry.total_transcript_bytes = registry
201                .total_transcript_bytes
202                .saturating_sub(state.transcript_bytes);
203        }
204        registry.owners.insert(
205            owner.clone(),
206            OwnerState {
207                current_turn: turn_id,
208                continuation: None,
209                updated_at: now,
210            },
211        );
212        evict_oldest(registry);
213        (state, superseded_turn)
214    };
215
216    continuation_candidate_from_state(owner, turn_id, body, state, superseded_turn, now)
217}
218
219fn continuation_candidate_from_state(
220    owner: &ConversationIdentity,
221    turn_id: u64,
222    body: &ResponsesRequest,
223    state: Option<ContinuationState>,
224    superseded_turn: bool,
225    now: u64,
226) -> ContinuationReservation {
227    let state = match state {
228        Some(state) if now.saturating_sub(state.updated_at) <= TTL_MS => state,
229        Some(_) | None => {
230            return ContinuationReservation::new(
231                ContinuationCandidate {
232                    turn_id: Some(turn_id),
233                    previous_response_id: None,
234                    input_delta: None,
235                    input_delta_count: body.input.len(),
236                    disabled_reason: Some(if superseded_turn {
237                        "superseded_turn".to_string()
238                    } else {
239                        "missing_state".to_string()
240                    }),
241                },
242                Some(owner.clone()),
243                None,
244            );
245        }
246    };
247
248    let signature = prompt_signature(body);
249    if signature != state.prompt_signature {
250        return ContinuationReservation::new(
251            ContinuationCandidate {
252                turn_id: Some(turn_id),
253                previous_response_id: None,
254                input_delta: None,
255                input_delta_count: body.input.len(),
256                disabled_reason: Some("prompt_changed".to_string()),
257            },
258            Some(owner.clone()),
259            None,
260        );
261    }
262
263    let Some(suffix) = input_suffix_after_prefix(&body.input, &state.transcript) else {
264        return ContinuationReservation::new(
265            ContinuationCandidate {
266                turn_id: Some(turn_id),
267                previous_response_id: None,
268                input_delta: None,
269                input_delta_count: body.input.len(),
270                disabled_reason: Some("not_append_only".to_string()),
271            },
272            Some(owner.clone()),
273            None,
274        );
275    };
276
277    if suffix.is_empty() {
278        return ContinuationReservation::new(
279            ContinuationCandidate {
280                turn_id: Some(turn_id),
281                previous_response_id: None,
282                input_delta: None,
283                input_delta_count: 0,
284                disabled_reason: Some("empty_delta".to_string()),
285            },
286            Some(owner.clone()),
287            None,
288        );
289    }
290
291    ContinuationReservation::new(
292        ContinuationCandidate {
293            turn_id: Some(turn_id),
294            previous_response_id: Some(state.response_id),
295            input_delta_count: suffix.len(),
296            input_delta: Some(suffix),
297            disabled_reason: None,
298        },
299        Some(owner.clone()),
300        Some(state.socket_id),
301    )
302}
303
304#[deprecated(note = "recording without typed socket provenance is not reusable")]
305pub fn record_continuation(
306    session_id: Option<&str>,
307    turn_id: Option<u64>,
308    request_body: &ResponsesRequest,
309    response_id: Option<&str>,
310    output_items: &[ResponsesInputItem],
311) {
312    let owner = session_id.map(|session_id| ConversationIdentity::Main(session_id.to_owned()));
313    let reservation = ContinuationReservation::new(
314        ContinuationCandidate {
315            turn_id,
316            previous_response_id: None,
317            input_delta: None,
318            input_delta_count: request_body.input.len(),
319            disabled_reason: Some("legacy_recording_without_socket".to_string()),
320        },
321        owner,
322        None,
323    );
324    record_continuation_for_owner(&reservation, request_body, response_id, None, output_items);
325}
326
327pub(crate) fn record_continuation_for_owner(
328    reservation: &ContinuationReservation,
329    request_body: &ResponsesRequest,
330    response_id: Option<&str>,
331    socket_id: Option<u64>,
332    output_items: &[ResponsesInputItem],
333) {
334    let (owner, turn_id) = match (reservation.owner(), reservation.turn_id()) {
335        (Some(owner), Some(turn_id)) => (owner, turn_id),
336        _ => return,
337    };
338
339    let (response_id, socket_id) = match (response_id, socket_id) {
340        (Some(response_id), Some(socket_id)) if socket_id != 0 => {
341            (response_id.to_string(), socket_id)
342        }
343        _ => {
344            abort_continuation_inner(Some(owner), Some(turn_id));
345            return;
346        }
347    };
348    let mut transcript: Vec<ResponsesInputItem> = request_body.input.clone();
349    transcript.extend_from_slice(output_items);
350
351    let transcript_json = serde_json::to_string(&transcript).unwrap_or_default();
352    let transcript_bytes = transcript_json.len() as u64;
353
354    if transcript_bytes > MAX_OWNER_TRANSCRIPT_BYTES {
355        abort_continuation_inner(Some(owner), Some(turn_id));
356        return;
357    }
358
359    let state = ContinuationState {
360        response_id,
361        socket_id,
362        prompt_signature: prompt_signature(request_body),
363        transcript,
364        transcript_bytes,
365        updated_at: now_ms(),
366    };
367
368    let mut guard = REGISTRY.lock().unwrap();
369    let Some(registry) = guard.as_mut() else {
370        return;
371    };
372    let Some(owner_state) = registry.owners.get_mut(owner) else {
373        return;
374    };
375    if owner_state.current_turn != turn_id {
376        return;
377    }
378    if let Some(existing) = owner_state.continuation.replace(state) {
379        registry.total_transcript_bytes = registry
380            .total_transcript_bytes
381            .saturating_sub(existing.transcript_bytes);
382    }
383    registry.total_transcript_bytes += transcript_bytes;
384    evict_oldest(registry);
385}
386
387#[deprecated(note = "use the owner-aware provider flow for typed conversation ownership")]
388pub fn abort_continuation(session_id: Option<&str>, turn_id: Option<u64>) {
389    let owner = session_id.map(|session_id| ConversationIdentity::Main(session_id.to_owned()));
390    abort_continuation_inner(owner.as_ref(), turn_id);
391}
392
393pub(crate) fn abort_continuation_for_owner(reservation: &ContinuationReservation) {
394    abort_continuation_inner(reservation.owner(), reservation.turn_id());
395}
396
397fn abort_continuation_inner(owner: Option<&ConversationIdentity>, turn_id: Option<u64>) {
398    let (Some(owner), Some(turn_id)) = (owner, turn_id) else {
399        return;
400    };
401    let mut guard = REGISTRY.lock().unwrap();
402    let Some(registry) = guard.as_mut() else {
403        return;
404    };
405    if registry
406        .owners
407        .get(owner)
408        .is_some_and(|state| state.current_turn == turn_id)
409        && let Some(state) = registry.owners.remove(owner)
410        && let Some(continuation) = state.continuation
411    {
412        registry.total_transcript_bytes = registry
413            .total_transcript_bytes
414            .saturating_sub(continuation.transcript_bytes);
415    }
416}
417
418#[deprecated(note = "use the owner-aware provider flow for typed conversation ownership")]
419pub fn if_current_turn<T>(
420    session_id: Option<&str>,
421    turn_id: Option<u64>,
422    action: impl FnOnce() -> T,
423) -> Option<T> {
424    let owner = session_id.map(|session_id| ConversationIdentity::Main(session_id.to_owned()));
425    if_current_turn_inner(owner.as_ref(), turn_id, action)
426}
427
428pub(crate) fn if_current_turn_for_owner<T>(
429    reservation: &ContinuationReservation,
430    action: impl FnOnce() -> T,
431) -> Option<T> {
432    if_current_turn_inner(reservation.owner(), reservation.turn_id(), action)
433}
434
435fn if_current_turn_inner<T>(
436    owner: Option<&ConversationIdentity>,
437    turn_id: Option<u64>,
438    action: impl FnOnce() -> T,
439) -> Option<T> {
440    let (Some(owner), Some(turn_id)) = (owner, turn_id) else {
441        return None;
442    };
443    let guard = REGISTRY.lock().unwrap();
444    let current = guard
445        .as_ref()
446        .and_then(|registry| registry.owners.get(owner))
447        .is_some_and(|state| state.current_turn == turn_id);
448    current.then(action)
449}
450
451#[deprecated(note = "use the owner-aware provider flow for typed conversation ownership")]
452pub fn with_current_turn(
453    session_id: Option<&str>,
454    turn_id: Option<u64>,
455    action: impl FnOnce(),
456) -> bool {
457    let owner = session_id.map(|session_id| ConversationIdentity::Main(session_id.to_owned()));
458    if_current_turn_inner(owner.as_ref(), turn_id, action).is_some()
459}
460
461pub(crate) fn with_current_turn_for_owner(
462    reservation: &ContinuationReservation,
463    action: impl FnOnce(),
464) -> bool {
465    if_current_turn_for_owner(reservation, action).is_some()
466}
467
468#[deprecated(note = "use the owner-aware provider flow for typed conversation ownership")]
469pub fn is_current_turn(session_id: Option<&str>, turn_id: Option<u64>) -> bool {
470    let owner = session_id.map(|session_id| ConversationIdentity::Main(session_id.to_owned()));
471    is_current_turn_inner(owner.as_ref(), turn_id)
472}
473
474#[allow(dead_code)]
475pub(crate) fn is_current_turn_for_owner(reservation: &ContinuationReservation) -> bool {
476    is_current_turn_inner(reservation.owner(), reservation.turn_id())
477}
478
479fn is_current_turn_inner(owner: Option<&ConversationIdentity>, turn_id: Option<u64>) -> bool {
480    let (Some(owner), Some(turn_id)) = (owner, turn_id) else {
481        return false;
482    };
483    let guard = REGISTRY.lock().unwrap();
484    guard
485        .as_ref()
486        .and_then(|registry| registry.owners.get(owner))
487        .is_some_and(|state| state.current_turn == turn_id)
488}
489
490#[deprecated(note = "use the owner-aware provider flow for typed conversation ownership")]
491pub fn clear_continuation(session_id: Option<&str>) {
492    let owner = session_id.map(|session_id| ConversationIdentity::Main(session_id.to_owned()));
493    clear_continuation_for_owner(owner.as_ref());
494}
495
496pub(crate) fn clear_continuation_for_owner(owner: Option<&ConversationIdentity>) {
497    let Some(owner) = owner else {
498        return;
499    };
500    let mut guard = REGISTRY.lock().unwrap();
501    let Some(registry) = guard.as_mut() else {
502        return;
503    };
504    if let Some(state) = registry.owners.remove(owner)
505        && let Some(continuation) = state.continuation
506    {
507        registry.total_transcript_bytes = registry
508            .total_transcript_bytes
509            .saturating_sub(continuation.transcript_bytes);
510    }
511}
512
513#[deprecated(note = "use the owner-aware test helper for typed conversation ownership")]
514pub fn has_continuation_for_tests(session_id: &str) -> bool {
515    let owner = ConversationIdentity::Main(session_id.to_owned());
516    has_continuation_for_owner_for_tests(&owner)
517}
518
519pub(crate) fn has_continuation_for_owner_for_tests(owner: &ConversationIdentity) -> bool {
520    let guard = REGISTRY.lock().unwrap();
521    guard
522        .as_ref()
523        .and_then(|registry| registry.owners.get(owner))
524        .is_some_and(|state| state.continuation.is_some())
525}
526
527pub fn clear_all_continuations_for_tests() {
528    let mut guard = REGISTRY.lock().unwrap();
529    *guard = None;
530}
531
532fn input_suffix_after_prefix(
533    input: &[ResponsesInputItem],
534    prefix: &[ResponsesInputItem],
535) -> Option<Vec<ResponsesInputItem>> {
536    if prefix.len() > input.len() {
537        return None;
538    }
539    for i in 0..prefix.len() {
540        let a = serde_json::to_value(&input[i]).unwrap_or_default();
541        let b = serde_json::to_value(&prefix[i]).unwrap_or_default();
542        if a != b {
543            return None;
544        }
545    }
546    Some(input[prefix.len()..].to_vec())
547}
548
549fn prompt_signature(body: &ResponsesRequest) -> String {
550    let value = serde_json::to_value(body).unwrap_or_default();
551    let obj = match value.as_object() {
552        Some(o) => o,
553        None => return String::new(),
554    };
555    let mut entries: Vec<(&String, &serde_json::Value)> =
556        obj.iter().filter(|(k, _)| *k != "input").collect();
557    entries.sort_by_key(|(a, _)| *a);
558    let mut sig = String::from("{");
559    for (i, (key, val)) in entries.iter().enumerate() {
560        if i > 0 {
561            sig.push(',');
562        }
563        sig.push_str(&format!("\"{}\":{}", key, stable_json(val)));
564    }
565    sig.push('}');
566    sig
567}
568
569fn stable_json(value: &serde_json::Value) -> String {
570    match value {
571        serde_json::Value::Null => "null".to_string(),
572        serde_json::Value::Bool(b) => b.to_string(),
573        serde_json::Value::Number(n) => n.to_string(),
574        serde_json::Value::String(s) => serde_json::to_string(s).unwrap_or_default(),
575        serde_json::Value::Array(arr) => {
576            let items: Vec<String> = arr.iter().map(stable_json).collect();
577            format!("[{}]", items.join(","))
578        }
579        serde_json::Value::Object(obj) => {
580            let mut entries: Vec<(&String, &serde_json::Value)> = obj.iter().collect();
581            entries.sort_by_key(|(a, _)| *a);
582            let items: Vec<String> = entries
583                .iter()
584                .map(|(k, v)| {
585                    format!(
586                        "{}:{}",
587                        serde_json::to_string(k).unwrap_or_default(),
588                        stable_json(v)
589                    )
590                })
591                .collect();
592            format!("{{{}}}", items.join(","))
593        }
594    }
595}
596
597fn evict_oldest(registry: &mut ContinuationRegistry) {
598    while registry.owners.len() > MAX_STATES
599        || registry.total_transcript_bytes > MAX_TOTAL_TRANSCRIPT_BYTES
600    {
601        let owner = registry
602            .owners
603            .iter()
604            .min_by_key(|(_, state)| state.updated_at)
605            .map(|(owner, _)| owner.clone());
606        let Some(owner) = owner else {
607            break;
608        };
609        if let Some(state) = registry.owners.remove(&owner)
610            && let Some(continuation) = state.continuation
611        {
612            registry.total_transcript_bytes = registry
613                .total_transcript_bytes
614                .saturating_sub(continuation.transcript_bytes);
615        }
616    }
617}
618
619#[cfg(test)]
620mod tests {
621    use super::*;
622    use serde_json::json;
623
624    fn lock_registry() -> tokio::sync::MutexGuard<'static, ()> {
625        let guard = lock_continuation_registry_for_tests();
626        clear_all_continuations_for_tests();
627        guard
628    }
629
630    fn main_owner(session_id: &str) -> ConversationIdentity {
631        ConversationIdentity::Main(session_id.to_string())
632    }
633
634    fn agent_owner(session_id: &str, agent_id: &str) -> ConversationIdentity {
635        ConversationIdentity::Agent(session_id.to_string(), agent_id.to_string())
636    }
637
638    fn input(text: &str) -> ResponsesInputItem {
639        ResponsesInputItem::Message {
640            role: "user".to_string(),
641            content: vec![
642                super::super::translate::request::ResponsesContentPart::InputText {
643                    text: text.to_string(),
644                },
645            ],
646        }
647    }
648
649    fn request_with_input(
650        input: Vec<ResponsesInputItem>,
651        extra: Option<serde_json::Value>,
652    ) -> ResponsesRequest {
653        let mut fields = serde_json::Map::new();
654        fields.insert("model".into(), json!("gpt-5.5"));
655        fields.insert("input".into(), json!(input));
656        fields.insert("store".into(), json!(false));
657        fields.insert("stream".into(), json!(true));
658        fields.insert("text".into(), json!({"verbosity": "low"}));
659        fields.insert("parallel_tool_calls".into(), json!(true));
660        if let Some(extras) = extra
661            && let Some(obj) = extras.as_object()
662        {
663            for (key, value) in obj {
664                fields.insert(key.clone(), value.clone());
665            }
666        }
667        serde_json::from_value(serde_json::Value::Object(fields)).unwrap()
668    }
669
670    fn start_and_record(
671        owner: &ConversationIdentity,
672        request: &ResponsesRequest,
673        response_id: &str,
674    ) {
675        let reservation = continuation_candidate_for_owner(Some(owner), request, true);
676        record_continuation_for_owner(&reservation, request, Some(response_id), Some(1), &[]);
677    }
678
679    #[test]
680    #[allow(deprecated)]
681    fn disabled_and_missing_identity_requests_are_stateless() {
682        let _registry_guard = lock_registry();
683        let request = request_with_input(vec![input("one")], None);
684        let owner = main_owner("session-a");
685
686        let disabled = continuation_candidate_for_owner(Some(&owner), &request, false);
687        assert_eq!(disabled.owner(), Some(&owner));
688        assert_eq!(disabled.turn_id(), None);
689        assert_eq!(disabled.candidate().input_delta_count, request.input.len());
690        assert_eq!(
691            disabled.candidate().disabled_reason.as_deref(),
692            Some("disabled")
693        );
694
695        let missing = continuation_candidate_for_owner(None, &request, true);
696        assert_eq!(missing.owner(), None);
697        assert_eq!(missing.turn_id(), None);
698        assert_eq!(missing.candidate().input_delta_count, request.input.len());
699        assert_eq!(
700            missing.candidate().disabled_reason.as_deref(),
701            Some("missing_identity")
702        );
703
704        let legacy_missing = continuation_candidate(None, &request, true);
705        assert_eq!(
706            legacy_missing.disabled_reason.as_deref(),
707            Some("missing_session")
708        );
709    }
710
711    #[test]
712    fn sibling_agents_reserve_and_publish_independently() {
713        let _registry_guard = lock_registry();
714        let sibling_one = agent_owner("session-a", "agent-one");
715        let sibling_two = agent_owner("session-a", "agent-two");
716        let first_request = request_with_input(vec![input("one")], None);
717
718        let first = continuation_candidate_for_owner(Some(&sibling_one), &first_request, true);
719        let second = continuation_candidate_for_owner(Some(&sibling_two), &first_request, true);
720        assert_ne!(first.turn_id(), second.turn_id());
721        record_continuation_for_owner(&first, &first_request, Some("resp_one"), Some(11), &[]);
722        record_continuation_for_owner(&second, &first_request, Some("resp_two"), Some(22), &[]);
723        assert!(has_continuation_for_owner_for_tests(&sibling_one));
724        assert!(has_continuation_for_owner_for_tests(&sibling_two));
725
726        let next_request = request_with_input(vec![input("one"), input("two")], None);
727        let first_next = continuation_candidate_for_owner(Some(&sibling_one), &next_request, true);
728        let second_next = continuation_candidate_for_owner(Some(&sibling_two), &next_request, true);
729        assert_eq!(
730            first_next.candidate().previous_response_id.as_deref(),
731            Some("resp_one")
732        );
733        assert_eq!(first_next.origin_socket_id(), Some(11));
734        assert_eq!(
735            second_next.candidate().previous_response_id.as_deref(),
736            Some("resp_two")
737        );
738        assert_eq!(second_next.origin_socket_id(), Some(22));
739    }
740
741    #[test]
742    fn different_owner_completion_order_cannot_interfere() {
743        let _registry_guard = lock_registry();
744        let main = main_owner("session-a");
745        let agent = agent_owner("session-a", "agent-a");
746        let request = request_with_input(vec![input("one")], None);
747        let main_reservation = continuation_candidate_for_owner(Some(&main), &request, true);
748        let agent_reservation = continuation_candidate_for_owner(Some(&agent), &request, true);
749
750        record_continuation_for_owner(
751            &agent_reservation,
752            &request,
753            Some("resp_agent"),
754            Some(1),
755            &[],
756        );
757        record_continuation_for_owner(&main_reservation, &request, Some("resp_main"), Some(1), &[]);
758        abort_continuation_for_owner(&ContinuationReservation::new(
759            ContinuationCandidate {
760                turn_id: agent_reservation.turn_id(),
761                previous_response_id: None,
762                input_delta: None,
763                input_delta_count: 0,
764                disabled_reason: None,
765            },
766            Some(main.clone()),
767            None,
768        ));
769
770        assert!(has_continuation_for_owner_for_tests(&main));
771        assert!(has_continuation_for_owner_for_tests(&agent));
772        let next = request_with_input(vec![input("one"), input("two")], None);
773        assert_eq!(
774            continuation_candidate_for_owner(Some(&agent), &next, true)
775                .candidate()
776                .previous_response_id
777                .as_deref(),
778            Some("resp_agent")
779        );
780    }
781
782    #[test]
783    fn missing_response_id_aborts_only_the_current_owner() {
784        let _registry_guard = lock_registry();
785        let owner = main_owner("session-a");
786        let sibling = agent_owner("session-a", "agent-a");
787        let request = request_with_input(vec![input("one")], None);
788        start_and_record(&owner, &request, "resp_main");
789        start_and_record(&sibling, &request, "resp_agent");
790
791        let reservation = continuation_candidate_for_owner(Some(&owner), &request, true);
792        record_continuation_for_owner(&reservation, &request, None, Some(1), &[]);
793
794        assert!(!has_continuation_for_owner_for_tests(&owner));
795        assert!(has_continuation_for_owner_for_tests(&sibling));
796    }
797
798    #[test]
799    fn missing_socket_id_does_not_publish_reusable_state() {
800        let _registry_guard = lock_registry();
801        let owner = main_owner("session-no-socket");
802        let request = request_with_input(vec![input("one")], None);
803        let reservation = continuation_candidate_for_owner(Some(&owner), &request, true);
804
805        record_continuation_for_owner(
806            &reservation,
807            &request,
808            Some("resp_without_socket"),
809            None,
810            &[],
811        );
812
813        assert!(!has_continuation_for_owner_for_tests(&owner));
814        let next = continuation_candidate_for_owner(Some(&owner), &request, true);
815        assert_eq!(next.candidate().previous_response_id, None);
816        assert_eq!(next.origin_socket_id(), None);
817    }
818
819    #[test]
820    #[allow(deprecated)]
821    fn legacy_recording_without_provenance_publishes_no_reusable_state() {
822        let _registry_guard = lock_registry();
823        let session_id = "legacy-no-provenance";
824        let request = request_with_input(vec![input("one")], None);
825        let candidate = continuation_candidate(Some(session_id), &request, true);
826
827        record_continuation(
828            Some(session_id),
829            candidate.turn_id,
830            &request,
831            Some("resp_legacy"),
832            &[],
833        );
834
835        assert!(!has_continuation_for_tests(session_id));
836        let next = continuation_candidate(Some(session_id), &request, true);
837        assert_eq!(next.previous_response_id, None);
838    }
839
840    #[test]
841    fn same_owner_stale_turn_cannot_publish_clear_or_run_actions() {
842        let _registry_guard = lock_registry();
843        let owner = main_owner("session-a");
844        let request = request_with_input(vec![input("one")], None);
845        start_and_record(&owner, &request, "resp_1");
846
847        let stale = continuation_candidate_for_owner(Some(&owner), &request, true);
848        let current = continuation_candidate_for_owner(Some(&owner), &request, true);
849        assert_eq!(
850            current.candidate().disabled_reason.as_deref(),
851            Some("superseded_turn")
852        );
853        record_continuation_for_owner(&stale, &request, Some("resp_stale"), Some(1), &[]);
854        assert!(!has_continuation_for_owner_for_tests(&owner));
855        record_continuation_for_owner(&current, &request, Some("resp_current"), Some(1), &[]);
856        assert!(has_continuation_for_owner_for_tests(&owner));
857        abort_continuation_for_owner(&stale);
858        assert!(has_continuation_for_owner_for_tests(&owner));
859
860        let mut ran = false;
861        assert_eq!(if_current_turn_for_owner(&stale, || ran = true), None);
862        assert!(!ran);
863    }
864
865    #[test]
866    #[allow(deprecated)]
867    fn missing_owner_or_turn_mutations_are_hard_noops() {
868        let _registry_guard = lock_registry();
869        let owner = main_owner("session-a");
870        let request = request_with_input(vec![input("one")], None);
871        start_and_record(&owner, &request, "resp_1");
872
873        let missing_owner = ContinuationReservation::new(
874            ContinuationCandidate {
875                turn_id: Some(1),
876                previous_response_id: None,
877                input_delta: None,
878                input_delta_count: 1,
879                disabled_reason: None,
880            },
881            None,
882            None,
883        );
884        let missing_turn = ContinuationReservation::new(
885            ContinuationCandidate {
886                turn_id: None,
887                previous_response_id: None,
888                input_delta: None,
889                input_delta_count: 1,
890                disabled_reason: None,
891            },
892            Some(owner.clone()),
893            None,
894        );
895        record_continuation_for_owner(&missing_owner, &request, Some("ignored"), Some(1), &[]);
896        record_continuation_for_owner(&missing_turn, &request, Some("ignored"), Some(1), &[]);
897        abort_continuation_for_owner(&missing_owner);
898        abort_continuation_for_owner(&missing_turn);
899        clear_continuation_for_owner(None);
900        assert!(has_continuation_for_owner_for_tests(&owner));
901
902        let mut runs = 0;
903        assert_eq!(
904            if_current_turn_for_owner(&missing_owner, || runs += 1),
905            None
906        );
907        assert_eq!(if_current_turn_for_owner(&missing_turn, || runs += 1), None);
908        assert!(!with_current_turn_for_owner(&missing_owner, || runs += 1));
909        assert!(!with_current_turn_for_owner(&missing_turn, || runs += 1));
910        assert_eq!(if_current_turn(None, Some(1), || runs += 1), None);
911        assert!(!with_current_turn(None, Some(1), || runs += 1));
912        assert_eq!(runs, 0);
913    }
914
915    #[test]
916    fn append_only_and_prompt_guards_stay_owner_scoped() {
917        let _registry_guard = lock_registry();
918        let owner = main_owner("session-a");
919        let request = request_with_input(vec![input("one")], None);
920        start_and_record(&owner, &request, "resp_1");
921
922        let appended = request_with_input(vec![input("one"), input("two")], None);
923        let reservation = continuation_candidate_for_owner(Some(&owner), &appended, true);
924        assert_eq!(
925            reservation.candidate().previous_response_id.as_deref(),
926            Some("resp_1")
927        );
928        assert_eq!(reservation.candidate().input_delta_count, 1);
929
930        let full_context = reservation.full_context_retry();
931        assert_eq!(full_context.owner(), Some(&owner));
932        assert_eq!(full_context.turn_id(), reservation.turn_id());
933        assert_eq!(full_context.candidate().previous_response_id, None);
934        assert!(full_context.candidate().input_delta.is_none());
935        assert_eq!(full_context.origin_socket_id(), None);
936
937        record_continuation_for_owner(&reservation, &appended, Some("resp_2"), Some(1), &[]);
938        let changed = request_with_input(
939            vec![input("one"), input("two"), input("three")],
940            Some(json!({"service_tier": "flex"})),
941        );
942        let reservation = continuation_candidate_for_owner(Some(&owner), &changed, true);
943        assert_eq!(
944            reservation.candidate().disabled_reason.as_deref(),
945            Some("prompt_changed")
946        );
947        assert!(!has_continuation_for_owner_for_tests(&owner));
948    }
949}