Skip to main content

lean_ctx/core/context_kernel/
proxy_bridge.rs

1//! Unified bridge between raw proxy requests and context-kernel services.
2
3use 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/// Raw data extracted from a proxy request; proxy-side types remain in the proxy.
17#[derive(Debug, Clone, Default)]
18pub struct ProxyRequestData {
19    /// Request headers used for identity and client-profile detection.
20    pub headers: Vec<(String, String)>,
21    /// Tokens supplied as model input.
22    pub input_tokens: usize,
23    /// Tokens produced by the model.
24    pub output_tokens: usize,
25    /// Tokens consumed by model reasoning.
26    pub reasoning_tokens: usize,
27    /// Tokens avoided through context optimization.
28    pub tokens_saved: usize,
29    /// Requested model name, when available.
30    pub model: Option<String>,
31    /// Upstream provider name, when available.
32    pub provider: Option<String>,
33    /// Whether this request retries an earlier attempt.
34    pub is_retry: bool,
35    /// Attempt number used to infer the request outcome.
36    pub request_count: usize,
37}
38
39/// Result of kernel processing for a single proxy request.
40#[derive(Debug, Clone)]
41pub struct ProxyKernelResult {
42    /// Caller identity resolved from request headers.
43    pub identity: CallerIdentity,
44    /// Effective coverage for the inline proxy path.
45    pub coverage: CoverageClass,
46    /// Stable machine-readable coverage label.
47    pub coverage_label: &'static str,
48    /// Whether the kernel can modify context on this path.
49    pub is_addressable: bool,
50    /// Optimization level available for the request.
51    pub optimization_level: OptimizationLevel,
52    /// Broker-computed token allocation.
53    pub kernel_budget: BrokerBudget,
54    /// Outcome inferred from observable proxy behavior.
55    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/// Resolves kernel context and records identity and ETPAO metrics for one request.
81#[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/// Returns aggregate identity attribution recorded by the proxy bridge.
131#[must_use]
132pub fn identity_summary() -> IdentityLedgerSummary {
133    lock_identity_ledger().summary()
134}
135
136/// Returns aggregate live ETPAO metrics recorded by the proxy bridge.
137#[must_use]
138pub fn etpao_summary() -> EtpaoSummary {
139    lock_etpao_tracker().summary()
140}
141
142/// Returns current effective tokens per accepted outcome.
143#[must_use]
144pub fn current_etpao() -> f64 {
145    lock_etpao_tracker().current_etpao()
146}
147
148/// Clears all process-wide proxy bridge metrics.
149pub 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}