lean_ctx/core/context_kernel/
proxy_bridge.rs1use std::sync::{Mutex, MutexGuard, OnceLock};
4
5use super::client_wiring::{self, OptimizationLevel};
6use super::context_broker::BrokerBudget;
7use super::coverage_class::{self, CoverageClass};
8use super::etpao_live::{EtpaoLive, EtpaoSummary, OutcomeMetrics, RequestMetrics};
9use super::identity::{CallerIdentity, IdentityLedger, IdentityLedgerSummary};
10use super::outcome_signal::{self, InferredOutcome, OutcomeSignal};
11use super::types::ReceiptOutcome;
12
13static IDENTITY_LEDGER: OnceLock<Mutex<IdentityLedger>> = OnceLock::new();
14static ETPAO_TRACKER: OnceLock<Mutex<EtpaoLive>> = OnceLock::new();
15
16#[derive(Debug, Clone, Default)]
18pub struct ProxyRequestData {
19 pub headers: Vec<(String, String)>,
21 pub input_tokens: usize,
23 pub output_tokens: usize,
25 pub reasoning_tokens: usize,
27 pub tokens_saved: usize,
29 pub model: Option<String>,
31 pub provider: Option<String>,
33 pub is_retry: bool,
35 pub request_count: usize,
37}
38
39#[derive(Debug, Clone)]
41pub struct ProxyKernelResult {
42 pub identity: CallerIdentity,
44 pub coverage: CoverageClass,
46 pub coverage_label: &'static str,
48 pub is_addressable: bool,
50 pub optimization_level: OptimizationLevel,
52 pub kernel_budget: BrokerBudget,
54 pub outcome_signal: InferredOutcome,
56}
57
58fn identity_ledger() -> &'static Mutex<IdentityLedger> {
59 IDENTITY_LEDGER.get_or_init(|| Mutex::new(IdentityLedger::new()))
60}
61
62fn etpao_tracker() -> &'static Mutex<EtpaoLive> {
63 ETPAO_TRACKER.get_or_init(|| Mutex::new(EtpaoLive::new()))
64}
65
66fn lock_identity_ledger() -> MutexGuard<'static, IdentityLedger> {
67 match identity_ledger().lock() {
68 Ok(ledger) => ledger,
69 Err(poisoned) => poisoned.into_inner(),
70 }
71}
72
73fn lock_etpao_tracker() -> MutexGuard<'static, EtpaoLive> {
74 match etpao_tracker().lock() {
75 Ok(tracker) => tracker,
76 Err(poisoned) => poisoned.into_inner(),
77 }
78}
79
80#[must_use]
82pub fn process_proxy_request(data: &ProxyRequestData) -> ProxyKernelResult {
83 let context = client_wiring::build_request_context(&data.headers, true, false, false);
84 let outcome =
85 outcome_signal::infer_outcome(data.request_count, data.is_retry, data.output_tokens);
86 let accepted = outcome.outcome == ReceiptOutcome::Accepted;
87 let client_id = context.profile.client_id.clone();
88
89 {
90 let mut tracker = lock_etpao_tracker();
91 tracker.record_request(RequestMetrics {
92 input_tokens: data.input_tokens,
93 output_tokens: data.output_tokens,
94 reasoning_tokens: data.reasoning_tokens,
95 schema_tokens: 0,
96 cache_write_tokens: 0,
97 retry_count: usize::from(data.is_retry),
98 client_id: client_id.clone(),
99 coverage_class: context.coverage,
100 });
101 tracker.record_outcome(OutcomeMetrics {
102 accepted,
103 quality_score: outcome.confidence,
104 first_pass: outcome.signal == OutcomeSignal::FirstPass,
105 client_id,
106 });
107 }
108
109 let consumed = data
110 .input_tokens
111 .saturating_add(data.output_tokens)
112 .saturating_add(data.reasoning_tokens);
113 lock_identity_ledger().record(&context.identity, consumed, data.tokens_saved, accepted);
114
115 let coverage = context.coverage;
116 let optimization_level = client_wiring::optimization_level(&context);
117 let kernel_budget = context.broker_budget;
118
119 ProxyKernelResult {
120 identity: context.identity,
121 coverage,
122 coverage_label: coverage_class::coverage_label(coverage),
123 is_addressable: coverage_class::is_addressable(coverage),
124 optimization_level,
125 kernel_budget,
126 outcome_signal: outcome,
127 }
128}
129
130#[must_use]
132pub fn identity_summary() -> IdentityLedgerSummary {
133 lock_identity_ledger().summary()
134}
135
136#[must_use]
138pub fn etpao_summary() -> EtpaoSummary {
139 lock_etpao_tracker().summary()
140}
141
142#[must_use]
144pub fn current_etpao() -> f64 {
145 lock_etpao_tracker().current_etpao()
146}
147
148pub fn reset_state() {
150 *lock_identity_ledger() = IdentityLedger::new();
151 *lock_etpao_tracker() = EtpaoLive::new();
152}
153
154#[cfg(test)]
155mod tests {
156 use std::sync::{Mutex, MutexGuard};
157
158 use super::{
159 ProxyRequestData, current_etpao, etpao_summary, identity_summary, process_proxy_request,
160 reset_state,
161 };
162 use crate::core::context_kernel::client_wiring::OptimizationLevel;
163 use crate::core::context_kernel::coverage_class::CoverageClass;
164 use crate::core::context_kernel::types::ReceiptOutcome;
165
166 static TEST_LOCK: Mutex<()> = Mutex::new(());
167
168 fn isolated_test() -> MutexGuard<'static, ()> {
169 let guard = match TEST_LOCK.lock() {
170 Ok(guard) => guard,
171 Err(poisoned) => poisoned.into_inner(),
172 };
173 reset_state();
174 guard
175 }
176
177 fn request() -> ProxyRequestData {
178 ProxyRequestData {
179 headers: vec![("x-user-id".to_owned(), "bridge-user".to_owned())],
180 input_tokens: 100,
181 output_tokens: 20,
182 request_count: 1,
183 ..ProxyRequestData::default()
184 }
185 }
186
187 #[test]
188 fn process_empty_request() {
189 let _guard = isolated_test();
190 let result = process_proxy_request(&ProxyRequestData::default());
191
192 assert_eq!(result.coverage, CoverageClass::FullInline);
193 assert!(result.kernel_budget.context_tokens > 0);
194 }
195
196 #[test]
197 fn process_records_identity() {
198 let _guard = isolated_test();
199 let _ = process_proxy_request(&request());
200
201 let summary = identity_summary();
202 assert_eq!(summary.total_users, 1);
203 assert_eq!(summary.total_tokens, 120);
204 }
205
206 #[test]
207 fn process_records_etpao() {
208 let _guard = isolated_test();
209 let _ = process_proxy_request(&request());
210
211 assert!(current_etpao() > 0.0);
212 assert_eq!(etpao_summary().accepted_outcomes, 1);
213 }
214
215 #[test]
216 fn coverage_is_full_inline() {
217 let _guard = isolated_test();
218 let result = process_proxy_request(&request());
219
220 assert_eq!(result.coverage, CoverageClass::FullInline);
221 assert_eq!(result.coverage_label, "full_inline");
222 assert!(result.is_addressable);
223 assert_eq!(result.optimization_level, OptimizationLevel::Full);
224 }
225
226 #[test]
227 fn outcome_first_pass() {
228 let _guard = isolated_test();
229 let result = process_proxy_request(&request());
230
231 assert_eq!(result.outcome_signal.outcome, ReceiptOutcome::Accepted);
232 }
233
234 #[test]
235 fn outcome_retry() {
236 let _guard = isolated_test();
237 let mut data = request();
238 data.is_retry = true;
239 data.request_count = 2;
240
241 let result = process_proxy_request(&data);
242 assert_eq!(result.outcome_signal.outcome, ReceiptOutcome::Rejected);
243 }
244
245 #[test]
246 fn reset_clears_state() {
247 let _guard = isolated_test();
248 let _ = process_proxy_request(&request());
249 reset_state();
250
251 assert_eq!(current_etpao(), 0.0);
252 assert_eq!(identity_summary().total_users, 0);
253 }
254}