Skip to main content

claude_codex/
monitor.rs

1use std::{
2    collections::{HashMap, VecDeque},
3    path::PathBuf,
4    sync::{Arc, Mutex},
5    time::{Duration, Instant, SystemTime},
6};
7
8const DEFAULT_RECENT_LIMIT: usize = 200;
9
10#[derive(Debug, Clone, Copy, PartialEq, Eq)]
11pub enum EndpointKind {
12    Messages,
13    CountTokens,
14}
15
16impl EndpointKind {
17    pub fn label(self) -> &'static str {
18        match self {
19            Self::Messages => "messages",
20            Self::CountTokens => "count_tokens",
21        }
22    }
23}
24
25#[derive(Debug, Clone, PartialEq, Eq)]
26pub enum RequestStatus {
27    Started,
28    ProviderSelected,
29    Upstream,
30    Streaming,
31    Completed,
32    Failed,
33}
34
35impl RequestStatus {
36    pub fn label(&self) -> &'static str {
37        match self {
38            Self::Started => "started",
39            Self::ProviderSelected => "selected",
40            Self::Upstream => "upstream",
41            Self::Streaming => "streaming",
42            Self::Completed => "completed",
43            Self::Failed => "failed",
44        }
45    }
46}
47
48#[derive(Debug, Clone)]
49pub enum MonitorEvent {
50    RequestStarted {
51        request_id: String,
52        session_id: Option<String>,
53        session_seq: Option<u64>,
54        endpoint: EndpointKind,
55    },
56    ProviderSelected {
57        request_id: String,
58        provider: String,
59        model: String,
60        effort: Option<String>,
61    },
62    ModelResolved {
63        request_id: String,
64        model: String,
65    },
66    UpstreamStarted {
67        request_id: String,
68    },
69    TrafficCapturePath {
70        request_id: String,
71        path: PathBuf,
72    },
73    StreamProgress {
74        request_id: String,
75        bytes: u64,
76        chunks: u64,
77        input_tokens: Option<u64>,
78        output_tokens: Option<u64>,
79    },
80    UsageUpdated {
81        request_id: String,
82        input_tokens: Option<u64>,
83        output_tokens: Option<u64>,
84    },
85    RequestCompleted {
86        request_id: String,
87        http_status: u16,
88        input_tokens: Option<u64>,
89        output_tokens: Option<u64>,
90    },
91    RequestFailed {
92        request_id: String,
93        http_status: Option<u16>,
94        error: String,
95    },
96    RequestAbandoned {
97        request_id: String,
98        error: String,
99    },
100}
101
102#[derive(Debug, Clone)]
103pub struct ActiveRequest {
104    pub request_id: String,
105    pub session_id: Option<String>,
106    pub session_seq: Option<u64>,
107    pub provider: Option<String>,
108    pub model: Option<String>,
109    pub effort: Option<String>,
110    pub endpoint: EndpointKind,
111    pub started_at: SystemTime,
112    started_instant: Instant,
113    pub status: RequestStatus,
114    pub streamed_bytes: u64,
115    pub stream_chunks: u64,
116    pub input_tokens: Option<u64>,
117    pub output_tokens: Option<u64>,
118    pub error: Option<String>,
119    pub traffic_capture_path: Option<PathBuf>,
120}
121
122impl ActiveRequest {
123    pub fn elapsed(&self) -> Duration {
124        self.started_instant.elapsed()
125    }
126
127    pub fn rate(&self) -> Throughput {
128        throughput(
129            self.output_tokens,
130            self.streamed_bytes,
131            self.stream_chunks,
132            self.elapsed(),
133        )
134    }
135}
136
137#[derive(Debug, Clone)]
138pub struct CompletedRequest {
139    pub request_id: String,
140    pub session_id: Option<String>,
141    pub session_seq: Option<u64>,
142    pub provider: Option<String>,
143    pub model: Option<String>,
144    pub effort: Option<String>,
145    pub endpoint: EndpointKind,
146    pub started_at: SystemTime,
147    pub finished_at: SystemTime,
148    pub status: RequestStatus,
149    pub http_status: Option<u16>,
150    pub latency: Duration,
151    pub streamed_bytes: u64,
152    pub stream_chunks: u64,
153    pub input_tokens: Option<u64>,
154    pub output_tokens: Option<u64>,
155    pub error: Option<String>,
156    pub traffic_capture_path: Option<PathBuf>,
157}
158
159impl CompletedRequest {
160    pub fn rate(&self) -> Throughput {
161        throughput(
162            self.output_tokens,
163            self.streamed_bytes,
164            self.stream_chunks,
165            self.latency,
166        )
167    }
168}
169
170#[derive(Debug, Clone, PartialEq)]
171pub enum Throughput {
172    TokensPerSecond(f64),
173    BytesPerSecond(f64),
174    EventsPerSecond(f64),
175    None,
176}
177
178impl Throughput {
179    pub fn label(&self) -> String {
180        match self {
181            Self::TokensPerSecond(value) => format!("{value:.1} tok/s"),
182            Self::BytesPerSecond(value) if *value >= 1024.0 => {
183                format!("{:.1} KB/s", value / 1024.0)
184            }
185            Self::BytesPerSecond(value) => format!("{value:.0} B/s"),
186            Self::EventsPerSecond(value) => format!("{value:.1} ev/s"),
187            Self::None => "-".to_string(),
188        }
189    }
190}
191
192#[derive(Debug, Clone)]
193pub struct MonitorState {
194    pub started_at: SystemTime,
195    pub sessions: Vec<SessionSummary>,
196    pub active: Vec<ActiveRequest>,
197    pub recent: Vec<CompletedRequest>,
198}
199
200#[derive(Debug, Clone)]
201pub struct SessionSummary {
202    pub session_id: Option<String>,
203    pub active_count: usize,
204    pub request_count: usize,
205    pub failure_count: usize,
206    pub provider: Option<String>,
207    pub model: Option<String>,
208    pub effort: Option<String>,
209    pub last_seen: SystemTime,
210    pub input_tokens: u64,
211    pub output_tokens: u64,
212    pub elapsed: Duration,
213    pub last_status: String,
214}
215
216impl SessionSummary {
217    pub fn rate(&self) -> Throughput {
218        throughput(
219            Some(self.output_tokens).filter(|tokens| *tokens > 0),
220            0,
221            0,
222            self.elapsed,
223        )
224    }
225
226    pub fn label(&self) -> String {
227        self.session_id
228            .clone()
229            .unwrap_or_else(|| "no-session".to_string())
230    }
231}
232
233#[derive(Debug)]
234struct MonitorStore {
235    started_at: SystemTime,
236    active: HashMap<String, ActiveRequest>,
237    recent: VecDeque<CompletedRequest>,
238    recent_limit: usize,
239}
240
241#[derive(Debug, Clone)]
242pub struct MonitorHandle {
243    store: Arc<Mutex<MonitorStore>>,
244}
245
246impl Default for MonitorHandle {
247    fn default() -> Self {
248        Self::new(DEFAULT_RECENT_LIMIT)
249    }
250}
251
252impl MonitorHandle {
253    pub fn new(recent_limit: usize) -> Self {
254        Self {
255            store: Arc::new(Mutex::new(MonitorStore {
256                started_at: SystemTime::now(),
257                active: HashMap::new(),
258                recent: VecDeque::new(),
259                recent_limit,
260            })),
261        }
262    }
263
264    pub fn publish(&self, event: MonitorEvent) {
265        if let Ok(mut store) = self.store.lock() {
266            store.apply(event);
267        }
268    }
269
270    pub fn snapshot(&self) -> MonitorState {
271        match self.store.lock() {
272            Ok(store) => store.snapshot(),
273            Err(_) => MonitorState {
274                started_at: SystemTime::now(),
275                sessions: Vec::new(),
276                active: Vec::new(),
277                recent: Vec::new(),
278            },
279        }
280    }
281
282    pub fn request_started(
283        &self,
284        request_id: impl Into<String>,
285        session_id: Option<String>,
286        session_seq: Option<u64>,
287        endpoint: EndpointKind,
288    ) {
289        self.publish(MonitorEvent::RequestStarted {
290            request_id: request_id.into(),
291            session_id,
292            session_seq,
293            endpoint,
294        });
295    }
296
297    pub fn provider_selected(
298        &self,
299        request_id: impl Into<String>,
300        provider: impl Into<String>,
301        model: impl Into<String>,
302        effort: Option<String>,
303    ) {
304        self.publish(MonitorEvent::ProviderSelected {
305            request_id: request_id.into(),
306            provider: provider.into(),
307            model: model.into(),
308            effort,
309        });
310    }
311
312    pub fn model_resolved(&self, request_id: impl Into<String>, model: impl Into<String>) {
313        self.publish(MonitorEvent::ModelResolved {
314            request_id: request_id.into(),
315            model: model.into(),
316        });
317    }
318
319    pub fn upstream_started(&self, request_id: impl Into<String>) {
320        self.publish(MonitorEvent::UpstreamStarted {
321            request_id: request_id.into(),
322        });
323    }
324
325    pub fn traffic_capture_path(&self, request_id: impl Into<String>, path: PathBuf) {
326        self.publish(MonitorEvent::TrafficCapturePath {
327            request_id: request_id.into(),
328            path,
329        });
330    }
331
332    pub fn stream_progress(
333        &self,
334        request_id: impl Into<String>,
335        bytes: u64,
336        chunks: u64,
337        input_tokens: Option<u64>,
338        output_tokens: Option<u64>,
339    ) {
340        self.publish(MonitorEvent::StreamProgress {
341            request_id: request_id.into(),
342            bytes,
343            chunks,
344            input_tokens,
345            output_tokens,
346        });
347    }
348
349    pub fn usage_updated(
350        &self,
351        request_id: impl Into<String>,
352        input_tokens: Option<u64>,
353        output_tokens: Option<u64>,
354    ) {
355        self.publish(MonitorEvent::UsageUpdated {
356            request_id: request_id.into(),
357            input_tokens,
358            output_tokens,
359        });
360    }
361
362    pub fn request_completed(
363        &self,
364        request_id: impl Into<String>,
365        http_status: u16,
366        input_tokens: Option<u64>,
367        output_tokens: Option<u64>,
368    ) {
369        self.publish(MonitorEvent::RequestCompleted {
370            request_id: request_id.into(),
371            http_status,
372            input_tokens,
373            output_tokens,
374        });
375    }
376
377    pub fn request_failed(
378        &self,
379        request_id: impl Into<String>,
380        http_status: Option<u16>,
381        error: impl Into<String>,
382    ) {
383        self.publish(MonitorEvent::RequestFailed {
384            request_id: request_id.into(),
385            http_status,
386            error: error.into(),
387        });
388    }
389
390    pub fn request_abandoned(&self, request_id: impl Into<String>, error: impl Into<String>) {
391        self.publish(MonitorEvent::RequestAbandoned {
392            request_id: request_id.into(),
393            error: error.into(),
394        });
395    }
396}
397
398impl MonitorStore {
399    fn apply(&mut self, event: MonitorEvent) {
400        match event {
401            MonitorEvent::RequestStarted {
402                request_id,
403                session_id,
404                session_seq,
405                endpoint,
406            } => {
407                self.active.insert(
408                    request_id.clone(),
409                    ActiveRequest {
410                        request_id,
411                        session_id,
412                        session_seq,
413                        provider: None,
414                        model: None,
415                        effort: None,
416                        endpoint,
417                        started_at: SystemTime::now(),
418                        started_instant: Instant::now(),
419                        status: RequestStatus::Started,
420                        streamed_bytes: 0,
421                        stream_chunks: 0,
422                        input_tokens: None,
423                        output_tokens: None,
424                        error: None,
425                        traffic_capture_path: None,
426                    },
427                );
428            }
429            MonitorEvent::ProviderSelected {
430                request_id,
431                provider,
432                model,
433                effort,
434            } => {
435                if let Some(active) = self.active.get_mut(&request_id) {
436                    active.provider = Some(provider);
437                    active.model = Some(model);
438                    active.effort = effort;
439                    active.status = RequestStatus::ProviderSelected;
440                }
441            }
442            MonitorEvent::ModelResolved { request_id, model } => {
443                if let Some(active) = self.active.get_mut(&request_id) {
444                    active.model = Some(match active.model.take() {
445                        Some(incoming) if incoming != model => format!("{incoming} → {model}"),
446                        Some(incoming) => incoming,
447                        None => model,
448                    });
449                }
450            }
451            MonitorEvent::UpstreamStarted { request_id } => {
452                if let Some(active) = self.active.get_mut(&request_id) {
453                    active.status = RequestStatus::Upstream;
454                }
455            }
456            MonitorEvent::TrafficCapturePath { request_id, path } => {
457                if let Some(active) = self.active.get_mut(&request_id) {
458                    active.traffic_capture_path = Some(path);
459                }
460            }
461            MonitorEvent::StreamProgress {
462                request_id,
463                bytes,
464                chunks,
465                input_tokens,
466                output_tokens,
467            } => {
468                if let Some(active) = self.active.get_mut(&request_id) {
469                    active.status = RequestStatus::Streaming;
470                    active.streamed_bytes = active.streamed_bytes.saturating_add(bytes);
471                    active.stream_chunks = active.stream_chunks.saturating_add(chunks);
472                    active.input_tokens = input_tokens.or(active.input_tokens);
473                    active.output_tokens = output_tokens.or(active.output_tokens);
474                } else if let Some(completed) = self
475                    .recent
476                    .iter_mut()
477                    .find(|request| request.request_id == request_id)
478                {
479                    completed.streamed_bytes = completed.streamed_bytes.saturating_add(bytes);
480                    completed.stream_chunks = completed.stream_chunks.saturating_add(chunks);
481                    completed.input_tokens = input_tokens.or(completed.input_tokens);
482                    completed.output_tokens = output_tokens.or(completed.output_tokens);
483                    completed.finished_at = SystemTime::now();
484                    completed.latency = completed
485                        .finished_at
486                        .duration_since(completed.started_at)
487                        .unwrap_or(completed.latency);
488                }
489            }
490            MonitorEvent::UsageUpdated {
491                request_id,
492                input_tokens,
493                output_tokens,
494            } => {
495                if let Some(active) = self.active.get_mut(&request_id) {
496                    active.input_tokens = input_tokens.or(active.input_tokens);
497                    active.output_tokens = output_tokens.or(active.output_tokens);
498                } else if let Some(completed) = self
499                    .recent
500                    .iter_mut()
501                    .find(|request| request.request_id == request_id)
502                {
503                    completed.input_tokens = input_tokens.or(completed.input_tokens);
504                    completed.output_tokens = output_tokens.or(completed.output_tokens);
505                }
506            }
507            MonitorEvent::RequestCompleted {
508                request_id,
509                http_status,
510                input_tokens,
511                output_tokens,
512            } => {
513                self.finish(
514                    &request_id,
515                    RequestStatus::Completed,
516                    Some(http_status),
517                    input_tokens,
518                    output_tokens,
519                    None,
520                );
521            }
522            MonitorEvent::RequestFailed {
523                request_id,
524                http_status,
525                error,
526            } => {
527                self.finish(
528                    &request_id,
529                    RequestStatus::Failed,
530                    http_status,
531                    None,
532                    None,
533                    Some(error),
534                );
535            }
536            MonitorEvent::RequestAbandoned { request_id, error } => {
537                self.finish_active(
538                    &request_id,
539                    RequestStatus::Failed,
540                    None,
541                    None,
542                    None,
543                    Some(error),
544                );
545            }
546        }
547    }
548
549    fn finish_active(
550        &mut self,
551        request_id: &str,
552        status: RequestStatus,
553        http_status: Option<u16>,
554        input_tokens: Option<u64>,
555        output_tokens: Option<u64>,
556        error: Option<String>,
557    ) {
558        if self.active.contains_key(request_id) {
559            self.finish(
560                request_id,
561                status,
562                http_status,
563                input_tokens,
564                output_tokens,
565                error,
566            );
567        }
568    }
569
570    fn finish(
571        &mut self,
572        request_id: &str,
573        status: RequestStatus,
574        http_status: Option<u16>,
575        input_tokens: Option<u64>,
576        output_tokens: Option<u64>,
577        error: Option<String>,
578    ) {
579        let active = self
580            .active
581            .remove(request_id)
582            .unwrap_or_else(|| ActiveRequest {
583                request_id: request_id.to_string(),
584                session_id: None,
585                session_seq: None,
586                provider: None,
587                model: None,
588                effort: None,
589                endpoint: EndpointKind::Messages,
590                started_at: SystemTime::now(),
591                started_instant: Instant::now(),
592                status: RequestStatus::Started,
593                streamed_bytes: 0,
594                stream_chunks: 0,
595                input_tokens: None,
596                output_tokens: None,
597                error: None,
598                traffic_capture_path: None,
599            });
600        let completed = CompletedRequest {
601            request_id: active.request_id,
602            session_id: active.session_id,
603            session_seq: active.session_seq,
604            provider: active.provider,
605            model: active.model,
606            effort: active.effort,
607            endpoint: active.endpoint,
608            started_at: active.started_at,
609            finished_at: SystemTime::now(),
610            status,
611            http_status,
612            latency: active.started_instant.elapsed(),
613            streamed_bytes: active.streamed_bytes,
614            stream_chunks: active.stream_chunks,
615            input_tokens: input_tokens.or(active.input_tokens),
616            output_tokens: output_tokens.or(active.output_tokens),
617            error: error.or(active.error),
618            traffic_capture_path: active.traffic_capture_path,
619        };
620        self.recent.push_front(completed);
621        while self.recent.len() > self.recent_limit {
622            self.recent.pop_back();
623        }
624    }
625
626    fn snapshot(&self) -> MonitorState {
627        let mut active: Vec<_> = self.active.values().cloned().collect();
628        active.sort_by_key(|request| request.started_at);
629        let sessions = session_summaries(&active, &self.recent);
630        MonitorState {
631            started_at: self.started_at,
632            sessions,
633            active,
634            recent: self.recent.iter().cloned().collect(),
635        }
636    }
637}
638
639fn session_summaries(
640    active: &[ActiveRequest],
641    recent: &VecDeque<CompletedRequest>,
642) -> Vec<SessionSummary> {
643    let mut sessions: HashMap<Option<String>, SessionSummary> = HashMap::new();
644    for request in recent.iter().rev() {
645        let entry = sessions
646            .entry(request.session_id.clone())
647            .or_insert_with(|| SessionSummary {
648                session_id: request.session_id.clone(),
649                active_count: 0,
650                request_count: 0,
651                failure_count: 0,
652                provider: None,
653                model: None,
654                effort: None,
655                last_seen: request.finished_at,
656                input_tokens: 0,
657                output_tokens: 0,
658                elapsed: Duration::ZERO,
659                last_status: "-".to_string(),
660            });
661        entry.request_count += 1;
662        if request.status == RequestStatus::Failed {
663            entry.failure_count += 1;
664        }
665        entry.provider = request.provider.clone().or(entry.provider.clone());
666        entry.model = request.model.clone().or(entry.model.clone());
667        entry.effort = request.effort.clone().or(entry.effort.clone());
668        entry.last_seen = max_system_time(entry.last_seen, request.finished_at);
669        entry.input_tokens = entry
670            .input_tokens
671            .saturating_add(request.input_tokens.unwrap_or(0));
672        entry.output_tokens = entry
673            .output_tokens
674            .saturating_add(request.output_tokens.unwrap_or(0));
675        entry.elapsed = entry.elapsed.saturating_add(request.latency);
676        entry.last_status = request.status.label().to_string();
677    }
678
679    for request in active {
680        let entry = sessions
681            .entry(request.session_id.clone())
682            .or_insert_with(|| SessionSummary {
683                session_id: request.session_id.clone(),
684                active_count: 0,
685                request_count: 0,
686                failure_count: 0,
687                provider: None,
688                model: None,
689                effort: None,
690                last_seen: request.started_at,
691                input_tokens: 0,
692                output_tokens: 0,
693                elapsed: Duration::ZERO,
694                last_status: "-".to_string(),
695            });
696        entry.active_count += 1;
697        entry.request_count += 1;
698        entry.provider = request.provider.clone().or(entry.provider.clone());
699        entry.model = request.model.clone().or(entry.model.clone());
700        entry.effort = request.effort.clone().or(entry.effort.clone());
701        entry.last_seen = max_system_time(entry.last_seen, request.started_at);
702        entry.input_tokens = entry
703            .input_tokens
704            .saturating_add(request.input_tokens.unwrap_or(0));
705        entry.output_tokens = entry
706            .output_tokens
707            .saturating_add(request.output_tokens.unwrap_or(0));
708        entry.elapsed = entry.elapsed.saturating_add(request.elapsed());
709        entry.last_status = request.status.label().to_string();
710    }
711
712    let mut out: Vec<_> = sessions.into_values().collect();
713    out.sort_by_key(SessionSummary::label);
714    out
715}
716
717fn max_system_time(left: SystemTime, right: SystemTime) -> SystemTime {
718    if right.duration_since(left).is_ok() {
719        right
720    } else {
721        left
722    }
723}
724
725pub fn throughput(
726    output_tokens: Option<u64>,
727    streamed_bytes: u64,
728    stream_chunks: u64,
729    elapsed: Duration,
730) -> Throughput {
731    let secs = elapsed.as_secs_f64();
732    if secs <= 0.0 {
733        return Throughput::None;
734    }
735    if let Some(tokens) = output_tokens.filter(|tokens| *tokens > 0) {
736        return Throughput::TokensPerSecond(tokens as f64 / secs);
737    }
738    if streamed_bytes > 0 {
739        return Throughput::BytesPerSecond(streamed_bytes as f64 / secs);
740    }
741    if stream_chunks > 0 {
742        return Throughput::EventsPerSecond(stream_chunks as f64 / secs);
743    }
744    Throughput::None
745}
746
747pub fn usage_from_anthropic_sse(bytes: &[u8]) -> (Option<u64>, Option<u64>) {
748    let text = String::from_utf8_lossy(bytes);
749    let mut input_tokens = None;
750    let mut output_tokens = None;
751    for line in text.lines() {
752        let Some(data) = line.strip_prefix("data:") else {
753            continue;
754        };
755        let Ok(value) = serde_json::from_str::<serde_json::Value>(data.trim()) else {
756            continue;
757        };
758        for usage in [
759            value.pointer("/usage"),
760            value.pointer("/delta/usage"),
761            value.pointer("/message/usage"),
762        ]
763        .into_iter()
764        .flatten()
765        {
766            if let Some(tokens) = usage.get("input_tokens").and_then(|value| value.as_u64()) {
767                input_tokens = Some(tokens);
768            }
769            if let Some(tokens) = usage.get("output_tokens").and_then(|value| value.as_u64()) {
770                output_tokens = Some(tokens);
771            }
772        }
773    }
774    (input_tokens, output_tokens)
775}
776
777#[cfg(test)]
778mod tests {
779    use super::*;
780
781    #[test]
782    fn started_requests_appear_active() {
783        let monitor = MonitorHandle::new(10);
784        monitor.request_started(
785            "r1",
786            Some("s1".to_string()),
787            Some(3),
788            EndpointKind::Messages,
789        );
790        let state = monitor.snapshot();
791        assert_eq!(state.active.len(), 1);
792        assert_eq!(state.active[0].request_id, "r1");
793        assert_eq!(state.active[0].session_id.as_deref(), Some("s1"));
794        assert_eq!(state.active[0].session_seq, Some(3));
795    }
796
797    #[test]
798    fn resolved_model_appends_to_incoming_alias() {
799        let monitor = MonitorHandle::new(10);
800        monitor.request_started("r1", None, None, EndpointKind::Messages);
801        monitor.provider_selected("r1", "codex", "claude-sonnet-4-6", None);
802        monitor.model_resolved("r1", "gpt-5.4");
803
804        let state = monitor.snapshot();
805        assert_eq!(
806            state.active[0].model.as_deref(),
807            Some("claude-sonnet-4-6 → gpt-5.4")
808        );
809    }
810
811    #[test]
812    fn identical_resolved_model_is_shown_once() {
813        let monitor = MonitorHandle::new(10);
814        monitor.request_started("r1", None, None, EndpointKind::Messages);
815        monitor.provider_selected("r1", "codex", "gpt-5.6-sol", None);
816        monitor.model_resolved("r1", "gpt-5.6-sol");
817
818        let state = monitor.snapshot();
819        assert_eq!(state.active[0].model.as_deref(), Some("gpt-5.6-sol"));
820    }
821
822    #[test]
823    fn active_stream_progress_reports_token_rate() {
824        let monitor = MonitorHandle::new(10);
825        monitor.request_started("r1", None, None, EndpointKind::Messages);
826        monitor.stream_progress("r1", 50, 1, Some(1_225), Some(141));
827
828        let state = monitor.snapshot();
829        assert_eq!(state.active.len(), 1);
830        assert!(matches!(
831            state.active[0].rate(),
832            Throughput::TokensPerSecond(_)
833        ));
834    }
835
836    #[test]
837    fn late_stream_usage_updates_completed_request() {
838        let monitor = MonitorHandle::new(10);
839        monitor.request_started("r1", None, None, EndpointKind::Messages);
840        monitor.stream_progress("r1", 100, 1, Some(0), Some(0));
841        monitor.request_completed("r1", 200, None, None);
842        monitor.stream_progress("r1", 50, 1, Some(1_225), Some(141));
843
844        let state = monitor.snapshot();
845        assert!(state.active.is_empty());
846        assert_eq!(state.recent[0].streamed_bytes, 150);
847        assert_eq!(state.recent[0].stream_chunks, 2);
848        assert_eq!(state.recent[0].input_tokens, Some(1_225));
849        assert_eq!(state.recent[0].output_tokens, Some(141));
850        assert!(matches!(
851            state.recent[0].rate(),
852            Throughput::TokensPerSecond(_)
853        ));
854    }
855
856    #[test]
857    fn completed_requests_leave_active_and_enter_recent() {
858        let monitor = MonitorHandle::new(10);
859        monitor.request_started("r1", None, None, EndpointKind::Messages);
860        monitor.provider_selected("r1", "codex", "gpt-5.5", Some("high".to_string()));
861        monitor.request_completed("r1", 200, Some(10), Some(20));
862        let state = monitor.snapshot();
863        assert!(state.active.is_empty());
864        assert_eq!(state.recent.len(), 1);
865        assert_eq!(state.recent[0].provider.as_deref(), Some("codex"));
866        assert_eq!(state.recent[0].effort.as_deref(), Some("high"));
867        assert_eq!(state.recent[0].output_tokens, Some(20));
868    }
869
870    #[test]
871    fn failed_requests_preserve_error_summary() {
872        let monitor = MonitorHandle::new(10);
873        monitor.request_started("r1", None, None, EndpointKind::Messages);
874        monitor.request_failed("r1", Some(400), "Unknown model");
875        let state = monitor.snapshot();
876        assert_eq!(state.recent[0].status, RequestStatus::Failed);
877        assert_eq!(state.recent[0].http_status, Some(400));
878        assert_eq!(state.recent[0].error.as_deref(), Some("Unknown model"));
879    }
880
881    #[test]
882    fn abandoned_requests_leave_active_once() {
883        let monitor = MonitorHandle::new(10);
884        monitor.request_started("r1", None, None, EndpointKind::Messages);
885        monitor.request_abandoned("r1", "request dropped");
886        monitor.request_abandoned("r1", "request dropped again");
887        let state = monitor.snapshot();
888        assert!(state.active.is_empty());
889        assert_eq!(state.recent.len(), 1);
890        assert_eq!(state.recent[0].status, RequestStatus::Failed);
891        assert_eq!(state.recent[0].http_status, None);
892        assert_eq!(state.recent[0].error.as_deref(), Some("request dropped"));
893    }
894
895    #[test]
896    fn completed_requests_ignore_late_abandonment() {
897        let monitor = MonitorHandle::new(10);
898        monitor.request_started("r1", None, None, EndpointKind::Messages);
899        monitor.request_completed("r1", 200, None, None);
900        monitor.request_abandoned("r1", "request dropped");
901        let state = monitor.snapshot();
902        assert!(state.active.is_empty());
903        assert_eq!(state.recent.len(), 1);
904        assert_eq!(state.recent[0].status, RequestStatus::Completed);
905    }
906
907    #[test]
908    fn bounded_recent_history_drops_oldest() {
909        let monitor = MonitorHandle::new(2);
910        for id in ["r1", "r2", "r3"] {
911            monitor.request_started(id, None, None, EndpointKind::Messages);
912            monitor.request_completed(id, 200, None, None);
913        }
914        let state = monitor.snapshot();
915        let ids: Vec<_> = state
916            .recent
917            .iter()
918            .map(|request| request.request_id.as_str())
919            .collect();
920        assert_eq!(ids, vec!["r3", "r2"]);
921    }
922
923    #[test]
924    fn throughput_selects_best_available_signal() {
925        let elapsed = Duration::from_secs(2);
926        assert_eq!(
927            throughput(Some(84), 1024, 10, elapsed),
928            Throughput::TokensPerSecond(42.0)
929        );
930        assert_eq!(
931            throughput(None, 2048, 10, elapsed),
932            Throughput::BytesPerSecond(1024.0)
933        );
934        assert_eq!(
935            throughput(None, 0, 36, elapsed),
936            Throughput::EventsPerSecond(18.0)
937        );
938    }
939
940    #[test]
941    fn sse_usage_extracts_final_message_delta_tokens() {
942        let sse = br#"event: message_start
943data: {"type":"message_start","message":{"usage":{"input_tokens":0,"output_tokens":0}}}
944
945event: message_delta
946data: {"type":"message_delta","delta":{"stop_reason":"end_turn"},"usage":{"input_tokens":12,"output_tokens":48}}
947
948"#;
949        assert_eq!(usage_from_anthropic_sse(sse), (Some(12), Some(48)));
950    }
951
952    #[test]
953    fn session_summaries_group_recent_and_active_requests() {
954        let monitor = MonitorHandle::new(10);
955        monitor.request_started(
956            "r1",
957            Some("s1".to_string()),
958            Some(1),
959            EndpointKind::Messages,
960        );
961        monitor.provider_selected("r1", "codex", "gpt-5.5", None);
962        monitor.request_completed("r1", 200, Some(10), Some(20));
963        monitor.request_started(
964            "r2",
965            Some("s1".to_string()),
966            Some(2),
967            EndpointKind::Messages,
968        );
969        monitor.provider_selected("r2", "codex", "gpt-5.5", Some("xhigh".to_string()));
970        let state = monitor.snapshot();
971        assert_eq!(state.sessions.len(), 1);
972        assert_eq!(state.sessions[0].label(), "s1");
973        assert_eq!(state.sessions[0].request_count, 2);
974        assert_eq!(state.sessions[0].active_count, 1);
975        assert_eq!(state.sessions[0].effort.as_deref(), Some("xhigh"));
976        assert_eq!(state.sessions[0].output_tokens, 20);
977    }
978
979    #[test]
980    fn session_order_is_stable_across_activity() {
981        let monitor = MonitorHandle::new(10);
982        monitor.request_started(
983            "r1",
984            Some("session-b".to_string()),
985            Some(1),
986            EndpointKind::Messages,
987        );
988        monitor.request_started(
989            "r2",
990            Some("session-a".to_string()),
991            Some(1),
992            EndpointKind::Messages,
993        );
994
995        let first: Vec<_> = monitor
996            .snapshot()
997            .sessions
998            .iter()
999            .map(SessionSummary::label)
1000            .collect();
1001        monitor.request_completed("r1", 200, None, None);
1002        monitor.request_started(
1003            "r3",
1004            Some("session-b".to_string()),
1005            Some(2),
1006            EndpointKind::Messages,
1007        );
1008        let second: Vec<_> = monitor
1009            .snapshot()
1010            .sessions
1011            .iter()
1012            .map(SessionSummary::label)
1013            .collect();
1014
1015        assert_eq!(first, vec!["session-a", "session-b"]);
1016        assert_eq!(second, first);
1017    }
1018}