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
8mod mock;
9
10pub use mock::{MockMonitor, mock_state};
11
12const DEFAULT_RECENT_LIMIT: usize = 200;
13pub const SESSION_TOKEN_BUCKET_SECS: u64 = 10;
14
15#[derive(Debug, Clone, Copy, PartialEq, Eq)]
16pub enum EndpointKind {
17    Messages,
18    CountTokens,
19    Responses,
20    ChatCompletions,
21    Images,
22    Transcriptions,
23}
24
25impl EndpointKind {
26    pub fn label(self) -> &'static str {
27        match self {
28            Self::Messages => "messages",
29            Self::CountTokens => "count_tokens",
30            Self::Responses => "responses",
31            Self::ChatCompletions => "chat_completions",
32            Self::Images => "images",
33            Self::Transcriptions => "transcriptions",
34        }
35    }
36}
37
38#[derive(Debug, Clone, PartialEq, Eq)]
39pub enum RequestStatus {
40    Started,
41    ProviderSelected,
42    Compacting,
43    Upstream,
44    Streaming,
45    Completed,
46    Failed,
47}
48
49impl RequestStatus {
50    pub fn label(&self) -> &'static str {
51        match self {
52            Self::Started => "started",
53            Self::ProviderSelected => "selected",
54            Self::Compacting => "compacting",
55            Self::Upstream => "upstream",
56            Self::Streaming => "streaming",
57            Self::Completed => "completed",
58            Self::Failed => "failed",
59        }
60    }
61}
62
63#[derive(Debug, Clone)]
64pub enum MonitorEvent {
65    RequestStarted {
66        request_id: String,
67        session_id: Option<String>,
68        session_seq: Option<u64>,
69        endpoint: EndpointKind,
70    },
71    ProjectResolved {
72        request_id: String,
73        project: String,
74    },
75    SessionSequenceResolved {
76        request_id: String,
77        session_seq: u64,
78    },
79    ProviderSelected {
80        request_id: String,
81        provider: String,
82        model: String,
83        effort: Option<String>,
84    },
85    ModelResolved {
86        request_id: String,
87        model: String,
88    },
89    CompactionStarted {
90        request_id: String,
91    },
92    UpstreamStarted {
93        request_id: String,
94    },
95    GenerationStarted {
96        request_id: String,
97    },
98    TrafficCapturePath {
99        request_id: String,
100        path: PathBuf,
101    },
102    StreamProgress {
103        request_id: String,
104        bytes: u64,
105        chunks: u64,
106        input_tokens: Option<u64>,
107        output_tokens: Option<u64>,
108    },
109    UsageUpdated {
110        request_id: String,
111        input_tokens: Option<u64>,
112        output_tokens: Option<u64>,
113    },
114    RequestCompleted {
115        request_id: String,
116        http_status: u16,
117        input_tokens: Option<u64>,
118        output_tokens: Option<u64>,
119    },
120    RequestFailed {
121        request_id: String,
122        http_status: Option<u16>,
123        error: String,
124    },
125    RequestAbandoned {
126        request_id: String,
127        error: String,
128    },
129}
130
131#[derive(Debug, Clone)]
132pub struct ActiveRequest {
133    pub request_id: String,
134    pub session_id: Option<String>,
135    pub session_seq: Option<u64>,
136    pub project: Option<String>,
137    pub provider: Option<String>,
138    pub model: Option<String>,
139    pub effort: Option<String>,
140    pub endpoint: EndpointKind,
141    pub started_at: SystemTime,
142    started_instant: Instant,
143    pub generation_started_at: Option<SystemTime>,
144    generation_started_instant: Option<Instant>,
145    generation_initial_output_tokens: u64,
146    pub generation_finished_at: Option<SystemTime>,
147    pub generation_duration: Option<Duration>,
148    pub status: RequestStatus,
149    pub streamed_bytes: u64,
150    pub stream_chunks: u64,
151    pub input_tokens: Option<u64>,
152    pub output_tokens: Option<u64>,
153    pub error: Option<String>,
154    pub traffic_capture_path: Option<PathBuf>,
155}
156
157impl ActiveRequest {
158    pub fn elapsed(&self) -> Duration {
159        self.started_instant.elapsed()
160    }
161
162    pub fn rate(&self) -> Throughput {
163        throughput(
164            self.output_tokens
165                .and_then(|tokens| tokens.checked_sub(self.generation_initial_output_tokens)),
166            self.streamed_bytes,
167            self.stream_chunks,
168            self.generation_duration.unwrap_or(Duration::ZERO),
169        )
170    }
171}
172
173#[derive(Debug, Clone)]
174pub struct CompletedRequest {
175    pub request_id: String,
176    pub session_id: Option<String>,
177    pub session_seq: Option<u64>,
178    pub project: Option<String>,
179    pub provider: Option<String>,
180    pub model: Option<String>,
181    pub effort: Option<String>,
182    pub endpoint: EndpointKind,
183    pub started_at: SystemTime,
184    pub finished_at: SystemTime,
185    pub generation_started_at: Option<SystemTime>,
186    generation_started_instant: Option<Instant>,
187    generation_initial_output_tokens: u64,
188    pub generation_finished_at: Option<SystemTime>,
189    pub generation_duration: Option<Duration>,
190    pub status: RequestStatus,
191    pub http_status: Option<u16>,
192    pub latency: Duration,
193    pub streamed_bytes: u64,
194    pub stream_chunks: u64,
195    pub input_tokens: Option<u64>,
196    pub output_tokens: Option<u64>,
197    pub error: Option<String>,
198    pub traffic_capture_path: Option<PathBuf>,
199}
200
201impl CompletedRequest {
202    pub fn rate(&self) -> Throughput {
203        throughput(
204            self.output_tokens
205                .and_then(|tokens| tokens.checked_sub(self.generation_initial_output_tokens)),
206            self.streamed_bytes,
207            self.stream_chunks,
208            self.generation_duration.unwrap_or(Duration::ZERO),
209        )
210    }
211}
212
213#[derive(Debug, Clone, PartialEq)]
214pub enum Throughput {
215    TokensPerSecond(f64),
216    BytesPerSecond(f64),
217    EventsPerSecond(f64),
218    None,
219}
220
221impl Throughput {
222    pub fn label(&self) -> String {
223        match self {
224            Self::TokensPerSecond(value) => format!("{value:.1} tok/s"),
225            Self::BytesPerSecond(value) if *value >= 1024.0 => {
226                format!("{:.1} KB/s", value / 1024.0)
227            }
228            Self::BytesPerSecond(value) => format!("{value:.0} B/s"),
229            Self::EventsPerSecond(value) => format!("{value:.1} ev/s"),
230            Self::None => "-".to_string(),
231        }
232    }
233}
234
235#[derive(Debug, Clone)]
236pub struct MonitorState {
237    pub started_at: SystemTime,
238    pub sessions: Vec<SessionSummary>,
239    pub active: Vec<ActiveRequest>,
240    pub recent: Vec<CompletedRequest>,
241}
242
243#[derive(Debug, Clone)]
244pub struct SessionSummary {
245    pub session_id: Option<String>,
246    pub project: Option<String>,
247    pub active_count: usize,
248    pub request_count: usize,
249    pub failure_count: usize,
250    pub provider: Option<String>,
251    pub model: Option<String>,
252    pub effort: Option<String>,
253    pub last_seen: SystemTime,
254    pub input_tokens: u64,
255    pub output_tokens: u64,
256    pub output_token_samples: Vec<(SystemTime, u64)>,
257    rate_output_tokens: u64,
258    pub generation_duration: Duration,
259    pub last_status: String,
260}
261
262impl SessionSummary {
263    pub fn rate(&self) -> Throughput {
264        throughput(
265            Some(self.rate_output_tokens).filter(|tokens| *tokens > 0),
266            0,
267            0,
268            self.generation_duration,
269        )
270    }
271
272    pub fn label(&self) -> String {
273        self.session_id
274            .clone()
275            .unwrap_or_else(|| "no-session".to_string())
276    }
277}
278
279#[derive(Debug)]
280struct MonitorStore {
281    started_at: SystemTime,
282    active: HashMap<String, ActiveRequest>,
283    recent: VecDeque<CompletedRequest>,
284    session_usage: HashMap<Option<String>, SessionUsage>,
285    session_output_buckets: HashMap<Option<String>, Vec<(u64, u64)>>,
286    recent_limit: usize,
287}
288
289#[derive(Debug, Default)]
290struct SessionUsage {
291    input_tokens: u64,
292    output_tokens: u64,
293}
294
295#[derive(Debug, Clone)]
296pub struct MonitorHandle {
297    store: Arc<Mutex<MonitorStore>>,
298}
299
300impl Default for MonitorHandle {
301    fn default() -> Self {
302        Self::new(DEFAULT_RECENT_LIMIT)
303    }
304}
305
306impl MonitorHandle {
307    pub fn new(recent_limit: usize) -> Self {
308        Self {
309            store: Arc::new(Mutex::new(MonitorStore {
310                started_at: SystemTime::now(),
311                active: HashMap::new(),
312                recent: VecDeque::new(),
313                session_usage: HashMap::new(),
314                session_output_buckets: HashMap::new(),
315                recent_limit,
316            })),
317        }
318    }
319
320    pub fn publish(&self, event: MonitorEvent) {
321        if let Ok(mut store) = self.store.lock() {
322            store.apply(event);
323        }
324    }
325
326    pub fn snapshot(&self) -> MonitorState {
327        match self.store.lock() {
328            Ok(store) => store.snapshot(),
329            Err(_) => MonitorState {
330                started_at: SystemTime::now(),
331                sessions: Vec::new(),
332                active: Vec::new(),
333                recent: Vec::new(),
334            },
335        }
336    }
337
338    pub fn request_started(
339        &self,
340        request_id: impl Into<String>,
341        session_id: Option<String>,
342        session_seq: Option<u64>,
343        endpoint: EndpointKind,
344    ) {
345        self.publish(MonitorEvent::RequestStarted {
346            request_id: request_id.into(),
347            session_id,
348            session_seq,
349            endpoint,
350        });
351    }
352
353    pub fn project_resolved(&self, request_id: impl Into<String>, project: impl Into<String>) {
354        self.publish(MonitorEvent::ProjectResolved {
355            request_id: request_id.into(),
356            project: project.into(),
357        });
358    }
359
360    pub fn session_sequence_resolved(&self, request_id: impl Into<String>, session_seq: u64) {
361        self.publish(MonitorEvent::SessionSequenceResolved {
362            request_id: request_id.into(),
363            session_seq,
364        });
365    }
366
367    pub fn provider_selected(
368        &self,
369        request_id: impl Into<String>,
370        provider: impl Into<String>,
371        model: impl Into<String>,
372        effort: Option<String>,
373    ) {
374        self.publish(MonitorEvent::ProviderSelected {
375            request_id: request_id.into(),
376            provider: provider.into(),
377            model: model.into(),
378            effort,
379        });
380    }
381
382    pub fn model_resolved(&self, request_id: impl Into<String>, model: impl Into<String>) {
383        self.publish(MonitorEvent::ModelResolved {
384            request_id: request_id.into(),
385            model: model.into(),
386        });
387    }
388
389    pub fn compaction_started(&self, request_id: impl Into<String>) {
390        self.publish(MonitorEvent::CompactionStarted {
391            request_id: request_id.into(),
392        });
393    }
394
395    pub fn upstream_started(&self, request_id: impl Into<String>) {
396        self.publish(MonitorEvent::UpstreamStarted {
397            request_id: request_id.into(),
398        });
399    }
400
401    pub fn generation_started(&self, request_id: impl Into<String>) {
402        self.publish(MonitorEvent::GenerationStarted {
403            request_id: request_id.into(),
404        });
405    }
406
407    pub fn traffic_capture_path(&self, request_id: impl Into<String>, path: PathBuf) {
408        self.publish(MonitorEvent::TrafficCapturePath {
409            request_id: request_id.into(),
410            path,
411        });
412    }
413
414    pub fn stream_progress(
415        &self,
416        request_id: impl Into<String>,
417        bytes: u64,
418        chunks: u64,
419        input_tokens: Option<u64>,
420        output_tokens: Option<u64>,
421    ) {
422        self.publish(MonitorEvent::StreamProgress {
423            request_id: request_id.into(),
424            bytes,
425            chunks,
426            input_tokens,
427            output_tokens,
428        });
429    }
430
431    pub fn usage_updated(
432        &self,
433        request_id: impl Into<String>,
434        input_tokens: Option<u64>,
435        output_tokens: Option<u64>,
436    ) {
437        self.publish(MonitorEvent::UsageUpdated {
438            request_id: request_id.into(),
439            input_tokens,
440            output_tokens,
441        });
442    }
443
444    pub fn request_completed(
445        &self,
446        request_id: impl Into<String>,
447        http_status: u16,
448        input_tokens: Option<u64>,
449        output_tokens: Option<u64>,
450    ) {
451        self.publish(MonitorEvent::RequestCompleted {
452            request_id: request_id.into(),
453            http_status,
454            input_tokens,
455            output_tokens,
456        });
457    }
458
459    pub fn request_failed(
460        &self,
461        request_id: impl Into<String>,
462        http_status: Option<u16>,
463        error: impl Into<String>,
464    ) {
465        self.publish(MonitorEvent::RequestFailed {
466            request_id: request_id.into(),
467            http_status,
468            error: error.into(),
469        });
470    }
471
472    pub fn request_abandoned(&self, request_id: impl Into<String>, error: impl Into<String>) {
473        self.publish(MonitorEvent::RequestAbandoned {
474            request_id: request_id.into(),
475            error: error.into(),
476        });
477    }
478}
479
480impl MonitorStore {
481    fn apply(&mut self, event: MonitorEvent) {
482        match event {
483            MonitorEvent::RequestStarted {
484                request_id,
485                session_id,
486                session_seq,
487                endpoint,
488            } => {
489                self.active.insert(
490                    request_id.clone(),
491                    ActiveRequest {
492                        request_id,
493                        session_id,
494                        session_seq,
495                        project: None,
496                        provider: None,
497                        model: None,
498                        effort: None,
499                        endpoint,
500                        started_at: SystemTime::now(),
501                        started_instant: Instant::now(),
502                        generation_started_at: None,
503                        generation_started_instant: None,
504                        generation_initial_output_tokens: 0,
505                        generation_finished_at: None,
506                        generation_duration: None,
507                        status: RequestStatus::Started,
508                        streamed_bytes: 0,
509                        stream_chunks: 0,
510                        input_tokens: None,
511                        output_tokens: None,
512                        error: None,
513                        traffic_capture_path: None,
514                    },
515                );
516            }
517            MonitorEvent::ProjectResolved {
518                request_id,
519                project,
520            } => {
521                if let Some(active) = self.active.get_mut(&request_id) {
522                    active.project = Some(project);
523                }
524            }
525            MonitorEvent::SessionSequenceResolved {
526                request_id,
527                session_seq,
528            } => {
529                if let Some(active) = self.active.get_mut(&request_id) {
530                    active.session_seq = Some(session_seq);
531                }
532            }
533            MonitorEvent::ProviderSelected {
534                request_id,
535                provider,
536                model,
537                effort,
538            } => {
539                if let Some(active) = self.active.get_mut(&request_id) {
540                    active.provider = Some(provider);
541                    active.model = Some(model);
542                    active.effort = effort;
543                    active.status = RequestStatus::ProviderSelected;
544                }
545            }
546            MonitorEvent::ModelResolved { request_id, model } => {
547                if let Some(active) = self.active.get_mut(&request_id) {
548                    active.model = Some(match active.model.take() {
549                        Some(incoming) if incoming != model => format!("{incoming} → {model}"),
550                        Some(incoming) => incoming,
551                        None => model,
552                    });
553                }
554            }
555            MonitorEvent::CompactionStarted { request_id } => {
556                if let Some(active) = self.active.get_mut(&request_id) {
557                    active.status = RequestStatus::Compacting;
558                }
559            }
560            MonitorEvent::UpstreamStarted { request_id } => {
561                if let Some(active) = self.active.get_mut(&request_id) {
562                    active.status = RequestStatus::Upstream;
563                }
564            }
565            MonitorEvent::GenerationStarted { request_id } => {
566                if let Some(active) = self.active.get_mut(&request_id) {
567                    active.generation_started_at = Some(SystemTime::now());
568                    active.generation_started_instant = Some(Instant::now());
569                    active.generation_initial_output_tokens = active.output_tokens.unwrap_or(0);
570                    active.generation_finished_at = None;
571                    active.generation_duration = None;
572                }
573            }
574            MonitorEvent::TrafficCapturePath { request_id, path } => {
575                if let Some(active) = self.active.get_mut(&request_id) {
576                    active.traffic_capture_path = Some(path);
577                }
578            }
579            MonitorEvent::StreamProgress {
580                request_id,
581                bytes,
582                chunks,
583                input_tokens,
584                output_tokens,
585            } => {
586                let mut usage_update = None;
587                let mut history_update = None;
588                if let Some(active) = self.active.get_mut(&request_id) {
589                    active.status = RequestStatus::Streaming;
590                    if active.generation_started_instant.is_none() {
591                        active.generation_started_at = Some(SystemTime::now());
592                        active.generation_started_instant = Some(Instant::now());
593                        active.generation_initial_output_tokens =
594                            output_tokens.or(active.output_tokens).unwrap_or(0);
595                    } else {
596                        active.generation_finished_at = Some(SystemTime::now());
597                        active.generation_duration = active
598                            .generation_started_instant
599                            .map(|started| started.elapsed());
600                    }
601                    active.streamed_bytes = active.streamed_bytes.saturating_add(bytes);
602                    active.stream_chunks = active.stream_chunks.saturating_add(chunks);
603                    let input_delta = update_token_count(&mut active.input_tokens, input_tokens);
604                    let output_delta = update_token_count(&mut active.output_tokens, output_tokens);
605                    usage_update = Some((active.session_id.clone(), input_delta, output_delta));
606                } else if let Some(completed) = self
607                    .recent
608                    .iter_mut()
609                    .find(|request| request.request_id == request_id)
610                {
611                    if let Some(started) = completed.generation_started_instant {
612                        completed.generation_finished_at = Some(SystemTime::now());
613                        completed.generation_duration = Some(started.elapsed());
614                    }
615                    completed.streamed_bytes = completed.streamed_bytes.saturating_add(bytes);
616                    completed.stream_chunks = completed.stream_chunks.saturating_add(chunks);
617                    let input_delta = update_token_count(&mut completed.input_tokens, input_tokens);
618                    let output_delta =
619                        update_token_count(&mut completed.output_tokens, output_tokens);
620                    usage_update = Some((completed.session_id.clone(), input_delta, output_delta));
621                    if output_delta > 0 {
622                        history_update = Some((
623                            completed.session_id.clone(),
624                            completed
625                                .generation_finished_at
626                                .unwrap_or(completed.finished_at),
627                            output_delta,
628                        ));
629                    }
630                }
631                if let Some((session_id, input_delta, output_delta)) = usage_update {
632                    self.record_session_usage(session_id, input_delta, output_delta);
633                }
634                if let Some((session_id, timestamp, tokens)) = history_update {
635                    self.record_session_output(session_id, timestamp, tokens);
636                }
637            }
638            MonitorEvent::UsageUpdated {
639                request_id,
640                input_tokens,
641                output_tokens,
642            } => {
643                let mut usage_update = None;
644                let mut history_update = None;
645                if let Some(active) = self.active.get_mut(&request_id) {
646                    if output_tokens.is_some()
647                        && let Some(started) = active.generation_started_instant
648                    {
649                        active.generation_finished_at = Some(SystemTime::now());
650                        active.generation_duration = Some(started.elapsed());
651                    }
652                    let input_delta = update_token_count(&mut active.input_tokens, input_tokens);
653                    let output_delta = update_token_count(&mut active.output_tokens, output_tokens);
654                    usage_update = Some((active.session_id.clone(), input_delta, output_delta));
655                } else if let Some(completed) = self
656                    .recent
657                    .iter_mut()
658                    .find(|request| request.request_id == request_id)
659                {
660                    if output_tokens.is_some()
661                        && let Some(started) = completed.generation_started_instant
662                    {
663                        completed.generation_finished_at = Some(SystemTime::now());
664                        completed.generation_duration = Some(started.elapsed());
665                    }
666                    let input_delta = update_token_count(&mut completed.input_tokens, input_tokens);
667                    let output_delta =
668                        update_token_count(&mut completed.output_tokens, output_tokens);
669                    usage_update = Some((completed.session_id.clone(), input_delta, output_delta));
670                    if output_delta > 0 {
671                        history_update = Some((
672                            completed.session_id.clone(),
673                            completed
674                                .generation_finished_at
675                                .unwrap_or(completed.finished_at),
676                            output_delta,
677                        ));
678                    }
679                }
680                if let Some((session_id, input_delta, output_delta)) = usage_update {
681                    self.record_session_usage(session_id, input_delta, output_delta);
682                }
683                if let Some((session_id, timestamp, tokens)) = history_update {
684                    self.record_session_output(session_id, timestamp, tokens);
685                }
686            }
687            MonitorEvent::RequestCompleted {
688                request_id,
689                http_status,
690                input_tokens,
691                output_tokens,
692            } => {
693                self.finish(
694                    &request_id,
695                    RequestStatus::Completed,
696                    Some(http_status),
697                    input_tokens,
698                    output_tokens,
699                    None,
700                );
701            }
702            MonitorEvent::RequestFailed {
703                request_id,
704                http_status,
705                error,
706            } => {
707                self.finish(
708                    &request_id,
709                    RequestStatus::Failed,
710                    http_status,
711                    None,
712                    None,
713                    Some(error),
714                );
715            }
716            MonitorEvent::RequestAbandoned { request_id, error } => {
717                self.finish_active(
718                    &request_id,
719                    RequestStatus::Failed,
720                    None,
721                    None,
722                    None,
723                    Some(error),
724                );
725            }
726        }
727    }
728
729    fn finish_active(
730        &mut self,
731        request_id: &str,
732        status: RequestStatus,
733        http_status: Option<u16>,
734        input_tokens: Option<u64>,
735        output_tokens: Option<u64>,
736        error: Option<String>,
737    ) {
738        if self.active.contains_key(request_id) {
739            self.finish(
740                request_id,
741                status,
742                http_status,
743                input_tokens,
744                output_tokens,
745                error,
746            );
747        }
748    }
749
750    fn finish(
751        &mut self,
752        request_id: &str,
753        status: RequestStatus,
754        http_status: Option<u16>,
755        input_tokens: Option<u64>,
756        output_tokens: Option<u64>,
757        error: Option<String>,
758    ) {
759        let mut active = self
760            .active
761            .remove(request_id)
762            .unwrap_or_else(|| ActiveRequest {
763                request_id: request_id.to_string(),
764                session_id: None,
765                session_seq: None,
766                project: None,
767                provider: None,
768                model: None,
769                effort: None,
770                endpoint: EndpointKind::Messages,
771                started_at: SystemTime::now(),
772                started_instant: Instant::now(),
773                generation_started_at: None,
774                generation_started_instant: None,
775                generation_initial_output_tokens: 0,
776                generation_finished_at: None,
777                generation_duration: None,
778                status: RequestStatus::Started,
779                streamed_bytes: 0,
780                stream_chunks: 0,
781                input_tokens: None,
782                output_tokens: None,
783                error: None,
784                traffic_capture_path: None,
785            });
786        if output_tokens.is_some()
787            && let Some(started) = active.generation_started_instant
788        {
789            active.generation_finished_at = Some(SystemTime::now());
790            active.generation_duration = Some(started.elapsed());
791        }
792        let input_delta = update_token_count(&mut active.input_tokens, input_tokens);
793        let output_delta = update_token_count(&mut active.output_tokens, output_tokens);
794        self.record_session_usage(active.session_id.clone(), input_delta, output_delta);
795        let completed = CompletedRequest {
796            request_id: active.request_id,
797            session_id: active.session_id,
798            session_seq: active.session_seq,
799            project: active.project,
800            provider: active.provider,
801            model: active.model,
802            effort: active.effort,
803            endpoint: active.endpoint,
804            started_at: active.started_at,
805            finished_at: SystemTime::now(),
806            generation_started_at: active.generation_started_at,
807            generation_started_instant: active.generation_started_instant,
808            generation_initial_output_tokens: active.generation_initial_output_tokens,
809            generation_finished_at: active.generation_finished_at,
810            generation_duration: active.generation_duration,
811            status,
812            http_status,
813            latency: active.started_instant.elapsed(),
814            streamed_bytes: active.streamed_bytes,
815            stream_chunks: active.stream_chunks,
816            input_tokens: active.input_tokens,
817            output_tokens: active.output_tokens,
818            error: error.or(active.error),
819            traffic_capture_path: active.traffic_capture_path,
820        };
821        if let Some(tokens) = completed.output_tokens.filter(|tokens| *tokens > 0) {
822            self.record_session_output(
823                completed.session_id.clone(),
824                completed
825                    .generation_finished_at
826                    .unwrap_or(completed.finished_at),
827                tokens,
828            );
829        }
830        self.recent.push_front(completed);
831        while self.recent.len() > self.recent_limit {
832            self.recent.pop_back();
833        }
834    }
835
836    fn record_session_usage(
837        &mut self,
838        session_id: Option<String>,
839        input_tokens: u64,
840        output_tokens: u64,
841    ) {
842        let usage = self.session_usage.entry(session_id).or_default();
843        usage.input_tokens = usage.input_tokens.saturating_add(input_tokens);
844        usage.output_tokens = usage.output_tokens.saturating_add(output_tokens);
845    }
846
847    fn record_session_output(
848        &mut self,
849        session_id: Option<String>,
850        timestamp: SystemTime,
851        tokens: u64,
852    ) {
853        let bucket = session_token_bucket(timestamp);
854        let buckets = self.session_output_buckets.entry(session_id).or_default();
855        match buckets.binary_search_by_key(&bucket, |(bucket, _)| *bucket) {
856            Ok(index) => buckets[index].1 = buckets[index].1.saturating_add(tokens),
857            Err(index) => buckets.insert(index, (bucket, tokens)),
858        }
859    }
860
861    fn snapshot(&self) -> MonitorState {
862        let mut active: Vec<_> = self.active.values().cloned().collect();
863        active.sort_by_key(|request| request.started_at);
864        let sessions = session_summaries(
865            &active,
866            &self.recent,
867            &self.session_usage,
868            &self.session_output_buckets,
869        );
870        MonitorState {
871            started_at: self.started_at,
872            sessions,
873            active,
874            recent: self.recent.iter().cloned().collect(),
875        }
876    }
877}
878
879fn session_summaries(
880    active: &[ActiveRequest],
881    recent: &VecDeque<CompletedRequest>,
882    session_usage: &HashMap<Option<String>, SessionUsage>,
883    session_output_buckets: &HashMap<Option<String>, Vec<(u64, u64)>>,
884) -> Vec<SessionSummary> {
885    let mut sessions: HashMap<Option<String>, SessionSummary> = HashMap::new();
886    for request in recent.iter().rev() {
887        let entry = sessions
888            .entry(request.session_id.clone())
889            .or_insert_with(|| SessionSummary {
890                session_id: request.session_id.clone(),
891                project: request.project.clone(),
892                active_count: 0,
893                request_count: 0,
894                failure_count: 0,
895                provider: None,
896                model: None,
897                effort: None,
898                last_seen: request.finished_at,
899                input_tokens: 0,
900                output_tokens: 0,
901                output_token_samples: Vec::new(),
902                rate_output_tokens: 0,
903                generation_duration: Duration::ZERO,
904                last_status: "-".to_string(),
905            });
906        entry.request_count += 1;
907        if request.status == RequestStatus::Failed {
908            entry.failure_count += 1;
909        }
910        entry.project = request.project.clone().or(entry.project.clone());
911        entry.provider = request.provider.clone().or(entry.provider.clone());
912        entry.model = request.model.clone().or(entry.model.clone());
913        entry.effort = request.effort.clone().or(entry.effort.clone());
914        entry.last_seen = max_system_time(entry.last_seen, request.finished_at);
915        if let (Some(tokens), Some(duration)) = (
916            request
917                .output_tokens
918                .and_then(|tokens| tokens.checked_sub(request.generation_initial_output_tokens))
919                .filter(|tokens| *tokens > 0),
920            request
921                .generation_duration
922                .filter(|duration| !duration.is_zero()),
923        ) {
924            entry.rate_output_tokens = entry.rate_output_tokens.saturating_add(tokens);
925            entry.generation_duration = entry.generation_duration.saturating_add(duration);
926        }
927        entry.last_status = request.status.label().to_string();
928    }
929
930    for request in active {
931        let entry = sessions
932            .entry(request.session_id.clone())
933            .or_insert_with(|| SessionSummary {
934                session_id: request.session_id.clone(),
935                project: request.project.clone(),
936                active_count: 0,
937                request_count: 0,
938                failure_count: 0,
939                provider: None,
940                model: None,
941                effort: None,
942                last_seen: request.started_at,
943                input_tokens: 0,
944                output_tokens: 0,
945                output_token_samples: Vec::new(),
946                rate_output_tokens: 0,
947                generation_duration: Duration::ZERO,
948                last_status: "-".to_string(),
949            });
950        entry.active_count += 1;
951        entry.request_count += 1;
952        entry.project = request.project.clone().or(entry.project.clone());
953        entry.provider = request.provider.clone().or(entry.provider.clone());
954        entry.model = request.model.clone().or(entry.model.clone());
955        entry.effort = request.effort.clone().or(entry.effort.clone());
956        entry.last_seen = max_system_time(entry.last_seen, request.started_at);
957        if let (Some(tokens), Some(duration)) = (
958            request
959                .output_tokens
960                .and_then(|tokens| tokens.checked_sub(request.generation_initial_output_tokens))
961                .filter(|tokens| *tokens > 0),
962            request
963                .generation_duration
964                .filter(|duration| !duration.is_zero()),
965        ) {
966            entry.rate_output_tokens = entry.rate_output_tokens.saturating_add(tokens);
967            entry.generation_duration = entry.generation_duration.saturating_add(duration);
968        }
969        entry.last_status = request.status.label().to_string();
970    }
971
972    for (session_id, session) in &mut sessions {
973        if let Some(usage) = session_usage.get(session_id) {
974            session.input_tokens = usage.input_tokens;
975            session.output_tokens = usage.output_tokens;
976        }
977        if let Some(buckets) = session_output_buckets.get(session_id) {
978            session.output_token_samples = buckets
979                .iter()
980                .map(|(bucket, tokens)| (session_token_bucket_start(*bucket), *tokens))
981                .collect();
982        }
983    }
984
985    let mut out: Vec<_> = sessions.into_values().collect();
986    out.sort_by_key(SessionSummary::label);
987    out
988}
989
990fn session_token_bucket(timestamp: SystemTime) -> u64 {
991    timestamp
992        .duration_since(SystemTime::UNIX_EPOCH)
993        .unwrap_or(Duration::ZERO)
994        .as_secs()
995        / SESSION_TOKEN_BUCKET_SECS
996}
997
998fn session_token_bucket_start(bucket: u64) -> SystemTime {
999    SystemTime::UNIX_EPOCH + Duration::from_secs(bucket.saturating_mul(SESSION_TOKEN_BUCKET_SECS))
1000}
1001
1002fn update_token_count(current: &mut Option<u64>, incoming: Option<u64>) -> u64 {
1003    let Some(incoming) = incoming else {
1004        return 0;
1005    };
1006    let previous = current.unwrap_or(0);
1007    if incoming > previous || current.is_none() {
1008        *current = Some(incoming);
1009    }
1010    incoming.saturating_sub(previous)
1011}
1012
1013fn max_system_time(left: SystemTime, right: SystemTime) -> SystemTime {
1014    if right.duration_since(left).is_ok() {
1015        right
1016    } else {
1017        left
1018    }
1019}
1020
1021pub fn throughput(
1022    output_tokens: Option<u64>,
1023    streamed_bytes: u64,
1024    stream_chunks: u64,
1025    elapsed: Duration,
1026) -> Throughput {
1027    let secs = elapsed.as_secs_f64();
1028    if secs <= 0.0 {
1029        return Throughput::None;
1030    }
1031    if let Some(tokens) = output_tokens.filter(|tokens| *tokens > 0) {
1032        return Throughput::TokensPerSecond(tokens as f64 / secs);
1033    }
1034    if streamed_bytes > 0 {
1035        return Throughput::BytesPerSecond(streamed_bytes as f64 / secs);
1036    }
1037    if stream_chunks > 0 {
1038        return Throughput::EventsPerSecond(stream_chunks as f64 / secs);
1039    }
1040    Throughput::None
1041}
1042
1043pub fn usage_from_anthropic_sse(bytes: &[u8]) -> (Option<u64>, Option<u64>) {
1044    let text = String::from_utf8_lossy(bytes);
1045    let mut input_tokens = None;
1046    let mut output_tokens = None;
1047    for line in text.lines() {
1048        let Some(data) = line.strip_prefix("data:") else {
1049            continue;
1050        };
1051        let Ok(value) = serde_json::from_str::<serde_json::Value>(data.trim()) else {
1052            continue;
1053        };
1054        for usage in [
1055            value.pointer("/usage"),
1056            value.pointer("/delta/usage"),
1057            value.pointer("/message/usage"),
1058        ]
1059        .into_iter()
1060        .flatten()
1061        {
1062            if let Some(tokens) = usage.get("input_tokens").and_then(|value| value.as_u64()) {
1063                input_tokens = Some(tokens);
1064            }
1065            if let Some(tokens) = usage.get("output_tokens").and_then(|value| value.as_u64()) {
1066                output_tokens = Some(tokens);
1067            }
1068        }
1069    }
1070    (input_tokens, output_tokens)
1071}
1072
1073#[cfg(test)]
1074mod tests {
1075    use super::*;
1076
1077    #[test]
1078    fn started_requests_appear_active() {
1079        let monitor = MonitorHandle::new(10);
1080        monitor.request_started(
1081            "r1",
1082            Some("s1".to_string()),
1083            Some(3),
1084            EndpointKind::Messages,
1085        );
1086        let state = monitor.snapshot();
1087        assert_eq!(state.active.len(), 1);
1088        assert_eq!(state.active[0].request_id, "r1");
1089        assert_eq!(state.active[0].session_id.as_deref(), Some("s1"));
1090        assert_eq!(state.active[0].session_seq, Some(3));
1091    }
1092
1093    #[test]
1094    fn resolved_model_appends_to_incoming_alias() {
1095        let monitor = MonitorHandle::new(10);
1096        monitor.request_started("r1", None, None, EndpointKind::Messages);
1097        monitor.provider_selected("r1", "codex", "claude-sonnet-4-6", None);
1098        monitor.model_resolved("r1", "gpt-5.4");
1099
1100        let state = monitor.snapshot();
1101        assert_eq!(
1102            state.active[0].model.as_deref(),
1103            Some("claude-sonnet-4-6 → gpt-5.4")
1104        );
1105    }
1106
1107    #[test]
1108    fn identical_resolved_model_is_shown_once() {
1109        let monitor = MonitorHandle::new(10);
1110        monitor.request_started("r1", None, None, EndpointKind::Messages);
1111        monitor.provider_selected("r1", "codex", "gpt-5.6-sol", None);
1112        monitor.model_resolved("r1", "gpt-5.6-sol");
1113
1114        let state = monitor.snapshot();
1115        assert_eq!(state.active[0].model.as_deref(), Some("gpt-5.6-sol"));
1116    }
1117
1118    #[test]
1119    fn compaction_started_marks_request_compacting() {
1120        let monitor = MonitorHandle::new(10);
1121        monitor.request_started("r1", None, None, EndpointKind::Messages);
1122        monitor.provider_selected("r1", "codex", "gpt-5.6-sol", None);
1123        monitor.compaction_started("r1");
1124
1125        let state = monitor.snapshot();
1126        assert_eq!(state.active[0].status, RequestStatus::Compacting);
1127        assert_eq!(state.sessions[0].last_status, "compacting");
1128    }
1129
1130    #[test]
1131    fn generation_baseline_pairs_total_usage_with_the_full_observed_interval() {
1132        let monitor = MonitorHandle::new(10);
1133        monitor.request_started("r1", None, None, EndpointKind::Messages);
1134        monitor.generation_started("r1");
1135        monitor.stream_progress("r1", 50, 1, Some(1_225), Some(141));
1136        monitor.request_completed("r1", 200, None, None);
1137
1138        let request = &monitor.snapshot().recent[0];
1139
1140        assert!(
1141            request
1142                .generation_duration
1143                .is_some_and(|duration| !duration.is_zero())
1144        );
1145        assert!(matches!(request.rate(), Throughput::TokensPerSecond(_)));
1146    }
1147
1148    #[test]
1149    fn first_stream_progress_has_no_rate_without_an_interval() {
1150        let monitor = MonitorHandle::new(10);
1151        monitor.request_started("r1", None, None, EndpointKind::Messages);
1152        monitor.stream_progress("r1", 50, 1, Some(1_225), Some(141));
1153
1154        let state = monitor.snapshot();
1155        assert_eq!(state.active.len(), 1);
1156        assert_eq!(state.active[0].rate(), Throughput::None);
1157    }
1158
1159    #[test]
1160    fn late_stream_progress_extends_stream_timing_without_extending_request_latency() {
1161        let monitor = MonitorHandle::new(10);
1162        monitor.request_started("r1", None, None, EndpointKind::Messages);
1163        monitor.stream_progress("r1", 100, 1, Some(0), Some(0));
1164        monitor.request_completed("r1", 200, None, None);
1165        let completed = monitor.snapshot().recent[0].clone();
1166        monitor.stream_progress("r1", 50, 1, Some(1_225), Some(141));
1167
1168        let state = monitor.snapshot();
1169        assert!(state.active.is_empty());
1170        assert_eq!(state.recent[0].streamed_bytes, 150);
1171        assert_eq!(state.recent[0].stream_chunks, 2);
1172        assert_eq!(state.recent[0].input_tokens, Some(1_225));
1173        assert_eq!(state.recent[0].output_tokens, Some(141));
1174        assert_eq!(state.recent[0].finished_at, completed.finished_at);
1175        assert_eq!(state.recent[0].latency, completed.latency);
1176        assert!(state.recent[0].generation_duration > completed.generation_duration);
1177        assert!(matches!(
1178            state.recent[0].rate(),
1179            Throughput::TokensPerSecond(_)
1180        ));
1181        assert_eq!(
1182            state.sessions[0]
1183                .output_token_samples
1184                .iter()
1185                .map(|(_, tokens)| *tokens)
1186                .sum::<u64>(),
1187            141
1188        );
1189    }
1190
1191    #[test]
1192    fn completed_requests_leave_active_and_enter_recent() {
1193        let monitor = MonitorHandle::new(10);
1194        monitor.request_started("r1", None, None, EndpointKind::Messages);
1195        monitor.provider_selected("r1", "codex", "gpt-5.5", Some("high".to_string()));
1196        monitor.request_completed("r1", 200, Some(10), Some(20));
1197        let state = monitor.snapshot();
1198        assert!(state.active.is_empty());
1199        assert_eq!(state.recent.len(), 1);
1200        assert_eq!(state.recent[0].provider.as_deref(), Some("codex"));
1201        assert_eq!(state.recent[0].effort.as_deref(), Some("high"));
1202        assert_eq!(state.recent[0].output_tokens, Some(20));
1203    }
1204
1205    #[test]
1206    fn failed_requests_preserve_error_summary() {
1207        let monitor = MonitorHandle::new(10);
1208        monitor.request_started("r1", None, None, EndpointKind::Messages);
1209        monitor.request_failed("r1", Some(400), "Unknown model");
1210        let state = monitor.snapshot();
1211        assert_eq!(state.recent[0].status, RequestStatus::Failed);
1212        assert_eq!(state.recent[0].http_status, Some(400));
1213        assert_eq!(state.recent[0].error.as_deref(), Some("Unknown model"));
1214    }
1215
1216    #[test]
1217    fn abandoned_requests_leave_active_once() {
1218        let monitor = MonitorHandle::new(10);
1219        monitor.request_started("r1", None, None, EndpointKind::Messages);
1220        monitor.request_abandoned("r1", "request dropped");
1221        monitor.request_abandoned("r1", "request dropped again");
1222        let state = monitor.snapshot();
1223        assert!(state.active.is_empty());
1224        assert_eq!(state.recent.len(), 1);
1225        assert_eq!(state.recent[0].status, RequestStatus::Failed);
1226        assert_eq!(state.recent[0].http_status, None);
1227        assert_eq!(state.recent[0].error.as_deref(), Some("request dropped"));
1228    }
1229
1230    #[test]
1231    fn completed_requests_ignore_late_abandonment() {
1232        let monitor = MonitorHandle::new(10);
1233        monitor.request_started("r1", None, None, EndpointKind::Messages);
1234        monitor.request_completed("r1", 200, None, None);
1235        monitor.request_abandoned("r1", "request dropped");
1236        let state = monitor.snapshot();
1237        assert!(state.active.is_empty());
1238        assert_eq!(state.recent.len(), 1);
1239        assert_eq!(state.recent[0].status, RequestStatus::Completed);
1240    }
1241
1242    #[test]
1243    fn bounded_recent_history_drops_oldest() {
1244        let monitor = MonitorHandle::new(2);
1245        for id in ["r1", "r2", "r3"] {
1246            monitor.request_started(id, None, None, EndpointKind::Messages);
1247            monitor.request_completed(id, 200, None, None);
1248        }
1249        let state = monitor.snapshot();
1250        let ids: Vec<_> = state
1251            .recent
1252            .iter()
1253            .map(|request| request.request_id.as_str())
1254            .collect();
1255        assert_eq!(ids, vec!["r3", "r2"]);
1256    }
1257
1258    #[test]
1259    fn throughput_selects_best_available_signal() {
1260        let elapsed = Duration::from_secs(2);
1261        assert_eq!(
1262            throughput(Some(84), 1024, 10, elapsed),
1263            Throughput::TokensPerSecond(42.0)
1264        );
1265        assert_eq!(
1266            throughput(None, 2048, 10, elapsed),
1267            Throughput::BytesPerSecond(1024.0)
1268        );
1269        assert_eq!(
1270            throughput(None, 0, 36, elapsed),
1271            Throughput::EventsPerSecond(18.0)
1272        );
1273    }
1274
1275    #[test]
1276    fn sse_usage_extracts_final_message_delta_tokens() {
1277        let sse = br#"event: message_start
1278data: {"type":"message_start","message":{"usage":{"input_tokens":0,"output_tokens":0}}}
1279
1280event: message_delta
1281data: {"type":"message_delta","delta":{"stop_reason":"end_turn"},"usage":{"input_tokens":12,"output_tokens":48}}
1282
1283"#;
1284        assert_eq!(usage_from_anthropic_sse(sse), (Some(12), Some(48)));
1285    }
1286
1287    fn completed_request(
1288        request_id: &str,
1289        session_id: &str,
1290        output_tokens: u64,
1291        latency: Duration,
1292        generation_duration: Option<Duration>,
1293    ) -> CompletedRequest {
1294        CompletedRequest {
1295            request_id: request_id.to_string(),
1296            session_id: Some(session_id.to_string()),
1297            session_seq: None,
1298            project: None,
1299            provider: Some("codex".to_string()),
1300            model: Some("gpt-5.6-sol".to_string()),
1301            effort: None,
1302            endpoint: EndpointKind::Messages,
1303            started_at: SystemTime::UNIX_EPOCH,
1304            finished_at: SystemTime::UNIX_EPOCH + latency,
1305            generation_started_at: generation_duration.map(|_| SystemTime::UNIX_EPOCH),
1306            generation_started_instant: None,
1307            generation_initial_output_tokens: 0,
1308            generation_finished_at: generation_duration
1309                .map(|duration| SystemTime::UNIX_EPOCH + duration),
1310            generation_duration,
1311            status: RequestStatus::Completed,
1312            http_status: Some(200),
1313            latency,
1314            streamed_bytes: 0,
1315            stream_chunks: 0,
1316            input_tokens: None,
1317            output_tokens: Some(output_tokens),
1318            error: None,
1319            traffic_capture_path: None,
1320        }
1321    }
1322
1323    fn session_summaries_for_requests(recent: &VecDeque<CompletedRequest>) -> Vec<SessionSummary> {
1324        let mut usage = HashMap::<Option<String>, SessionUsage>::new();
1325        for request in recent {
1326            let entry = usage.entry(request.session_id.clone()).or_default();
1327            entry.input_tokens = entry
1328                .input_tokens
1329                .saturating_add(request.input_tokens.unwrap_or(0));
1330            entry.output_tokens = entry
1331                .output_tokens
1332                .saturating_add(request.output_tokens.unwrap_or(0));
1333        }
1334        session_summaries(&[], recent, &usage, &HashMap::new())
1335    }
1336
1337    #[test]
1338    fn completed_request_rate_uses_stream_interval_instead_of_request_latency() {
1339        let request = completed_request(
1340            "r1",
1341            "s1",
1342            120,
1343            Duration::from_secs(30),
1344            Some(Duration::from_secs(4)),
1345        );
1346
1347        assert_eq!(request.rate(), Throughput::TokensPerSecond(30.0));
1348    }
1349
1350    #[test]
1351    fn request_rate_uses_token_delta_from_the_initial_observation() {
1352        let mut request = completed_request(
1353            "r1",
1354            "s1",
1355            120,
1356            Duration::from_secs(30),
1357            Some(Duration::from_secs(4)),
1358        );
1359        request.generation_initial_output_tokens = 20;
1360
1361        assert_eq!(request.rate(), Throughput::TokensPerSecond(25.0));
1362    }
1363
1364    #[test]
1365    fn session_rate_combines_request_tokens_and_generation_intervals() {
1366        let recent = VecDeque::from([
1367            completed_request(
1368                "r2",
1369                "s1",
1370                50,
1371                Duration::from_secs(40),
1372                Some(Duration::from_secs(1)),
1373            ),
1374            completed_request(
1375                "r1",
1376                "s1",
1377                100,
1378                Duration::from_secs(20),
1379                Some(Duration::from_secs(4)),
1380            ),
1381        ]);
1382
1383        let sessions = session_summaries_for_requests(&recent);
1384
1385        assert_eq!(sessions[0].output_tokens, 150);
1386        assert_eq!(sessions[0].generation_duration, Duration::from_secs(5));
1387        assert_eq!(sessions[0].rate(), Throughput::TokensPerSecond(30.0));
1388    }
1389
1390    #[test]
1391    fn output_without_observed_stream_interval_has_no_output_rate() {
1392        let request = completed_request("r1", "s1", 120, Duration::from_secs(30), None);
1393        let recent = VecDeque::from([request.clone()]);
1394
1395        assert_eq!(request.rate(), Throughput::None);
1396        assert_eq!(
1397            session_summaries_for_requests(&recent)[0].rate(),
1398            Throughput::None
1399        );
1400    }
1401
1402    #[test]
1403    fn session_rate_excludes_interval_without_output_usage() {
1404        let mut tokenless = completed_request(
1405            "tokenless",
1406            "s1",
1407            0,
1408            Duration::from_secs(30),
1409            Some(Duration::from_secs(100)),
1410        );
1411        tokenless.output_tokens = None;
1412        let recent = VecDeque::from([
1413            tokenless,
1414            completed_request(
1415                "measured",
1416                "s1",
1417                100,
1418                Duration::from_secs(20),
1419                Some(Duration::from_secs(4)),
1420            ),
1421        ]);
1422
1423        assert_eq!(
1424            session_summaries_for_requests(&recent)[0].rate(),
1425            Throughput::TokensPerSecond(25.0)
1426        );
1427    }
1428
1429    #[test]
1430    fn session_rate_excludes_output_without_a_matching_stream_interval() {
1431        let recent = VecDeque::from([
1432            completed_request("buffered", "s1", 900, Duration::from_secs(30), None),
1433            completed_request(
1434                "streamed",
1435                "s1",
1436                100,
1437                Duration::from_secs(20),
1438                Some(Duration::from_secs(4)),
1439            ),
1440        ]);
1441
1442        let session = &session_summaries_for_requests(&recent)[0];
1443
1444        assert_eq!(session.output_tokens, 1_000);
1445        assert_eq!(session.rate(), Throughput::TokensPerSecond(25.0));
1446    }
1447
1448    #[test]
1449    fn session_summaries_group_recent_and_active_requests() {
1450        let monitor = MonitorHandle::new(10);
1451        monitor.request_started(
1452            "r1",
1453            Some("s1".to_string()),
1454            Some(1),
1455            EndpointKind::Messages,
1456        );
1457        monitor.project_resolved("r1", "example");
1458        monitor.provider_selected("r1", "codex", "gpt-5.5", None);
1459        monitor.request_completed("r1", 200, Some(10), Some(20));
1460        monitor.request_started(
1461            "r2",
1462            Some("s1".to_string()),
1463            Some(2),
1464            EndpointKind::Messages,
1465        );
1466        monitor.provider_selected("r2", "codex", "gpt-5.5", Some("xhigh".to_string()));
1467        let state = monitor.snapshot();
1468        assert_eq!(state.sessions.len(), 1);
1469        assert_eq!(state.sessions[0].label(), "s1");
1470        assert_eq!(state.sessions[0].project.as_deref(), Some("example"));
1471        assert_eq!(state.sessions[0].request_count, 2);
1472        assert_eq!(state.sessions[0].active_count, 1);
1473        assert_eq!(state.sessions[0].effort.as_deref(), Some("xhigh"));
1474        assert_eq!(state.sessions[0].output_tokens, 20);
1475        assert_eq!(
1476            state.sessions[0]
1477                .output_token_samples
1478                .iter()
1479                .map(|(_, tokens)| *tokens)
1480                .collect::<Vec<_>>(),
1481            vec![20]
1482        );
1483    }
1484
1485    #[test]
1486    fn session_output_history_survives_request_eviction() {
1487        let monitor = MonitorHandle::new(1);
1488        for (request_id, tokens) in [("oldest", 20), ("newest", 80)] {
1489            monitor.request_started(
1490                request_id,
1491                Some("s1".to_string()),
1492                None,
1493                EndpointKind::Messages,
1494            );
1495            monitor.request_completed(request_id, 200, Some(tokens * 10), Some(tokens));
1496        }
1497
1498        let state = monitor.snapshot();
1499
1500        assert_eq!(state.recent.len(), 1);
1501        assert_eq!(state.sessions[0].input_tokens, 1_000);
1502        assert_eq!(state.sessions[0].output_tokens, 100);
1503        assert_eq!(
1504            state.sessions[0]
1505                .output_token_samples
1506                .iter()
1507                .map(|(_, tokens)| *tokens)
1508                .sum::<u64>(),
1509            100
1510        );
1511    }
1512
1513    #[test]
1514    fn session_usage_ignores_decreasing_request_observations() {
1515        let monitor = MonitorHandle::new(10);
1516        monitor.request_started(
1517            "r1",
1518            Some("s1".to_string()),
1519            Some(1),
1520            EndpointKind::Messages,
1521        );
1522        monitor.usage_updated("r1", Some(100), Some(20));
1523        monitor.usage_updated("r1", Some(90), Some(15));
1524        monitor.usage_updated("r1", Some(120), Some(25));
1525
1526        let active = monitor.snapshot();
1527        assert_eq!(active.active[0].input_tokens, Some(120));
1528        assert_eq!(active.active[0].output_tokens, Some(25));
1529        assert_eq!(active.sessions[0].input_tokens, 120);
1530        assert_eq!(active.sessions[0].output_tokens, 25);
1531
1532        monitor.request_completed("r1", 200, Some(80), Some(10));
1533        let completed = monitor.snapshot();
1534        assert_eq!(completed.recent[0].input_tokens, Some(120));
1535        assert_eq!(completed.recent[0].output_tokens, Some(25));
1536        assert_eq!(completed.sessions[0].input_tokens, 120);
1537        assert_eq!(completed.sessions[0].output_tokens, 25);
1538    }
1539
1540    #[test]
1541    fn compaction_preserves_cumulative_session_usage() {
1542        let monitor = MonitorHandle::new(10);
1543        monitor.request_started(
1544            "before",
1545            Some("s1".to_string()),
1546            Some(1),
1547            EndpointKind::Messages,
1548        );
1549        monitor.request_completed("before", 200, Some(100), Some(20));
1550        monitor.request_started(
1551            "compact",
1552            Some("s1".to_string()),
1553            Some(2),
1554            EndpointKind::Messages,
1555        );
1556        monitor.compaction_started("compact");
1557        monitor.request_completed("compact", 200, Some(40), Some(10));
1558
1559        let state = monitor.snapshot();
1560        assert_eq!(state.sessions[0].input_tokens, 140);
1561        assert_eq!(state.sessions[0].output_tokens, 30);
1562    }
1563
1564    #[test]
1565    fn session_sequence_restart_preserves_cumulative_usage() {
1566        let monitor = MonitorHandle::new(10);
1567        monitor.request_started(
1568            "before",
1569            Some("s1".to_string()),
1570            None,
1571            EndpointKind::Messages,
1572        );
1573        monitor.session_sequence_resolved("before", 7);
1574        monitor.request_completed("before", 200, Some(100), Some(20));
1575
1576        monitor.request_started(
1577            "after",
1578            Some("s1".to_string()),
1579            None,
1580            EndpointKind::Messages,
1581        );
1582        monitor.session_sequence_resolved("after", 1);
1583        monitor.request_completed("after", 200, Some(25), Some(5));
1584
1585        let state = monitor.snapshot();
1586        assert_eq!(state.sessions[0].input_tokens, 125);
1587        assert_eq!(state.sessions[0].output_tokens, 25);
1588    }
1589
1590    #[test]
1591    fn session_order_is_stable_across_activity() {
1592        let monitor = MonitorHandle::new(10);
1593        monitor.request_started(
1594            "r1",
1595            Some("session-b".to_string()),
1596            Some(1),
1597            EndpointKind::Messages,
1598        );
1599        monitor.request_started(
1600            "r2",
1601            Some("session-a".to_string()),
1602            Some(1),
1603            EndpointKind::Messages,
1604        );
1605
1606        let first: Vec<_> = monitor
1607            .snapshot()
1608            .sessions
1609            .iter()
1610            .map(SessionSummary::label)
1611            .collect();
1612        monitor.request_completed("r1", 200, None, None);
1613        monitor.request_started(
1614            "r3",
1615            Some("session-b".to_string()),
1616            Some(2),
1617            EndpointKind::Messages,
1618        );
1619        let second: Vec<_> = monitor
1620            .snapshot()
1621            .sessions
1622            .iter()
1623            .map(SessionSummary::label)
1624            .collect();
1625
1626        assert_eq!(first, vec!["session-a", "session-b"]);
1627        assert_eq!(second, first);
1628    }
1629}