Skip to main content

icydb_core/db/session/
request.rs

1//! Module: db::session::request
2//! Responsibility: one monotonic aggregate execution scope per request entry.
3//! Does not own: caller authorization, per-execution limits, or physical charging sites.
4//! Boundary: request roots issue shared scope handles that every derived session retains.
5
6use std::{
7    cell::{Cell, RefCell},
8    rc::Rc,
9};
10
11use crate::db::executor::budget::{
12    ExecutionBudgetExceeded, HardExecutionBudget, HardExecutionContext,
13    HardExecutionFailureHeadroom, resource_index,
14};
15#[cfg(feature = "diagnostics")]
16use crate::db::{
17    diagnostics::{
18        RequestDiagnosticResourceUsage, RequestDiagnostics, RequestDiagnosticsState,
19        RequestQueryPlanEvidence,
20    },
21    session::query::QueryPlanCacheAttribution,
22};
23use icydb_diagnostic_code::{DiagnosticExecutionBudgetResource, DiagnosticExecutionBudgetScope};
24
25const REQUEST_FAILURE_HEADROOM: HardExecutionFailureHeadroom =
26    HardExecutionFailureHeadroom::new(500_000_000, 64 * 1_024);
27const REQUEST_HARD_BUDGET: HardExecutionBudget = HardExecutionBudget::new(
28    [
29        256,                 // query executions
30        1_024,               // planning operations
31        256,                 // plan compilations
32        250_000,             // key/index entries visited
33        250_000,             // rows visited
34        128 * 1_024 * 1_024, // stored bytes read
35        16_000_000,          // predicate/expression steps
36        16_000_000,          // nested value steps
37        128 * 1_024 * 1_024, // decoded bytes
38        128 * 1_024 * 1_024, // materialized bytes
39        250_000,             // sort entries
40        32_000_000,          // sort comparisons
41        128 * 1_024 * 1_024, // sort temporary bytes
42        100_000,             // group/distinct entries
43        128 * 1_024 * 1_024, // group/distinct state bytes
44        1_000_000,           // cursor steps
45        128 * 1_024 * 1_024, // temporary bytes
46        1_000_000,           // diagnostic steps
47        100_000,             // result rows
48        64 * 1_024 * 1_024,  // result bytes
49        4_500_000_000,       // instruction units
50    ],
51    REQUEST_FAILURE_HEADROOM,
52);
53
54thread_local! {
55    static CURRENT_REQUEST_SCOPE: RefCell<Option<RequestExecutionScope>> =
56        const { RefCell::new(None) };
57}
58
59/// Non-cloneable capability owning one request's aggregate database counters.
60///
61/// Construct this once at request entry and derive every database session used
62/// by the request from it. Sessions retain the counters, so dropping this
63/// value does not reset work already attached to a derived session.
64pub struct RequestExecutionRoot {
65    scope: RequestExecutionScope,
66}
67
68impl RequestExecutionRoot {
69    /// Mint the fixed production request profile.
70    ///
71    /// This constructor is runtime wiring for generated and guarded facade
72    /// request entry. It is intentionally not a budget-policy configuration
73    /// surface.
74    #[doc(hidden)]
75    #[must_use]
76    pub fn __new_runtime_root() -> Self {
77        Self::from_budget(REQUEST_HARD_BUDGET)
78    }
79
80    /// Reuse the active synchronous request scope or mint the production root.
81    ///
82    /// This is runtime wiring for the public scoped-entry helper. Re-entering
83    /// that helper inside one active database segment must retain the existing
84    /// counters instead of creating a budget-reset escape hatch.
85    #[doc(hidden)]
86    #[must_use]
87    pub fn __new_or_current_runtime_root() -> Self {
88        current_request_scope().map_or_else(Self::__new_runtime_root, |scope| Self { scope })
89    }
90
91    /// Make this root current only while one synchronous call tree executes.
92    ///
93    /// The previous scope is restored before this method returns, including
94    /// during host unwinding. The scope is never retained ambiently across an
95    /// async suspension point.
96    #[doc(hidden)]
97    pub fn __with_current_scope<T>(&self, run: impl FnOnce() -> T) -> T {
98        let _guard = CurrentRequestScopeGuard::enter(self.scope());
99        run()
100    }
101
102    /// Whether no root is active or this root owns the active counters.
103    ///
104    /// Generated facade wiring uses this before accepting an explicit root.
105    /// A different active root would reset aggregate accounting inside a
106    /// request and must fail closed.
107    #[doc(hidden)]
108    #[must_use]
109    pub fn __is_compatible_with_current(&self) -> bool {
110        match current_request_scope() {
111            Some(current) => current.same_counters(&self.scope),
112            None => true,
113        }
114    }
115
116    /// Whether this root owns the counters currently installed for this poll.
117    #[doc(hidden)]
118    #[must_use]
119    pub fn __is_current(&self) -> bool {
120        current_request_scope().is_some_and(|current| current.same_counters(&self.scope))
121    }
122
123    #[cfg(test)]
124    #[must_use]
125    pub(in crate::db) fn new_for_tests(budget: HardExecutionBudget) -> Self {
126        Self::from_budget(budget)
127    }
128
129    fn from_budget(budget: HardExecutionBudget) -> Self {
130        Self {
131            scope: RequestExecutionScope {
132                counters: Rc::new(RequestExecutionCounters {
133                    budget,
134                    observed: [const { Cell::new(0) };
135                        DiagnosticExecutionBudgetResource::ALL.len()],
136                    #[cfg(feature = "diagnostics")]
137                    diagnostics: RefCell::new(None),
138                }),
139            },
140        }
141    }
142
143    pub(in crate::db) fn scope(&self) -> RequestExecutionScope {
144        self.scope.clone()
145    }
146
147    #[cfg(test)]
148    #[must_use]
149    pub(in crate::db) fn observed(&self, resource: DiagnosticExecutionBudgetResource) -> u64 {
150        self.scope.observed(resource)
151    }
152}
153
154pub(in crate::db) fn current_request_scope() -> Option<RequestExecutionScope> {
155    CURRENT_REQUEST_SCOPE.with(|current| current.borrow().clone())
156}
157
158struct CurrentRequestScopeGuard {
159    previous: Option<RequestExecutionScope>,
160}
161
162impl CurrentRequestScopeGuard {
163    fn enter(scope: RequestExecutionScope) -> Self {
164        let previous = CURRENT_REQUEST_SCOPE.with(|current| current.replace(Some(scope)));
165        Self { previous }
166    }
167}
168
169impl Drop for CurrentRequestScopeGuard {
170    fn drop(&mut self) {
171        CURRENT_REQUEST_SCOPE.with(|current| {
172            current.replace(self.previous.take());
173        });
174    }
175}
176
177/// Shared internal handle retained by every session derived from one root.
178#[derive(Clone)]
179pub(in crate::db) struct RequestExecutionScope {
180    counters: Rc<RequestExecutionCounters>,
181}
182
183impl RequestExecutionScope {
184    fn same_counters(&self, other: &Self) -> bool {
185        Rc::ptr_eq(&self.counters, &other.counters)
186    }
187
188    pub(in crate::db) fn charge(
189        &self,
190        context: HardExecutionContext,
191        resource: DiagnosticExecutionBudgetResource,
192        amount: u64,
193    ) -> Result<(), ExecutionBudgetExceeded> {
194        self.counters.charge(context, resource, amount)
195    }
196
197    /// Return how many complete equal-cost units remain across every named
198    /// request resource without mutating request accounting.
199    pub(in crate::db) fn remaining_budget_units(
200        &self,
201        per_unit: &[(DiagnosticExecutionBudgetResource, u64)],
202    ) -> u64 {
203        per_unit
204            .iter()
205            .filter(|(_resource, amount)| *amount != 0)
206            .map(|(resource, amount)| {
207                let index = resource_index(*resource);
208                self.counters
209                    .budget
210                    .limit(*resource)
211                    .saturating_sub(self.counters.observed[index].get())
212                    / amount
213            })
214            .min()
215            .unwrap_or(u64::MAX)
216    }
217
218    /// Preflight one fixed resource bundle without retaining a rejected or
219    /// speculative charge in the request scope.
220    pub(in crate::db) fn can_charge_budget_bundle(
221        &self,
222        charges: &[(DiagnosticExecutionBudgetResource, u64)],
223    ) -> bool {
224        charges.iter().all(|(resource, amount)| {
225            let index = resource_index(*resource);
226            self.counters.observed[index]
227                .get()
228                .checked_add(*amount)
229                .is_some_and(|observed| observed <= self.counters.budget.limit(*resource))
230        })
231    }
232
233    /// Commit one complete bundle only when every request resource fits.
234    /// `false` leaves every request counter unchanged.
235    pub(in crate::db) fn try_commit_budget_bundle(
236        &self,
237        charges: &[(DiagnosticExecutionBudgetResource, u64)],
238    ) -> bool {
239        if !self.can_charge_budget_bundle(charges) {
240            return false;
241        }
242        for (resource, amount) in charges {
243            let counter = &self.counters.observed[resource_index(*resource)];
244            counter.set(counter.get().saturating_add(*amount));
245        }
246
247        true
248    }
249
250    #[cfg(feature = "diagnostics")]
251    pub(in crate::db) fn enable_diagnostics(&self) -> bool {
252        let mut diagnostics = self.counters.diagnostics.borrow_mut();
253        if diagnostics.is_some() {
254            return false;
255        }
256        *diagnostics = Some(RequestDiagnosticsState::default());
257        true
258    }
259
260    #[cfg(feature = "diagnostics")]
261    pub(in crate::db) fn diagnostics_enabled(&self) -> bool {
262        self.counters.diagnostics.borrow().is_some()
263    }
264
265    #[cfg(feature = "diagnostics")]
266    pub(in crate::db) fn diagnostics_snapshot(&self) -> Option<RequestDiagnostics> {
267        let mut snapshot = self
268            .counters
269            .diagnostics
270            .borrow()
271            .as_ref()
272            .map(RequestDiagnosticsState::snapshot)?;
273        let response_bytes = request_diagnostics_bytes_estimate(&snapshot);
274        let context = HardExecutionContext::new(
275            DiagnosticExecutionBudgetScope::Request,
276            icydb_diagnostic_code::DiagnosticExecutionLane::TrustedRead,
277            0,
278        );
279        let charged = self.counters.charge_fail_soft(
280            context,
281            DiagnosticExecutionBudgetResource::DiagnosticSteps,
282            1,
283        ) && self.counters.charge_fail_soft(
284            context,
285            DiagnosticExecutionBudgetResource::ResultBytes,
286            response_bytes,
287        );
288        if !charged {
289            self.suppress_diagnostics(1);
290            snapshot.shapes.clear();
291            snapshot.warnings.clear();
292            snapshot.suppressed_observations = snapshot.suppressed_observations.saturating_add(1);
293        }
294        Some(snapshot)
295    }
296
297    #[cfg(feature = "diagnostics")]
298    pub(in crate::db) fn record_query_plan(
299        &self,
300        evidence: RequestQueryPlanEvidence,
301        cache: QueryPlanCacheAttribution,
302    ) {
303        if !self.diagnostics_enabled() {
304            return;
305        }
306        let context = HardExecutionContext::new(
307            DiagnosticExecutionBudgetScope::Request,
308            icydb_diagnostic_code::DiagnosticExecutionLane::TrustedRead,
309            evidence.normalized_shape_fingerprint_prefix,
310        );
311        let retained_bytes = evidence.retained_bytes_estimate();
312        let diagnostic_steps = evidence.work_steps_estimate();
313        if !self.counters.charge_fail_soft(
314            context,
315            DiagnosticExecutionBudgetResource::DiagnosticSteps,
316            diagnostic_steps,
317        ) || !self.counters.charge_fail_soft(
318            context,
319            DiagnosticExecutionBudgetResource::TemporaryBytes,
320            retained_bytes,
321        ) {
322            self.suppress_diagnostics(1);
323            return;
324        }
325        if let Some(diagnostics) = self.counters.diagnostics.borrow_mut().as_mut() {
326            diagnostics.observe_plan(evidence, cache.hits, cache.misses);
327        }
328    }
329
330    #[cfg(feature = "diagnostics")]
331    pub(in crate::db) fn record_execution(
332        &self,
333        context: HardExecutionContext,
334        usage: RequestDiagnosticResourceUsage,
335    ) {
336        if !self.diagnostics_enabled() {
337            return;
338        }
339        if !self.counters.charge_fail_soft(
340            context,
341            DiagnosticExecutionBudgetResource::DiagnosticSteps,
342            1,
343        ) {
344            self.suppress_diagnostics(1);
345            return;
346        }
347        if let Some(diagnostics) = self.counters.diagnostics.borrow_mut().as_mut() {
348            diagnostics.observe_execution(context.normalized_shape_fingerprint_prefix(), usage);
349        }
350    }
351
352    #[cfg(feature = "diagnostics")]
353    pub(in crate::db) fn record_exact_key_hashes(
354        &self,
355        context: HardExecutionContext,
356        hashes: &[[u8; 16]],
357    ) {
358        if hashes.is_empty() || !self.diagnostics_enabled() {
359            return;
360        }
361        let steps = u64::try_from(hashes.len()).unwrap_or(u64::MAX);
362        let retained_bytes = steps.saturating_mul(16);
363        if !self.counters.charge_fail_soft(
364            context,
365            DiagnosticExecutionBudgetResource::DiagnosticSteps,
366            steps,
367        ) || !self.counters.charge_fail_soft(
368            context,
369            DiagnosticExecutionBudgetResource::TemporaryBytes,
370            retained_bytes,
371        ) {
372            self.suppress_diagnostics(steps);
373            return;
374        }
375        if let Some(diagnostics) = self.counters.diagnostics.borrow_mut().as_mut() {
376            diagnostics
377                .observe_exact_key_hashes(context.normalized_shape_fingerprint_prefix(), hashes);
378        }
379    }
380
381    #[cfg(feature = "diagnostics")]
382    fn suppress_diagnostics(&self, count: u64) {
383        if let Some(diagnostics) = self.counters.diagnostics.borrow_mut().as_mut() {
384            diagnostics.suppress(count);
385        }
386    }
387
388    #[cfg(test)]
389    fn observed(&self, resource: DiagnosticExecutionBudgetResource) -> u64 {
390        self.counters.observed[resource_index(resource)].get()
391    }
392}
393
394struct RequestExecutionCounters {
395    budget: HardExecutionBudget,
396    observed: [Cell<u64>; DiagnosticExecutionBudgetResource::ALL.len()],
397    #[cfg(feature = "diagnostics")]
398    diagnostics: RefCell<Option<RequestDiagnosticsState>>,
399}
400
401impl RequestExecutionCounters {
402    fn charge(
403        &self,
404        context: HardExecutionContext,
405        resource: DiagnosticExecutionBudgetResource,
406        amount: u64,
407    ) -> Result<(), ExecutionBudgetExceeded> {
408        let index = resource_index(resource);
409        let counter = &self.observed[index];
410        let current = counter.get();
411        let (observed, overflowed) = current.overflowing_add(amount);
412        let observed = if overflowed { u64::MAX } else { observed };
413        counter.set(observed);
414        let limit = self.budget.limit(resource);
415        if overflowed || observed > limit {
416            return Err(ExecutionBudgetExceeded::new(
417                resource,
418                limit,
419                observed,
420                context.with_scope(DiagnosticExecutionBudgetScope::Request),
421            ));
422        }
423
424        Ok(())
425    }
426
427    #[cfg(feature = "diagnostics")]
428    fn charge_fail_soft(
429        &self,
430        _context: HardExecutionContext,
431        resource: DiagnosticExecutionBudgetResource,
432        amount: u64,
433    ) -> bool {
434        let index = resource_index(resource);
435        let counter = &self.observed[index];
436        let current = counter.get();
437        let limit = self.budget.limit(resource);
438        let Some(observed) = current.checked_add(amount) else {
439            counter.set(limit);
440            return false;
441        };
442        if observed > limit {
443            counter.set(limit);
444            return false;
445        }
446        counter.set(observed);
447        true
448    }
449}
450
451#[cfg(feature = "diagnostics")]
452fn request_diagnostics_bytes_estimate(diagnostics: &RequestDiagnostics) -> u64 {
453    let shape_bytes = diagnostics.shapes.iter().fold(0_u64, |total, shape| {
454        total
455            .saturating_add(u64::try_from(shape.entity.len()).unwrap_or(u64::MAX))
456            .saturating_add(
457                u64::try_from(shape.selected_index.as_ref().map_or(0, String::len))
458                    .unwrap_or(u64::MAX),
459            )
460            .saturating_add(
461                shape
462                    .residual_fields
463                    .iter()
464                    .chain(shape.compound_index_candidate.iter())
465                    .fold(0_u64, |bytes, field| {
466                        bytes.saturating_add(u64::try_from(field.len()).unwrap_or(u64::MAX))
467                    }),
468            )
469            .saturating_add(256)
470    });
471    diagnostics
472        .warnings
473        .iter()
474        .fold(shape_bytes, |total, warning| {
475            total
476                .saturating_add(u64::try_from(warning.message.len()).unwrap_or(u64::MAX))
477                .saturating_add(32)
478        })
479}
480
481#[cfg(test)]
482mod tests {
483    use super::*;
484
485    #[test]
486    fn synchronous_scope_is_installed_then_removed() {
487        assert!(current_request_scope().is_none());
488        let root = RequestExecutionRoot::new_for_tests(REQUEST_HARD_BUDGET);
489
490        root.__with_current_scope(|| {
491            assert!(current_request_scope().is_some());
492        });
493
494        assert!(current_request_scope().is_none());
495    }
496
497    #[test]
498    fn nested_entry_reuses_current_counters() {
499        let resource = DiagnosticExecutionBudgetResource::QueryExecutions;
500        let budget = REQUEST_HARD_BUDGET.with_limit_for_tests(resource, 1);
501        let root = RequestExecutionRoot::new_for_tests(budget);
502        let context = HardExecutionContext::new(
503            DiagnosticExecutionBudgetScope::Execution,
504            icydb_diagnostic_code::DiagnosticExecutionLane::PublicRead,
505            0,
506        );
507
508        root.__with_current_scope(|| {
509            let nested = RequestExecutionRoot::__new_or_current_runtime_root();
510            nested
511                .scope()
512                .charge(context, resource, 1)
513                .expect("first nested charge should fit");
514            let exhausted = root
515                .scope()
516                .charge(context, resource, 1)
517                .expect_err("parent should observe the nested charge");
518
519            assert_eq!(exhausted.scope(), DiagnosticExecutionBudgetScope::Request);
520            assert_eq!(exhausted.observed(), 2);
521        });
522    }
523
524    #[test]
525    fn budget_bundle_preflight_is_non_mutating_and_uses_remaining_capacity() {
526        let entries = DiagnosticExecutionBudgetResource::GroupDistinctEntries;
527        let bytes = DiagnosticExecutionBudgetResource::GroupDistinctStateBytes;
528        let budget = REQUEST_HARD_BUDGET
529            .with_limit_for_tests(entries, 5)
530            .with_limit_for_tests(bytes, 60);
531        let root = RequestExecutionRoot::new_for_tests(budget);
532        let scope = root.scope();
533        let context = HardExecutionContext::new(
534            DiagnosticExecutionBudgetScope::Execution,
535            icydb_diagnostic_code::DiagnosticExecutionLane::PublicRead,
536            0,
537        );
538        scope
539            .charge(context, entries, 2)
540            .expect("initial request charge should fit");
541        scope
542            .charge(context, bytes, 20)
543            .expect("initial request charge should fit");
544
545        assert_eq!(
546            scope.remaining_budget_units(&[(entries, 1), (bytes, 10)]),
547            3
548        );
549        assert!(scope.can_charge_budget_bundle(&[(entries, 3), (bytes, 30)]));
550        assert!(!scope.can_charge_budget_bundle(&[(entries, 4), (bytes, 30)]));
551        assert_eq!(
552            root.observed(entries),
553            2,
554            "preflight must not charge entries"
555        );
556        assert_eq!(root.observed(bytes), 20, "preflight must not charge bytes");
557        assert!(!scope.try_commit_budget_bundle(&[(entries, 4), (bytes, 30)]));
558        assert_eq!(root.observed(entries), 2);
559        assert_eq!(root.observed(bytes), 20);
560        assert!(scope.try_commit_budget_bundle(&[(entries, 3), (bytes, 30)]));
561        assert_eq!(root.observed(entries), 5);
562        assert_eq!(root.observed(bytes), 50);
563    }
564
565    #[test]
566    fn explicit_root_compatibility_rejects_a_different_active_root() {
567        let first = RequestExecutionRoot::new_for_tests(REQUEST_HARD_BUDGET);
568        let second = RequestExecutionRoot::new_for_tests(REQUEST_HARD_BUDGET);
569
570        assert!(first.__is_compatible_with_current());
571        assert!(!first.__is_current());
572        first.__with_current_scope(|| {
573            assert!(first.__is_current());
574            assert!(first.__is_compatible_with_current());
575            assert!(!second.__is_compatible_with_current());
576        });
577        assert!(second.__is_compatible_with_current());
578    }
579
580    #[cfg(feature = "diagnostics")]
581    #[test]
582    fn request_diagnostic_work_is_charged_to_the_shared_root() {
583        let root = RequestExecutionRoot::new_for_tests(REQUEST_HARD_BUDGET);
584        let scope = root.scope();
585        assert!(scope.enable_diagnostics());
586        scope.record_query_plan(
587            RequestQueryPlanEvidence::bounded(
588                12,
589                "Token",
590                crate::db::RequestDiagnosticAccessPath::ByKey,
591                None,
592                Vec::new(),
593                Vec::new(),
594                vec![[1; 16]],
595            ),
596            QueryPlanCacheAttribution {
597                hits: 1,
598                ..QueryPlanCacheAttribution::default()
599            },
600        );
601
602        assert_eq!(
603            root.observed(DiagnosticExecutionBudgetResource::DiagnosticSteps),
604            2,
605        );
606        assert!(root.observed(DiagnosticExecutionBudgetResource::TemporaryBytes) >= 21);
607    }
608
609    #[cfg(feature = "diagnostics")]
610    #[test]
611    fn exhausted_diagnostic_allowance_suppresses_detail_without_an_error() {
612        let budget = REQUEST_HARD_BUDGET
613            .with_limit_for_tests(DiagnosticExecutionBudgetResource::DiagnosticSteps, 0);
614        let root = RequestExecutionRoot::new_for_tests(budget);
615        let scope = root.scope();
616        assert!(scope.enable_diagnostics());
617        scope.record_query_plan(
618            RequestQueryPlanEvidence::bounded(
619                13,
620                "Token",
621                crate::db::RequestDiagnosticAccessPath::ByKey,
622                None,
623                Vec::new(),
624                Vec::new(),
625                Vec::new(),
626            ),
627            QueryPlanCacheAttribution::default(),
628        );
629
630        let snapshot = scope
631            .diagnostics_snapshot()
632            .expect("enabled diagnostics should still return a bounded snapshot");
633        assert!(snapshot.shapes.is_empty());
634        assert!(snapshot.suppressed_observations >= 2);
635    }
636}