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