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}