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(¤t, &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}