1use std::collections::{HashMap, HashSet};
14use std::time::{SystemTime, UNIX_EPOCH};
15
16use serde::{Deserialize, Serialize};
17
18#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
20#[serde(rename_all = "snake_case")]
21pub enum InferenceTask {
22 Generate,
23 Embed,
24 Classify,
25 Code,
26 Reasoning,
27}
28
29impl std::fmt::Display for InferenceTask {
30 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
31 match self {
32 InferenceTask::Generate => write!(f, "generate"),
33 InferenceTask::Embed => write!(f, "embed"),
34 InferenceTask::Classify => write!(f, "classify"),
35 InferenceTask::Code => write!(f, "code"),
36 InferenceTask::Reasoning => write!(f, "reasoning"),
37 }
38 }
39}
40
41#[derive(Debug, Clone, Serialize, Deserialize)]
43pub struct InferenceOutcome {
44 pub trace_id: String,
46 pub model_id: String,
48 pub task: InferenceTask,
50 pub routing_reason: String,
52 pub latency_ms: u64,
54 pub input_tokens: usize,
57 pub output_tokens: usize,
59 #[serde(default)]
62 pub cache_read_input_tokens: usize,
63 #[serde(default)]
66 pub cache_creation_input_tokens: usize,
67 pub inferred_outcome: Option<InferredOutcome>,
69 pub code_outcome: Option<CodeOutcome>,
71 pub error: Option<String>,
73 pub timestamp: u64,
75 #[serde(skip)]
82 pub success_credited: bool,
83}
84
85#[derive(Debug, Clone, Serialize, Deserialize)]
87#[serde(tag = "type", rename_all = "snake_case")]
88pub enum InferredOutcome {
89 Accepted { confidence: f64 },
91 AcceptedWithEdits { confidence: f64 },
93 Rejected { confidence: f64 },
95 Inconclusive,
97}
98
99impl InferredOutcome {
100 pub fn quality_score(&self) -> Option<f64> {
102 match self {
103 InferredOutcome::Accepted { confidence } => Some(*confidence),
104 InferredOutcome::AcceptedWithEdits { confidence } => Some(confidence * 0.7),
105 InferredOutcome::Rejected { confidence } => Some((1.0 - confidence) * 0.3),
106 InferredOutcome::Inconclusive => None,
107 }
108 }
109
110 pub fn is_success(&self) -> Option<bool> {
111 match self {
112 InferredOutcome::Accepted { .. } => Some(true),
113 InferredOutcome::AcceptedWithEdits { .. } => Some(true),
114 InferredOutcome::Rejected { .. } => Some(false),
115 InferredOutcome::Inconclusive => None,
116 }
117 }
118}
119
120#[derive(Debug, Clone, Serialize, Deserialize)]
122#[serde(tag = "type", rename_all = "snake_case")]
123pub enum CodeOutcome {
124 Applied,
126 Modified,
128 Ignored,
130 SignatureChanged,
132 BodyModified,
134 SymbolAdded,
136}
137
138impl CodeOutcome {
139 pub fn quality_score(&self) -> f64 {
140 match self {
141 CodeOutcome::Applied => 1.0,
142 CodeOutcome::SignatureChanged => 0.8,
143 CodeOutcome::BodyModified => 0.7,
144 CodeOutcome::SymbolAdded => 0.7,
145 CodeOutcome::Modified => 0.6,
146 CodeOutcome::Ignored => 0.1,
147 }
148 }
149
150 pub fn is_success(&self) -> bool {
151 !matches!(self, CodeOutcome::Ignored)
152 }
153}
154
155fn backfill_quality_observations(p: &mut ModelProfile) {
161 if p.quality_observations == 0 && p.fail_count > 0 {
162 p.quality_observations = p.fail_count;
163 }
164 for ts in p.task_stats.values_mut() {
165 if ts.quality_observations == 0 && ts.failures > 0 {
166 ts.quality_observations = ts.failures;
167 }
168 }
169}
170
171#[derive(Debug, Clone, Default, Serialize, Deserialize)]
173pub struct TaskStats {
174 pub calls: u64,
175 pub successes: u64,
176 pub failures: u64,
177 pub avg_latency_ms: f64,
179 pub ema_quality: f64,
181 #[serde(default)]
186 pub prior_sample_size: usize,
187 #[serde(default)]
194 pub quality_observations: u64,
195}
196
197impl TaskStats {
198 pub fn success_rate(&self) -> f64 {
199 let total = self.successes + self.failures;
200 if total == 0 {
201 return 0.5;
202 } self.successes as f64 / total as f64
204 }
205}
206
207#[derive(Debug, Clone, Serialize, Deserialize)]
209pub struct ModelProfile {
210 pub model_id: String,
211 pub total_calls: u64,
212 pub success_count: u64,
213 pub fail_count: u64,
214 pub total_latency_ms: u64,
215 #[serde(default)]
219 pub total_input_tokens: u64,
220 #[serde(default)]
222 pub total_output_tokens: u64,
223 #[serde(default)]
225 pub total_cache_read_input_tokens: u64,
226 #[serde(default)]
228 pub total_cache_creation_input_tokens: u64,
229 pub task_stats: HashMap<String, TaskStats>,
231 pub ema_quality: f64,
233 #[serde(default)]
236 pub prior_sample_size: usize,
237 #[serde(default)]
240 pub quality_observations: u64,
241 #[serde(default)]
246 pub quality_per_1k_tokens: f64,
247 pub updated_at: u64,
249}
250
251impl ModelProfile {
252 pub fn new(model_id: String) -> Self {
253 Self {
254 model_id,
255 total_calls: 0,
256 success_count: 0,
257 fail_count: 0,
258 total_latency_ms: 0,
259 total_input_tokens: 0,
260 total_output_tokens: 0,
261 total_cache_read_input_tokens: 0,
262 total_cache_creation_input_tokens: 0,
263 task_stats: HashMap::new(),
264 ema_quality: 0.5, prior_sample_size: 0,
266 quality_observations: 0,
267 quality_per_1k_tokens: 0.0,
268 updated_at: now_unix(),
269 }
270 }
271
272 pub fn success_rate_resolved(&self) -> Option<f64> {
281 let resolved = self.success_count + self.fail_count;
282 if resolved == 0 {
283 None
284 } else {
285 Some(self.success_count as f64 / resolved as f64)
286 }
287 }
288
289 pub fn avg_latency_ms(&self) -> f64 {
290 if self.total_calls == 0 {
291 return 0.0;
292 }
293 self.total_latency_ms as f64 / self.total_calls as f64
294 }
295
296 pub fn should_degrade(&self, threshold: u64) -> bool {
298 self.fail_count > self.success_count + threshold
299 }
300
301 pub fn task_stats(&self, task: InferenceTask) -> Option<&TaskStats> {
303 self.task_stats.get(&task.to_string())
304 }
305
306 pub fn total_tokens(&self) -> u64 {
310 self.total_input_tokens
311 + self.total_cache_read_input_tokens
312 + self.total_cache_creation_input_tokens
313 + self.total_output_tokens
314 }
315
316 pub fn compute_quality_per_1k_tokens(&self) -> f64 {
319 let total = self.total_tokens();
320 if total == 0 {
321 return 0.0;
322 }
323 self.ema_quality * 1000.0 / total as f64
324 }
325
326 pub fn tokens_per_success(&self) -> Option<f64> {
333 if self.success_count == 0 {
334 None
335 } else {
336 Some(self.total_tokens() as f64 / self.success_count as f64)
337 }
338 }
339
340 pub fn usd_per_success(
347 &self,
348 input_per_mtok: f64,
349 output_per_mtok: f64,
350 cache: CacheRates,
351 ) -> Option<f64> {
352 if self.success_count == 0 {
353 return None;
354 }
355 let usd = priced_input_usd(
356 self.total_input_tokens,
357 self.total_cache_read_input_tokens,
358 self.total_cache_creation_input_tokens,
359 input_per_mtok,
360 cache,
361 ) + self.total_output_tokens as f64 * output_per_mtok / 1_000_000.0;
362 Some(usd / self.success_count as f64)
363 }
364}
365
366#[derive(Debug, Clone, Copy, PartialEq)]
374pub struct CacheRates {
375 pub read_mult: f64,
377 pub write_mult: f64,
382}
383
384impl CacheRates {
385 pub const ANTHROPIC: Self = Self {
387 read_mult: 0.1,
388 write_mult: 1.25,
389 };
390 pub const OPENAI: Self = Self {
392 read_mult: 0.5,
393 write_mult: 0.0,
394 };
395 pub const NONE: Self = Self {
399 read_mult: 0.0,
400 write_mult: 0.0,
401 };
402}
403
404impl Default for CacheRates {
405 fn default() -> Self {
408 Self::ANTHROPIC
409 }
410}
411
412pub fn priced_input_usd(
418 uncached_input_tokens: u64,
419 cache_read_input_tokens: u64,
420 cache_creation_input_tokens: u64,
421 input_per_mtok: f64,
422 cache: CacheRates,
423) -> f64 {
424 (uncached_input_tokens as f64 * input_per_mtok
425 + cache_read_input_tokens as f64 * input_per_mtok * cache.read_mult
426 + cache_creation_input_tokens as f64 * input_per_mtok * cache.write_mult)
427 / 1_000_000.0
428}
429
430const EMA_ALPHA: f64 = 0.2;
432
433#[derive(Debug, Clone, Serialize, Deserialize)]
440pub struct OutcomeLedgerEntry {
441 pub trace_id: String,
442 pub model_id: String,
443 pub task: InferenceTask,
444 pub routing_reason: String,
445 pub latency_ms: u64,
446 pub input_tokens: usize,
447 pub output_tokens: usize,
448 #[serde(default, skip_serializing_if = "is_zero_usize")]
452 pub cache_read_input_tokens: usize,
453 #[serde(default, skip_serializing_if = "is_zero_usize")]
455 pub cache_creation_input_tokens: usize,
456 #[serde(default, skip_serializing_if = "Option::is_none")]
457 pub success: Option<bool>,
458 #[serde(default, skip_serializing_if = "Option::is_none")]
459 pub quality: Option<f64>,
460 #[serde(default, skip_serializing_if = "Option::is_none")]
461 pub error: Option<String>,
462 #[serde(default, skip_serializing_if = "Option::is_none")]
466 pub project_id: Option<String>,
467 #[serde(default, skip_serializing_if = "Option::is_none")]
470 pub intent: Option<String>,
471 pub timestamp: u64,
472}
473
474fn is_zero_usize(n: &usize) -> bool {
478 *n == 0
479}
480
481const MAX_LEDGER_BUFFER: usize = 5000;
484
485const MAX_LEDGER_ERROR_CHARS: usize = 256;
490
491fn redact_error(error: &str) -> String {
493 if error.chars().count() <= MAX_LEDGER_ERROR_CHARS {
494 return error.to_string();
495 }
496 let truncated: String = error.chars().take(MAX_LEDGER_ERROR_CHARS).collect();
497 format!("{truncated}…")
498}
499
500fn is_no_answer_failure(error: &str) -> bool {
526 let e = error.to_ascii_lowercase();
527 const NEEDLES: &[&str] = &[
528 "timeout",
530 "timed out",
531 "deadline",
532 "connection",
533 "connect error",
534 "reset by peer",
535 "broken pipe",
536 "network",
537 "unreachable",
538 "eof",
539 "stream closed",
540 "socket",
541 "503",
542 "502",
543 "504",
544 "internal server error",
551 "service unavailable",
552 "temporarily unavailable",
553 "overloaded",
554 "unauthorized",
563 "unauthenticated",
564 "forbidden",
565 "permission denied",
566 "authentication",
567 "invalid_client",
568 "invalid api key",
569 "bad request",
570 ];
571 NEEDLES.iter().any(|n| e.contains(n))
572}
573
574pub fn prune_ledger(path: &std::path::Path, max_entries: usize) -> std::io::Result<()> {
578 if max_entries == 0 || !path.exists() {
579 return Ok(());
580 }
581 let entries = read_ledger(path, 0);
582 if entries.len() <= max_entries {
583 return Ok(());
584 }
585 let keep = &entries[entries.len() - max_entries..];
586 let mut body = String::new();
587 for e in keep {
588 let line = serde_json::to_string(e).map_err(std::io::Error::other)?;
589 body.push_str(&line);
590 body.push('\n');
591 }
592 let tmp = path.with_extension("jsonl.tmp");
593 std::fs::write(&tmp, body)?;
594 std::fs::rename(&tmp, path)
595}
596
597pub fn append_ledger_entries(
602 path: &std::path::Path,
603 entries: &[OutcomeLedgerEntry],
604) -> std::io::Result<()> {
605 if entries.is_empty() {
606 return Ok(());
607 }
608 if let Some(parent) = path.parent() {
609 std::fs::create_dir_all(parent)?;
610 }
611 use std::io::Write;
612 let mut opts = std::fs::OpenOptions::new();
613 opts.create(true).append(true);
614 #[cfg(unix)]
617 {
618 use std::os::unix::fs::OpenOptionsExt;
619 opts.mode(0o600);
620 }
621 let existed = path.exists();
624 let mut f = opts.open(path)?;
625 if !existed {
626 car_secrets::harden_owner_only(path);
627 }
628 for e in entries {
629 let line = serde_json::to_string(e).map_err(std::io::Error::other)?;
630 f.write_all(line.as_bytes())?;
631 f.write_all(b"\n")?;
632 }
633 Ok(())
634}
635
636pub fn read_ledger(path: &std::path::Path, limit: usize) -> Vec<OutcomeLedgerEntry> {
640 let Ok(content) = std::fs::read_to_string(path) else {
641 return Vec::new();
642 };
643 let mut out: Vec<OutcomeLedgerEntry> = content
644 .lines()
645 .filter(|l| !l.trim().is_empty())
646 .filter_map(|l| serde_json::from_str(l).ok())
647 .collect();
648 if limit > 0 && out.len() > limit {
649 out = out.split_off(out.len() - limit);
650 }
651 out
652}
653
654pub struct OutcomeTracker {
656 profiles: HashMap<String, ModelProfile>,
658 pending: HashMap<String, InferenceOutcome>,
661 trace_counter: u64,
663 excluded: HashSet<String>,
665 dirty: bool,
671 ledger_buffer: std::collections::VecDeque<OutcomeLedgerEntry>,
674}
675
676impl OutcomeTracker {
677 pub fn new() -> Self {
678 Self {
679 profiles: HashMap::new(),
680 pending: HashMap::new(),
681 trace_counter: 0,
682 excluded: HashSet::new(),
683 dirty: false,
684 ledger_buffer: std::collections::VecDeque::new(),
685 }
686 }
687
688 fn push_ledger(&mut self, entry: OutcomeLedgerEntry) {
690 if self.ledger_buffer.len() >= MAX_LEDGER_BUFFER {
691 self.ledger_buffer.pop_front();
692 }
693 self.ledger_buffer.push_back(entry);
694 }
695
696 pub fn drain_ledger(&mut self) -> Vec<OutcomeLedgerEntry> {
698 self.ledger_buffer.drain(..).collect()
699 }
700
701 pub fn is_excluded(&self, model_id: &str) -> bool {
703 self.excluded.contains(model_id)
704 }
705
706 pub fn record_start(
708 &mut self,
709 model_id: &str,
710 task: InferenceTask,
711 routing_reason: &str,
712 ) -> String {
713 self.trace_counter += 1;
714 let trace_id = format!("t-{}-{}", now_unix(), self.trace_counter);
715
716 let outcome = InferenceOutcome {
717 trace_id: trace_id.clone(),
718 model_id: model_id.to_string(),
719 task,
720 routing_reason: redact_error(routing_reason),
725 latency_ms: 0,
726 input_tokens: 0,
727 output_tokens: 0,
728 cache_read_input_tokens: 0,
729 cache_creation_input_tokens: 0,
730 inferred_outcome: None,
731 code_outcome: None,
732 error: None,
733 timestamp: now_unix(),
734 success_credited: false,
735 };
736
737 self.pending.insert(trace_id.clone(), outcome);
738 trace_id
739 }
740
741 pub fn record_complete(
743 &mut self,
744 trace_id: &str,
745 latency_ms: u64,
746 input_tokens: usize,
747 output_tokens: usize,
748 ) {
749 self.record_complete_cached(trace_id, latency_ms, input_tokens, output_tokens, 0, 0);
750 }
751
752 pub fn record_complete_cached(
760 &mut self,
761 trace_id: &str,
762 latency_ms: u64,
763 input_tokens: usize,
764 output_tokens: usize,
765 cache_read_input_tokens: usize,
766 cache_creation_input_tokens: usize,
767 ) {
768 if let Some(outcome) = self.pending.get_mut(trace_id) {
769 outcome.latency_ms = latency_ms;
770 outcome.input_tokens = input_tokens;
771 outcome.output_tokens = output_tokens;
772 outcome.cache_read_input_tokens = cache_read_input_tokens;
773 outcome.cache_creation_input_tokens = cache_creation_input_tokens;
774
775 let mechanical_success = output_tokens > 0 && outcome.error.is_none();
788 if mechanical_success {
789 outcome.success_credited = true;
790 }
791 let model_id = outcome.model_id.clone();
792 let task_key = outcome.task.to_string();
793
794 let profile = self
797 .profiles
798 .entry(model_id.clone())
799 .or_insert_with(|| ModelProfile::new(model_id));
800
801 profile.total_calls += 1;
802 profile.total_latency_ms += latency_ms;
803 profile.total_input_tokens += input_tokens as u64;
804 profile.total_output_tokens += output_tokens as u64;
805 profile.total_cache_read_input_tokens += cache_read_input_tokens as u64;
806 profile.total_cache_creation_input_tokens += cache_creation_input_tokens as u64;
807 if mechanical_success {
808 profile.success_count += 1;
809 }
810
811 let ts = profile.task_stats.entry(task_key).or_default();
812 ts.calls += 1;
813 if mechanical_success {
814 ts.successes += 1;
815 }
816 ts.avg_latency_ms =
817 ts.avg_latency_ms + (latency_ms as f64 - ts.avg_latency_ms) / ts.calls as f64;
818
819 profile.updated_at = now_unix();
820 self.dirty = true;
821 }
822 }
823
824 pub fn record_failure(&mut self, trace_id: &str, error: &str) {
826 let mut ledger_entry = None;
827 if let Some(outcome) = self.pending.get_mut(trace_id) {
828 outcome.error = Some(error.to_string());
829
830 let profile = self
831 .profiles
832 .entry(outcome.model_id.clone())
833 .or_insert_with(|| ModelProfile::new(outcome.model_id.clone()));
834
835 profile.total_calls += 1;
841 profile.fail_count += 1;
842
843 let is_rate_limited = error.contains("429") || error.contains("RESOURCE_EXHAUSTED");
846 let is_no_answer = !is_rate_limited && is_no_answer_failure(error);
854 if is_rate_limited {
855 self.excluded.insert(outcome.model_id.clone());
857 profile.ema_quality *= 0.1;
858 profile.quality_observations += 1;
860 } else if !is_no_answer {
861 profile.ema_quality = profile.ema_quality * (1.0 - EMA_ALPHA) + 0.0 * EMA_ALPHA;
862 profile.quality_observations += 1;
865 }
866
867 let task_key = outcome.task.to_string();
868 let ts = profile.task_stats.entry(task_key).or_default();
869 ts.failures += 1;
872 if is_rate_limited {
873 ts.ema_quality *= 0.1;
874 ts.quality_observations += 1;
875 } else if !is_no_answer {
876 ts.ema_quality *= 1.0 - EMA_ALPHA;
877 ts.quality_observations += 1;
878 }
879
880 profile.updated_at = now_unix();
881 self.dirty = true;
882
883 ledger_entry = Some(OutcomeLedgerEntry {
884 trace_id: outcome.trace_id.clone(),
885 model_id: outcome.model_id.clone(),
886 task: outcome.task,
887 routing_reason: outcome.routing_reason.clone(),
888 latency_ms: outcome.latency_ms,
889 input_tokens: outcome.input_tokens,
890 output_tokens: outcome.output_tokens,
891 cache_read_input_tokens: outcome.cache_read_input_tokens,
892 cache_creation_input_tokens: outcome.cache_creation_input_tokens,
893 success: Some(false),
894 quality: if is_rate_limited || is_no_answer {
915 None
916 } else {
917 Some(0.0)
918 },
919 error: Some(redact_error(error)),
920 project_id: None,
921 intent: None,
922 timestamp: now_unix(),
923 });
924 }
925
926 if let Some(entry) = ledger_entry {
927 self.push_ledger(entry);
928 }
929
930 self.pending.remove(trace_id);
932 }
933
934 pub fn record_capability_rejection(&mut self, trace_id: &str, error: &str) {
941 self.record_unattributed(trace_id, error);
942 }
943
944 pub fn record_account_rejection(&mut self, trace_id: &str, error: &str) {
953 self.record_unattributed(trace_id, error);
954 }
955
956 fn record_unattributed(&mut self, trace_id: &str, error: &str) {
960 if let Some(outcome) = self.pending.remove(trace_id) {
961 self.push_ledger(OutcomeLedgerEntry {
962 trace_id: outcome.trace_id,
963 model_id: outcome.model_id,
964 task: outcome.task,
965 routing_reason: outcome.routing_reason,
966 latency_ms: outcome.latency_ms,
967 input_tokens: outcome.input_tokens,
968 output_tokens: outcome.output_tokens,
969 cache_read_input_tokens: outcome.cache_read_input_tokens,
970 cache_creation_input_tokens: outcome.cache_creation_input_tokens,
971 success: None,
972 quality: None,
973 error: Some(redact_error(error)),
974 project_id: None,
975 intent: None,
976 timestamp: now_unix(),
977 });
978 }
979 }
980
981 pub fn record_inferred_outcome(&mut self, trace_id: &str, outcome: InferredOutcome) {
983 if let Some(pending) = self.pending.remove(trace_id) {
984 self.apply_outcome(&pending, outcome.quality_score(), outcome.is_success());
985 }
986 }
987
988 pub fn record_code_outcome(&mut self, trace_id: &str, outcome: CodeOutcome) {
990 if let Some(pending) = self.pending.remove(trace_id) {
991 self.apply_outcome(
992 &pending,
993 Some(outcome.quality_score()),
994 Some(outcome.is_success()),
995 );
996 }
997 }
998
999 pub fn resolve_pending_from_signals(&mut self, outcomes: Vec<(String, InferredOutcome)>) {
1002 for (trace_id, inferred) in outcomes {
1003 self.record_inferred_outcome(&trace_id, inferred);
1004 }
1005 }
1006
1007 pub fn infer_outcomes_from_action_sequence(
1015 &self,
1016 action_results: &[(String, bool, f64, String)], ) -> Vec<(String, InferredOutcome)> {
1018 let mut outcomes = Vec::new();
1019
1020 for (i, (trace_id, success, confidence, output)) in action_results.iter().enumerate() {
1021 if trace_id.is_empty() {
1022 continue; }
1024
1025 if !success {
1026 outcomes.push((
1027 trace_id.clone(),
1028 InferredOutcome::Rejected {
1029 confidence: *confidence,
1030 },
1031 ));
1032 continue;
1033 }
1034
1035 let next_succeeded = action_results
1037 .get(i + 1)
1038 .map(|(_, s, _, _)| *s)
1039 .unwrap_or(true); let has_output = !output.trim().is_empty();
1042
1043 if has_output && next_succeeded {
1044 outcomes.push((
1045 trace_id.clone(),
1046 InferredOutcome::Accepted {
1047 confidence: *confidence,
1048 },
1049 ));
1050 } else if has_output && !next_succeeded {
1051 outcomes.push((
1053 trace_id.clone(),
1054 InferredOutcome::AcceptedWithEdits {
1055 confidence: confidence * 0.7,
1056 },
1057 ));
1058 } else {
1059 outcomes.push((trace_id.clone(), InferredOutcome::Inconclusive));
1060 }
1061 }
1062
1063 outcomes
1064 }
1065
1066 pub fn profile(&self, model_id: &str) -> Option<&ModelProfile> {
1068 self.profiles.get(model_id)
1069 }
1070
1071 pub fn has_pending(&self, trace_id: &str) -> bool {
1076 self.pending.contains_key(trace_id)
1077 }
1078
1079 pub fn all_profiles(&self) -> &HashMap<String, ModelProfile> {
1081 &self.profiles
1082 }
1083
1084 pub fn pending_trace_ids(&self) -> Vec<String> {
1086 self.pending.keys().cloned().collect()
1087 }
1088
1089 pub fn get_pending(&self, trace_id: &str) -> Option<&InferenceOutcome> {
1091 self.pending.get(trace_id)
1092 }
1093
1094 pub fn export_profiles(&self) -> Vec<ModelProfile> {
1098 self.profiles
1099 .values()
1100 .cloned()
1101 .map(|mut p| {
1102 p.quality_per_1k_tokens = p.compute_quality_per_1k_tokens();
1103 p
1104 })
1105 .collect()
1106 }
1107
1108 pub fn import_profiles(&mut self, profiles: Vec<ModelProfile>) {
1115 for p in profiles {
1116 self.profiles.insert(p.model_id.clone(), p);
1117 }
1118 self.dirty = true;
1119 }
1120
1121 pub fn save_to_file(&self, path: &std::path::Path) -> Result<(), std::io::Error> {
1129 let profiles = self.export_profiles();
1130 let json = serde_json::to_string_pretty(&profiles).map_err(std::io::Error::other)?;
1131 if let Some(parent) = path.parent() {
1132 std::fs::create_dir_all(parent)?;
1133 }
1134 let tmp = path.with_extension("json.tmp");
1135 std::fs::write(&tmp, json)?;
1136 std::fs::rename(&tmp, path)
1137 }
1138
1139 pub fn is_dirty(&self) -> bool {
1141 self.dirty
1142 }
1143
1144 pub fn save_if_dirty(&mut self, path: &std::path::Path) -> Result<bool, std::io::Error> {
1148 if !self.dirty {
1149 return Ok(false);
1150 }
1151 self.save_to_file(path)?;
1152 self.dirty = false;
1153 Ok(true)
1154 }
1155
1156 pub fn load_from_file(&mut self, path: &std::path::Path) -> Result<usize, std::io::Error> {
1158 if !path.exists() {
1159 return Ok(0);
1160 }
1161 let json = std::fs::read_to_string(path)?;
1162 let profiles: Vec<ModelProfile> = serde_json::from_str(&json)
1163 .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidData, e))?;
1164 let count = profiles.len();
1165 for mut p in profiles {
1169 backfill_quality_observations(&mut p);
1178 self.profiles.insert(p.model_id.clone(), p);
1179 }
1180 Ok(count)
1181 }
1182
1183 fn apply_outcome(
1185 &mut self,
1186 pending: &InferenceOutcome,
1187 quality: Option<f64>,
1188 success: Option<bool>,
1189 ) {
1190 let profile = self
1191 .profiles
1192 .entry(pending.model_id.clone())
1193 .or_insert_with(|| ModelProfile::new(pending.model_id.clone()));
1194
1195 if let Some(q) = quality {
1196 profile.ema_quality = profile.ema_quality * (1.0 - EMA_ALPHA) + q * EMA_ALPHA;
1197 profile.quality_observations += 1;
1198
1199 let task_key = pending.task.to_string();
1200 let ts = profile.task_stats.entry(task_key).or_default();
1201 ts.ema_quality = ts.ema_quality * (1.0 - EMA_ALPHA) + q * EMA_ALPHA;
1202 ts.quality_observations += 1;
1203 }
1204
1205 if let Some(ok) = success {
1206 let already_credited = pending.success_credited;
1209 let task_key = pending.task.to_string();
1210 if ok {
1211 if !already_credited {
1212 profile.success_count += 1;
1213 let ts = profile.task_stats.entry(task_key).or_default();
1214 ts.successes += 1;
1215 }
1216 } else {
1218 if already_credited {
1222 profile.success_count = profile.success_count.saturating_sub(1);
1223 let ts = profile.task_stats.entry(task_key.clone()).or_default();
1224 ts.successes = ts.successes.saturating_sub(1);
1225 }
1226 profile.fail_count += 1;
1227 let ts = profile.task_stats.entry(task_key).or_default();
1228 ts.failures += 1;
1229 }
1230 }
1231
1232 profile.updated_at = now_unix();
1233 self.dirty = true;
1234
1235 self.push_ledger(OutcomeLedgerEntry {
1236 trace_id: pending.trace_id.clone(),
1237 model_id: pending.model_id.clone(),
1238 task: pending.task,
1239 routing_reason: pending.routing_reason.clone(),
1240 latency_ms: pending.latency_ms,
1241 input_tokens: pending.input_tokens,
1242 output_tokens: pending.output_tokens,
1243 cache_read_input_tokens: pending.cache_read_input_tokens,
1244 cache_creation_input_tokens: pending.cache_creation_input_tokens,
1245 success,
1246 quality,
1247 error: None,
1248 project_id: None,
1249 intent: None,
1250 timestamp: now_unix(),
1251 });
1252 }
1253
1254 pub fn sweep_pending(&mut self, ttl_secs: u64) -> usize {
1264 self.sweep_pending_at(ttl_secs, now_unix())
1265 }
1266
1267 fn sweep_pending_at(&mut self, ttl_secs: u64, now: u64) -> usize {
1269 let cutoff = now.saturating_sub(ttl_secs);
1270 let expired: Vec<String> = self
1271 .pending
1272 .iter()
1273 .filter(|(_, o)| o.timestamp < cutoff)
1274 .map(|(id, _)| id.clone())
1275 .collect();
1276 for id in &expired {
1277 if let Some(o) = self.pending.remove(id) {
1278 if o.latency_ms > 0 {
1282 let credit_now =
1300 o.output_tokens > 0 && o.error.is_none() && !o.success_credited;
1301 let was_success = o.success_credited || credit_now;
1302 if credit_now {
1303 let profile = self
1304 .profiles
1305 .entry(o.model_id.clone())
1306 .or_insert_with(|| ModelProfile::new(o.model_id.clone()));
1307 profile.success_count += 1;
1308 let ts = profile.task_stats.entry(o.task.to_string()).or_default();
1309 ts.successes += 1;
1310 profile.updated_at = now_unix();
1311 self.dirty = true;
1312 }
1313 self.push_ledger(OutcomeLedgerEntry {
1314 trace_id: o.trace_id,
1315 model_id: o.model_id,
1316 task: o.task,
1317 routing_reason: o.routing_reason,
1318 latency_ms: o.latency_ms,
1319 input_tokens: o.input_tokens,
1320 output_tokens: o.output_tokens,
1321 cache_read_input_tokens: o.cache_read_input_tokens,
1322 cache_creation_input_tokens: o.cache_creation_input_tokens,
1323 success: if was_success { Some(true) } else { None },
1324 quality: None,
1325 error: None,
1326 project_id: None,
1327 intent: None,
1328 timestamp: now_unix(),
1329 });
1330 }
1331 }
1332 }
1333 expired.len()
1334 }
1335
1336 pub fn check_git_outcomes(&mut self, repo_dir: &std::path::Path) {
1344 let diff = match std::process::Command::new("git")
1345 .args(["diff", "--no-color"])
1346 .current_dir(repo_dir)
1347 .output()
1348 {
1349 Ok(output) => String::from_utf8_lossy(&output.stdout).to_string(),
1350 Err(_) => return,
1351 };
1352
1353 let staged_diff = match std::process::Command::new("git")
1354 .args(["diff", "--cached", "--no-color"])
1355 .current_dir(repo_dir)
1356 .output()
1357 {
1358 Ok(output) => String::from_utf8_lossy(&output.stdout).to_string(),
1359 Err(_) => String::new(),
1360 };
1361
1362 let combined_diff = format!("{}\n{}", diff, staged_diff);
1363
1364 if combined_diff.trim().is_empty() {
1365 return; }
1367
1368 #[cfg(feature = "ast")]
1370 let ast_outcome = Self::check_git_outcomes_ast(repo_dir);
1371
1372 let code_traces: Vec<(String, String)> = self
1373 .pending
1374 .iter()
1375 .filter(|(_, o)| matches!(o.task, InferenceTask::Code))
1376 .map(|(id, o)| (id.clone(), o.model_id.clone()))
1377 .collect();
1378
1379 for (trace_id, _model_id) in code_traces {
1380 if let Some(pending) = self.pending.get(&trace_id) {
1381 #[cfg(feature = "ast")]
1383 if let Some(ref ast_out) = ast_outcome {
1384 let pending_clone = pending.clone();
1385 self.apply_outcome(
1386 &pending_clone,
1387 Some(ast_out.quality_score()),
1388 Some(ast_out.is_success()),
1389 );
1390 continue;
1391 }
1392
1393 let output_tokens: Vec<&str> = pending
1395 .routing_reason
1396 .split_whitespace()
1397 .filter(|t| t.len() > 5)
1398 .collect();
1399
1400 let outcome = if output_tokens.iter().any(|t| combined_diff.contains(t)) {
1401 CodeOutcome::Applied
1402 } else {
1403 CodeOutcome::Modified
1404 };
1405
1406 let pending_clone = pending.clone();
1407 self.apply_outcome(
1408 &pending_clone,
1409 Some(outcome.quality_score()),
1410 Some(outcome.is_success()),
1411 );
1412 }
1413 }
1414 }
1415
1416 #[cfg(feature = "ast")]
1418 fn check_git_outcomes_ast(repo_dir: &std::path::Path) -> Option<CodeOutcome> {
1419 let name_only = std::process::Command::new("git")
1421 .args(["diff", "--name-only"])
1422 .current_dir(repo_dir)
1423 .output()
1424 .ok()?;
1425 let changed_files: Vec<&str> = std::str::from_utf8(&name_only.stdout)
1426 .ok()?
1427 .lines()
1428 .filter(|f| !f.is_empty())
1429 .collect();
1430
1431 if changed_files.is_empty() {
1432 return None;
1433 }
1434
1435 let mut has_sig_change = false;
1436 let mut has_body_change = false;
1437 let mut has_addition = false;
1438
1439 for file in &changed_files {
1440 if car_ast::Language::from_filename(file).is_none() {
1442 continue;
1443 }
1444
1445 let old_content = std::process::Command::new("git")
1447 .args(["show", &format!("HEAD:{}", file)])
1448 .current_dir(repo_dir)
1449 .output()
1450 .ok()
1451 .and_then(|o| {
1452 if o.status.success() {
1453 String::from_utf8(o.stdout).ok()
1454 } else {
1455 None
1456 }
1457 });
1458
1459 let new_path = repo_dir.join(file);
1461 let new_content = std::fs::read_to_string(&new_path).ok();
1462
1463 match (old_content, new_content) {
1464 (Some(old), Some(new)) => {
1465 let old_parsed = car_ast::parse_file(&old, file);
1466 let new_parsed = car_ast::parse_file(&new, file);
1467
1468 if let (Some(old_p), Some(new_p)) = (old_parsed, new_parsed) {
1469 let changes = car_ast::diff_symbols(&old_p, &new_p);
1470 for change in &changes {
1471 match change {
1472 car_ast::SymbolChange::Added(_) => has_addition = true,
1473 car_ast::SymbolChange::Modified {
1474 signature_changed, ..
1475 } => {
1476 if *signature_changed {
1477 has_sig_change = true;
1478 } else {
1479 has_body_change = true;
1480 }
1481 }
1482 car_ast::SymbolChange::Removed(_) => has_sig_change = true,
1483 }
1484 }
1485 }
1486 }
1487 (None, Some(_)) => has_addition = true, _ => {}
1489 }
1490 }
1491
1492 if has_sig_change {
1494 Some(CodeOutcome::SignatureChanged)
1495 } else if has_body_change {
1496 Some(CodeOutcome::BodyModified)
1497 } else if has_addition {
1498 Some(CodeOutcome::SymbolAdded)
1499 } else {
1500 None }
1502 }
1503}
1504
1505impl Default for OutcomeTracker {
1506 fn default() -> Self {
1507 Self::new()
1508 }
1509}
1510
1511fn now_unix() -> u64 {
1512 SystemTime::now()
1513 .duration_since(UNIX_EPOCH)
1514 .unwrap_or_default()
1515 .as_secs()
1516}
1517
1518#[cfg(test)]
1519mod tests {
1520 use super::*;
1521
1522 #[test]
1523 fn lifecycle() {
1524 let mut tracker = OutcomeTracker::new();
1525
1526 let trace = tracker.record_start(
1528 "qwen/qwen3-4b:q4_k_m",
1529 InferenceTask::Code,
1530 "Code task -> Qwen3-4B",
1531 );
1532
1533 tracker.record_complete(&trace, 1200, 100, 50);
1535
1536 let profile = tracker.profile("qwen/qwen3-4b:q4_k_m").unwrap();
1538 assert_eq!(profile.total_calls, 1);
1539 assert_eq!(profile.avg_latency_ms(), 1200.0);
1540
1541 tracker.record_inferred_outcome(&trace, InferredOutcome::Accepted { confidence: 0.9 });
1543
1544 let profile = tracker.profile("qwen/qwen3-4b:q4_k_m").unwrap();
1545 assert_eq!(profile.success_count, 1);
1546 assert!(profile.ema_quality > 0.5); }
1548
1549 #[test]
1550 fn failure_degrades() {
1551 let mut tracker = OutcomeTracker::new();
1556 for _ in 0..5 {
1557 let trace = tracker.record_start("bad-model", InferenceTask::Generate, "test");
1558 tracker.record_failure(&trace, "model produced malformed output");
1561 }
1562
1563 let profile = tracker.profile("bad-model").unwrap();
1564 assert_eq!(profile.fail_count, 5);
1565 assert_eq!(profile.success_count, 0);
1566 assert!(profile.should_degrade(2)); assert!(profile.ema_quality < 0.3); }
1569
1570 #[test]
1571 fn code_outcome_ground_truth() {
1572 let mut tracker = OutcomeTracker::new();
1573
1574 let trace = tracker.record_start("qwen/qwen3-4b:q4_k_m", InferenceTask::Code, "code");
1575 tracker.record_complete(&trace, 500, 200, 100);
1576 tracker.record_code_outcome(&trace, CodeOutcome::Applied);
1577
1578 let profile = tracker.profile("qwen/qwen3-4b:q4_k_m").unwrap();
1579 assert_eq!(profile.success_count, 1);
1580 assert!((profile.ema_quality - 0.6).abs() < 0.01);
1582 }
1583
1584 #[test]
1585 fn per_task_stats() {
1586 let mut tracker = OutcomeTracker::new();
1587
1588 for _ in 0..2 {
1590 let trace = tracker.record_start("m1", InferenceTask::Code, "code");
1591 tracker.record_complete(&trace, 1000, 100, 50);
1592 tracker.record_inferred_outcome(&trace, InferredOutcome::Accepted { confidence: 0.8 });
1593 }
1594 let trace = tracker.record_start("m1", InferenceTask::Generate, "gen");
1595 tracker.record_complete(&trace, 500, 50, 25);
1596 tracker.record_inferred_outcome(&trace, InferredOutcome::Rejected { confidence: 0.9 });
1597
1598 let profile = tracker.profile("m1").unwrap();
1599 assert_eq!(profile.total_calls, 3);
1600
1601 let code_stats = profile.task_stats(InferenceTask::Code).unwrap();
1602 assert_eq!(code_stats.calls, 2);
1603 assert_eq!(code_stats.successes, 2);
1604
1605 let gen_stats = profile.task_stats(InferenceTask::Generate).unwrap();
1606 assert_eq!(gen_stats.calls, 1);
1607 assert_eq!(gen_stats.failures, 1);
1608 }
1609
1610 #[test]
1611 fn export_populates_quality_per_1k_tokens() {
1612 let mut tracker = OutcomeTracker::new();
1613 let trace = tracker.record_start("m1", InferenceTask::Generate, "test");
1614 tracker.record_complete(&trace, 100, 800, 200); tracker.record_inferred_outcome(&trace, InferredOutcome::Accepted { confidence: 1.0 });
1616
1617 let exported = tracker.export_profiles();
1618 assert_eq!(exported.len(), 1);
1619 let p = &exported[0];
1620 assert!(
1623 (p.quality_per_1k_tokens - 0.6).abs() < 1e-6,
1624 "got {}",
1625 p.quality_per_1k_tokens
1626 );
1627 }
1628
1629 #[test]
1630 fn quality_per_1k_tokens_zero_without_tokens() {
1631 let profile = ModelProfile::new("x".into());
1632 assert_eq!(profile.compute_quality_per_1k_tokens(), 0.0);
1633 }
1634
1635 #[test]
1636 fn tokens_per_success_is_outcome_denominated() {
1637 let mut p = ModelProfile::new("x".into());
1638 assert_eq!(p.tokens_per_success(), None);
1640 p.total_input_tokens = 600;
1642 p.total_output_tokens = 300;
1643 p.success_count = 3;
1644 assert_eq!(p.tokens_per_success(), Some(300.0));
1645 let usd = p.usd_per_success(1.0, 2.0, CacheRates::ANTHROPIC).unwrap();
1647 assert!((usd - ((600.0 * 1.0 + 300.0 * 2.0) / 1_000_000.0 / 3.0)).abs() < 1e-12);
1648 assert_eq!(
1649 ModelProfile::new("y".into()).usd_per_success(1.0, 2.0, CacheRates::ANTHROPIC),
1650 None
1651 );
1652 }
1653
1654 #[test]
1655 fn usd_per_success_prices_cache_buckets_separately() {
1656 let mut p = ModelProfile::new("cached".into());
1657 p.success_count = 1;
1658 p.total_input_tokens = 100; p.total_cache_read_input_tokens = 1000; p.total_cache_creation_input_tokens = 200; p.total_output_tokens = 50; let expected_input = (100.0 * 1.0 + 1000.0 * 1.0 * 0.1 + 200.0 * 1.0 * 1.25) / 1_000_000.0;
1664 let expected = expected_input + 50.0 * 2.0 / 1_000_000.0;
1665 let usd = p.usd_per_success(1.0, 2.0, CacheRates::ANTHROPIC).unwrap();
1666 assert!((usd - expected).abs() < 1e-12, "got {usd}, want {expected}");
1667 assert_eq!(p.total_tokens(), 100 + 1000 + 200 + 50);
1669 let mut naive = ModelProfile::new("naive".into());
1672 naive.success_count = 1;
1673 naive.total_input_tokens = 1300; naive.total_output_tokens = 50;
1675 assert!(
1676 naive
1677 .usd_per_success(1.0, 2.0, CacheRates::ANTHROPIC)
1678 .unwrap()
1679 > usd
1680 );
1681 }
1682
1683 #[test]
1684 fn openai_cache_rates_price_reads_at_half_no_write_premium() {
1685 let mut p = ModelProfile::new("gpt".into());
1688 p.success_count = 1;
1689 p.total_input_tokens = 100; p.total_cache_read_input_tokens = 1000; p.total_cache_creation_input_tokens = 0; let usd = p.usd_per_success(1.0, 2.0, CacheRates::OPENAI).unwrap();
1693 let expected = (100.0 * 1.0 + 1000.0 * 1.0 * 0.5) / 1_000_000.0;
1694 assert!((usd - expected).abs() < 1e-12, "got {usd}, want {expected}");
1695 let anthropic = p.usd_per_success(1.0, 2.0, CacheRates::ANTHROPIC).unwrap();
1697 assert!(
1698 usd > anthropic,
1699 "OpenAI 0.5× read must exceed Anthropic 0.1×"
1700 );
1701 let mut naive = ModelProfile::new("naive".into());
1703 naive.success_count = 1;
1704 naive.total_input_tokens = 1100;
1705 assert!(naive.usd_per_success(1.0, 2.0, CacheRates::OPENAI).unwrap() > usd);
1706 }
1707
1708 #[test]
1709 fn record_complete_cached_accumulates_cache_totals() {
1710 let mut tracker = OutcomeTracker::new();
1711 let trace = tracker.record_start("m", InferenceTask::Generate, "gen");
1712 tracker.record_complete_cached(&trace, 100, 40, 20, 800, 120);
1713 let p = tracker.profile("m").unwrap();
1714 assert_eq!(p.total_input_tokens, 40, "uncached prefix only");
1715 assert_eq!(p.total_cache_read_input_tokens, 800);
1716 assert_eq!(p.total_cache_creation_input_tokens, 120);
1717 }
1718
1719 #[test]
1720 fn dirty_flag_and_save_if_dirty() {
1721 let dir = std::env::temp_dir().join("car-outcome-dirty-test");
1722 let _ = std::fs::remove_dir_all(&dir);
1723 let path = dir.join("outcome_profiles.json");
1724
1725 let mut tracker = OutcomeTracker::new();
1726 assert!(!tracker.is_dirty());
1728 assert!(!tracker.save_if_dirty(&path).unwrap());
1729 assert!(!path.exists());
1730
1731 let trace = tracker.record_start("m1", InferenceTask::Generate, "router");
1733 tracker.record_complete(&trace, 100, 10, 20);
1734 assert!(tracker.is_dirty());
1735
1736 assert!(tracker.save_if_dirty(&path).unwrap());
1738 assert!(path.exists());
1739 assert!(!tracker.is_dirty());
1740
1741 assert!(!tracker.save_if_dirty(&path).unwrap());
1743
1744 let mut fresh = OutcomeTracker::new();
1746 fresh.load_from_file(&path).unwrap();
1747 assert!(!fresh.is_dirty());
1748
1749 fresh.import_profiles(vec![ModelProfile::new("seeded".into())]);
1752 assert!(fresh.is_dirty());
1753
1754 let _ = std::fs::remove_dir_all(&dir);
1755 }
1756
1757 #[test]
1758 fn ledger_captures_resolved_outcomes() {
1759 let dir = std::env::temp_dir().join("car-outcome-ledger-test");
1760 let _ = std::fs::remove_dir_all(&dir);
1761 let path = dir.join("outcome_ledger.jsonl");
1762
1763 let mut tracker = OutcomeTracker::new();
1764
1765 let t1 = tracker.record_start("good-model", InferenceTask::Generate, "router:test");
1767 tracker.record_complete(&t1, 1200, 50, 100);
1768 tracker.record_inferred_outcome(&t1, InferredOutcome::Accepted { confidence: 0.9 });
1769
1770 let t2 = tracker.record_start("bad-model", InferenceTask::Code, "router:test");
1772 tracker.record_failure(&t2, "boom: 500");
1773
1774 let drained = tracker.drain_ledger();
1775 assert_eq!(drained.len(), 2);
1776 assert!(tracker.drain_ledger().is_empty(), "drain clears the buffer");
1777
1778 append_ledger_entries(&path, &drained).unwrap();
1779 let read = read_ledger(&path, 0);
1780 assert_eq!(read.len(), 2);
1781
1782 let good = read.iter().find(|e| e.model_id == "good-model").unwrap();
1783 assert_eq!(good.success, Some(true));
1784 assert!(good.quality.is_some());
1785 assert_eq!(good.routing_reason, "router:test");
1786 assert_eq!(good.latency_ms, 1200);
1787
1788 let bad = read.iter().find(|e| e.model_id == "bad-model").unwrap();
1789 assert_eq!(bad.success, Some(false));
1790 assert_eq!(bad.error.as_deref(), Some("boom: 500"));
1791
1792 assert_eq!(read_ledger(&path, 1).len(), 1);
1794
1795 let _ = std::fs::remove_dir_all(&dir);
1796 }
1797
1798 #[test]
1799 fn ledger_redacts_long_errors_and_prunes() {
1800 let dir = std::env::temp_dir().join("car-outcome-privacy-test");
1801 let _ = std::fs::remove_dir_all(&dir);
1802 let path = dir.join("outcome_ledger.jsonl");
1803
1804 let mut tracker = OutcomeTracker::new();
1806 let t = tracker.record_start("m", InferenceTask::Generate, "r");
1807 let huge = "x".repeat(5000);
1808 tracker.record_failure(&t, &huge);
1809 let drained = tracker.drain_ledger();
1810 let err = drained[0].error.as_ref().unwrap();
1811 assert!(
1812 err.chars().count() <= MAX_LEDGER_ERROR_CHARS + 1,
1813 "error truncated"
1814 );
1815
1816 let entries: Vec<OutcomeLedgerEntry> = (0..10)
1818 .map(|i| OutcomeLedgerEntry {
1819 trace_id: format!("t{i}"),
1820 model_id: "m".into(),
1821 task: InferenceTask::Generate,
1822 routing_reason: "r".into(),
1823 latency_ms: 1,
1824 input_tokens: 1,
1825 output_tokens: 1,
1826 cache_read_input_tokens: 0,
1827 cache_creation_input_tokens: 0,
1828 success: Some(true),
1829 quality: Some(1.0),
1830 error: None,
1831 project_id: None,
1832 intent: None,
1833 timestamp: i,
1834 })
1835 .collect();
1836 append_ledger_entries(&path, &entries).unwrap();
1837 prune_ledger(&path, 3).unwrap();
1838 let kept = read_ledger(&path, 0);
1839 assert_eq!(kept.len(), 3);
1840 assert_eq!(kept[0].trace_id, "t7"); assert_eq!(kept[2].trace_id, "t9");
1842
1843 let _ = std::fs::remove_dir_all(&dir);
1844 }
1845
1846 #[test]
1847 fn sweep_pending_credits_mechanical_success_or_inconclusive() {
1848 let mut tracker = OutcomeTracker::new();
1849
1850 let t1 = tracker.record_start("m", InferenceTask::Generate, "r");
1852 tracker.record_complete(&t1, 500, 10, 20);
1853 let t2 = tracker.record_start("m", InferenceTask::Generate, "r");
1855 tracker.record_complete(&t2, 300, 5, 0);
1856 let _t3 = tracker.record_start("m", InferenceTask::Generate, "r");
1858
1859 let swept = tracker.sweep_pending_at(0, now_unix() + 10);
1862 assert_eq!(swept, 3, "all pending entries evicted");
1863
1864 let mut receipts = tracker.drain_ledger();
1866 receipts.sort_by_key(|r| r.latency_ms);
1867 assert_eq!(receipts.len(), 2);
1868 assert_eq!(receipts[0].latency_ms, 300);
1870 assert_eq!(receipts[0].success, None);
1871 assert_eq!(receipts[1].latency_ms, 500);
1873 assert_eq!(receipts[1].success, Some(true));
1874 assert_eq!(receipts[1].quality, None);
1875
1876 let p = tracker.profile("m").expect("profile exists");
1879 assert_eq!(p.success_count, 1);
1880 assert_eq!(p.ema_quality, 0.5);
1881 }
1882
1883 #[test]
1884 fn record_complete_credits_success_immediately() {
1885 let mut tracker = OutcomeTracker::new();
1890 let t = tracker.record_start("m", InferenceTask::Generate, "r");
1891 tracker.record_complete(&t, 500, 12, 20);
1892
1893 let p = tracker.profile("m").expect("profile exists");
1894 assert_eq!(p.success_count, 1, "success credited at completion");
1895 assert_eq!(p.total_calls, 1);
1896 assert_eq!(p.total_input_tokens, 12, "input tokens recorded, not 0");
1897 assert_eq!(p.fail_count, 0);
1898
1899 tracker.sweep_pending_at(0, now_unix() + 10);
1901 let p = tracker.profile("m").unwrap();
1902 assert_eq!(p.success_count, 1, "sweep does not re-credit");
1903 }
1904
1905 #[test]
1906 fn success_rate_resolved_distinguishes_no_signal_from_measured() {
1907 let mut p = ModelProfile::new("m".to_string());
1912 assert_eq!(
1913 p.success_rate_resolved(),
1914 None,
1915 "no resolved signal, not a fabricated 0.5"
1916 );
1917
1918 p.success_count = 3;
1919 p.fail_count = 1;
1920 assert_eq!(
1921 p.success_rate_resolved(),
1922 Some(0.75),
1923 "real rate once resolved"
1924 );
1925 }
1926
1927 #[test]
1928 fn quality_observations_count_only_graded_signals() {
1929 let mut tracker = OutcomeTracker::new();
1933
1934 let t = tracker.record_start("m", InferenceTask::Generate, "r");
1936 tracker.record_complete(&t, 500, 12, 20);
1937 let p = tracker.profile("m").unwrap();
1938 assert_eq!(p.success_count, 1);
1939 assert_eq!(
1940 p.quality_observations, 0,
1941 "mechanical success is not a graded quality observation"
1942 );
1943 assert!(
1944 (p.ema_quality - 0.5).abs() < 1e-9,
1945 "EMA untouched by mechanical success"
1946 );
1947
1948 tracker.record_inferred_outcome(&t, InferredOutcome::Accepted { confidence: 0.9 });
1950 let p = tracker.profile("m").unwrap();
1951 assert_eq!(p.quality_observations, 1, "graded accept signal counts");
1952 assert!(p.ema_quality > 0.5, "graded accept moved the EMA up");
1953 assert_eq!(
1954 p.task_stats(InferenceTask::Generate)
1955 .unwrap()
1956 .quality_observations,
1957 1,
1958 "per-task graded count tracked too"
1959 );
1960
1961 let t2 = tracker.record_start("m", InferenceTask::Generate, "r");
1963 tracker.record_failure(&t2, "boom");
1964 let p = tracker.profile("m").unwrap();
1965 assert_eq!(p.quality_observations, 2, "failure is a graded observation");
1966 }
1967
1968 #[test]
1969 fn record_complete_no_output_is_not_a_success() {
1970 let mut tracker = OutcomeTracker::new();
1971 let t = tracker.record_start("m", InferenceTask::Generate, "r");
1972 tracker.record_complete(&t, 300, 5, 0); let p = tracker.profile("m").unwrap();
1974 assert_eq!(p.success_count, 0, "no output -> no mechanical success");
1975 assert_eq!(p.total_calls, 1);
1976 }
1977
1978 #[test]
1979 fn record_failure_counts_total_calls() {
1980 let mut tracker = OutcomeTracker::new();
1984 for _ in 0..3 {
1985 let t = tracker.record_start("m", InferenceTask::Generate, "r");
1986 tracker.record_failure(&t, "boom 500");
1987 }
1988 let p = tracker.profile("m").unwrap();
1989 assert_eq!(p.fail_count, 3);
1990 assert_eq!(p.total_calls, 3, "failures counted in total_calls");
1991 assert!(p.fail_count <= p.total_calls);
1992 }
1993
1994 #[test]
1995 fn capability_rejection_does_not_change_model_profile() {
1996 let mut tracker = OutcomeTracker::new();
1997 let trace = tracker.record_start("m", InferenceTask::Generate, "r");
1998
1999 tracker.record_capability_rejection(&trace, "unsupported mode: json schema");
2000
2001 assert!(tracker.profile("m").is_none());
2002 assert!(!tracker.has_pending(&trace));
2003 assert_eq!(tracker.ledger_buffer.len(), 1);
2004 }
2005
2006 #[test]
2007 fn transport_failures_are_quality_neutral() {
2008 let mut tracker = OutcomeTracker::new();
2015 let t = tracker.record_start("m", InferenceTask::Code, "r");
2018 tracker.record_failure(&t, "model produced malformed output");
2019 let baseline_ema = tracker.profile("m").unwrap().ema_quality;
2020 let baseline_qobs = tracker.profile("m").unwrap().quality_observations;
2021
2022 for _ in 0..5 {
2027 let t = tracker.record_start("m", InferenceTask::Code, "r");
2028 tracker.record_failure(&t, "daemon read timeout on infer after 30s");
2029 }
2030 for _ in 0..5 {
2031 let t = tracker.record_start("m", InferenceTask::Code, "r");
2032 tracker.record_failure(
2033 &t,
2034 "inference failed: Parslee org lookup failed: HTTP 401 Unauthorized: Authentication",
2035 );
2036 }
2037
2038 let p = tracker.profile("m").unwrap();
2039 assert_eq!(
2040 p.fail_count, 11,
2041 "no-answer failures still counted for availability"
2042 );
2043 assert_eq!(p.total_calls, 11, "every call counted (failures included)");
2044 assert!(
2045 (p.ema_quality - baseline_ema).abs() < 1e-9,
2046 "timeouts AND auth rejections must not move the quality EMA ({baseline_ema} -> {})",
2047 p.ema_quality
2048 );
2049 assert_eq!(
2050 p.quality_observations, baseline_qobs,
2051 "no-answer failures (transport + auth) are not graded quality evidence"
2052 );
2053 let t = tracker.record_start("m", InferenceTask::Code, "r");
2055 tracker.record_failure(&t, "model produced malformed output");
2056 assert!(
2057 tracker.profile("m").unwrap().ema_quality < baseline_ema,
2058 "a non-transport failure must still move the quality EMA down"
2059 );
2060 }
2061
2062 #[test]
2063 fn is_no_answer_failure_precision_boundary() {
2064 for s in [
2069 "daemon read timeout on infer after 30s",
2070 "connection reset by peer",
2071 "upstream returned 503 service unavailable",
2072 "HTTP 500 Internal Server Error",
2073 "stream closed unexpectedly (eof)",
2074 "HTTP 401 Unauthorized",
2075 "unauthenticated request",
2076 "HTTP 403 Forbidden",
2077 "permission denied by gateway",
2078 "Parslee org lookup failed: Authentication required",
2079 "oauth: invalid_client",
2080 "invalid api key",
2081 "HTTP 400 Bad Request: max context length exceeded",
2082 ] {
2083 assert!(
2084 is_no_answer_failure(s),
2085 "expected no-answer (quality-neutral): {s:?}"
2086 );
2087 }
2088 for s in [
2092 "model produced malformed output",
2093 "field 'user' not found in model response", "tokenizer: invalid token id 99999", "response failed JSON schema validation",
2096 "tool call arguments did not validate",
2097 ] {
2098 assert!(
2099 !is_no_answer_failure(s),
2100 "expected generation failure (graded): {s:?}"
2101 );
2102 }
2103 }
2104
2105 #[test]
2106 fn ledger_quality_distinguishes_generation_from_transport_failure() {
2107 let mut tracker = OutcomeTracker::new();
2114
2115 let t = tracker.record_start("m", InferenceTask::Code, "r");
2116 tracker.record_failure(&t, "model produced malformed output"); let t = tracker.record_start("m", InferenceTask::Code, "r");
2118 tracker.record_failure(&t, "daemon read timeout on infer after 30s"); let t = tracker.record_start("m", InferenceTask::Code, "r");
2120 tracker.record_failure(&t, "429 RESOURCE_EXHAUSTED"); let t = tracker.record_start("m", InferenceTask::Code, "r");
2122 tracker.record_failure(&t, "upstream: 503 service unavailable"); let t = tracker.record_start("m", InferenceTask::Code, "r");
2124 tracker.record_failure(&t, "HTTP 500 Internal Server Error"); let t = tracker.record_start("m", InferenceTask::Code, "r");
2126 tracker.record_failure(
2128 &t,
2129 "inference failed: Parslee org lookup failed: HTTP 401 Unauthorized: Authentication",
2130 );
2131
2132 let entries = tracker.drain_ledger();
2133 assert_eq!(entries.len(), 6);
2134 assert!(entries.iter().all(|e| e.success == Some(false)));
2136 assert_eq!(
2138 entries[0].quality,
2139 Some(0.0),
2140 "generation failure -> graded 0.0 answer quality"
2141 );
2142 assert_eq!(
2143 entries[1].quality, None,
2144 "transport failure produced no answer -> quality unknown (None)"
2145 );
2146 assert_eq!(
2147 entries[2].quality, None,
2148 "429 produced no answer -> quality unknown (None), an availability event"
2149 );
2150 assert_eq!(
2151 entries[3].quality, None,
2152 "5xx infra failure produced no answer -> None"
2153 );
2154 assert_eq!(
2155 entries[4].quality, None,
2156 "an 'internal server error' (500) is a server fault, not a bad answer -> None"
2157 );
2158 assert_eq!(
2159 entries[5].quality, None,
2160 "a 401 auth rejection produced no answer -> None (not a bad-answer 0.0)"
2161 );
2162 }
2163
2164 #[test]
2165 fn real_failure_signal_reclassifies_mechanical_success() {
2166 let mut tracker = OutcomeTracker::new();
2169 let t = tracker.record_start("m", InferenceTask::Generate, "r");
2170 tracker.record_complete(&t, 100, 8, 15);
2171 assert_eq!(tracker.profile("m").unwrap().success_count, 1);
2172
2173 tracker.record_inferred_outcome(&t, InferredOutcome::Rejected { confidence: 0.9 });
2174 let p = tracker.profile("m").unwrap();
2175 assert_eq!(p.success_count, 0, "mechanical success undone");
2176 assert_eq!(p.fail_count, 1, "failure booked");
2177 }
2178
2179 #[test]
2180 fn real_success_signal_does_not_double_count() {
2181 let mut tracker = OutcomeTracker::new();
2182 let t = tracker.record_start("m", InferenceTask::Generate, "r");
2183 tracker.record_complete(&t, 100, 8, 15);
2184 tracker.record_inferred_outcome(&t, InferredOutcome::Accepted { confidence: 0.9 });
2185 let p = tracker.profile("m").unwrap();
2186 assert_eq!(
2187 p.success_count, 1,
2188 "Accepted on an already-credited call is not +2"
2189 );
2190 }
2191
2192 #[test]
2193 fn export_import() {
2194 let mut tracker = OutcomeTracker::new();
2195 let trace = tracker.record_start("m1", InferenceTask::Generate, "test");
2196 tracker.record_complete(&trace, 100, 10, 5);
2197 tracker.record_inferred_outcome(&trace, InferredOutcome::Accepted { confidence: 0.9 });
2198
2199 let exported = tracker.export_profiles();
2200 assert_eq!(exported.len(), 1);
2201
2202 let mut tracker2 = OutcomeTracker::new();
2203 tracker2.import_profiles(exported);
2204 assert!(tracker2.profile("m1").is_some());
2205 }
2206}