Skip to main content

claude_codex/
session.rs

1use crate::config::AliasProvider;
2use crate::registry::normalize_incoming_model;
3use std::collections::{HashMap, VecDeque};
4use std::sync::{LazyLock, Mutex};
5
6const SESSION_IDLE_TTL_MS: u64 = 30 * 60 * 1000;
7pub const MAX_SESSIONS: usize = 10_000;
8
9#[derive(Debug, Clone)]
10pub struct SessionState {
11    pub seq: u64,
12    pub affinity_provider: Option<AliasProvider>,
13    pub last_seen: u64,
14}
15
16#[derive(Default)]
17struct SessionStore {
18    map: HashMap<String, SessionState>,
19    order: VecDeque<String>,
20}
21
22static SESSIONS: LazyLock<Mutex<SessionStore>> =
23    LazyLock::new(|| Mutex::new(SessionStore::default()));
24
25fn now_millis() -> u64 {
26    use std::time::{SystemTime, UNIX_EPOCH};
27    let dur = SystemTime::now()
28        .duration_since(UNIX_EPOCH)
29        .unwrap_or_default();
30    dur.as_millis() as u64
31}
32
33pub fn existing_session(session_id: Option<&str>, now: u64) -> Option<SessionState> {
34    let id = session_id?;
35    let mut store = SESSIONS.lock().expect("session lock");
36    let state = store.map.get(id).cloned()?;
37    if now.saturating_sub(state.last_seen) > SESSION_IDLE_TTL_MS {
38        store.map.remove(id);
39        store.order.retain(|item| item != id);
40        return None;
41    }
42    Some(state)
43}
44
45pub fn existing_session_now(session_id: Option<&str>) -> Option<SessionState> {
46    existing_session(session_id, now_millis())
47}
48
49pub fn record_session_request(
50    session_id: Option<&str>,
51    prior: Option<&SessionState>,
52    provider_name: &str,
53    model: &str,
54    now: u64,
55) -> Option<SessionState> {
56    record_session_request_with_affinity_update(session_id, prior, provider_name, model, true, now)
57}
58
59pub(crate) fn record_session_request_with_affinity_update(
60    session_id: Option<&str>,
61    prior: Option<&SessionState>,
62    provider_name: &str,
63    model: &str,
64    update_affinity: bool,
65    now: u64,
66) -> Option<SessionState> {
67    let id = session_id?;
68    let mut store = SESSIONS.lock().expect("session lock");
69    let stored = store
70        .map
71        .get(id)
72        .cloned()
73        .filter(|state| now.saturating_sub(state.last_seen) <= SESSION_IDLE_TTL_MS);
74    if stored.is_none() && store.map.remove(id).is_some() {
75        store.order.retain(|item| item != id);
76    }
77    let mut next = stored
78        .or_else(|| {
79            prior
80                .filter(|state| now.saturating_sub(state.last_seen) <= SESSION_IDLE_TTL_MS)
81                .cloned()
82        })
83        .unwrap_or(SessionState {
84            seq: 0,
85            affinity_provider: None,
86            last_seen: now,
87        });
88    next.seq += 1;
89    next.last_seen = now;
90    if update_affinity
91        && is_alias_routable_provider(provider_name)
92        && !crate::registry::is_anthropic_alias(normalize_incoming_model(model).as_str())
93    {
94        next.affinity_provider = Some(match provider_name {
95            "codex" => AliasProvider::Codex,
96            "kimi" => AliasProvider::Kimi,
97            _ => next.affinity_provider.unwrap_or(AliasProvider::Codex),
98        });
99    }
100
101    if !store.map.contains_key(id) {
102        store.order.push_back(id.to_string());
103    }
104    store.map.insert(id.to_string(), next.clone());
105
106    while store.order.len() > MAX_SESSIONS {
107        if let Some(evict) = store.order.pop_front() {
108            store.map.remove(&evict);
109        } else {
110            break;
111        }
112    }
113
114    Some(next)
115}
116
117fn is_alias_routable_provider(name: &str) -> bool {
118    matches!(name, "codex" | "kimi")
119}
120
121#[cfg(test)]
122pub fn reset_sessions_for_test() {
123    let mut store = SESSIONS.lock().expect("session lock");
124    store.map.clear();
125    store.order.clear();
126}
127
128pub fn affinity_provider_from_session(session: &SessionState) -> Option<AliasProvider> {
129    session.affinity_provider
130}
131
132#[cfg(test)]
133mod tests {
134    use super::*;
135
136    #[test]
137    fn concurrent_requests_increment_latest_sequence() {
138        let session_id = "session-concurrent-sequence-test";
139        let initial = record_session_request(Some(session_id), None, "codex", "gpt-5.6-sol", 1)
140            .expect("initial session");
141        let handles: Vec<_> = (0..16)
142            .map(|offset| {
143                let stale = initial.clone();
144                std::thread::spawn(move || {
145                    record_session_request(
146                        Some(session_id),
147                        Some(&stale),
148                        "codex",
149                        "gpt-5.6-sol",
150                        2 + offset,
151                    )
152                    .expect("recorded session")
153                    .seq
154                })
155            })
156            .collect();
157        let mut sequences: Vec<_> = handles
158            .into_iter()
159            .map(|handle| handle.join().expect("session thread"))
160            .collect();
161        sequences.sort_unstable();
162        assert_eq!(sequences, (2..=17).collect::<Vec<_>>());
163    }
164
165    #[test]
166    fn auxiliary_request_does_not_change_session_affinity() {
167        let session_id = "session-affinity-auxiliary-request-test";
168        let initial = record_session_request(Some(session_id), None, "codex", "gpt-5.6-sol", 1)
169            .expect("initial session");
170        assert_eq!(initial.affinity_provider, Some(AliasProvider::Codex));
171
172        let after_review = record_session_request_with_affinity_update(
173            Some(session_id),
174            Some(&initial),
175            "kimi",
176            "kimi-for-coding",
177            false,
178            2,
179        )
180        .expect("updated session");
181        assert_eq!(after_review.seq, 2);
182        assert_eq!(after_review.affinity_provider, Some(AliasProvider::Codex));
183    }
184}