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}