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};
15use icydb_diagnostic_code::{DiagnosticExecutionBudgetResource, DiagnosticExecutionBudgetScope};
16
17const REQUEST_FAILURE_HEADROOM: HardExecutionFailureHeadroom =
18    HardExecutionFailureHeadroom::new(500_000_000, 64 * 1_024);
19const REQUEST_HARD_BUDGET: HardExecutionBudget = HardExecutionBudget::new(
20    [
21        256,                 // query executions
22        1_024,               // planning operations
23        256,                 // plan compilations
24        250_000,             // key/index entries visited
25        250_000,             // rows visited
26        128 * 1_024 * 1_024, // stored bytes read
27        16_000_000,          // predicate/expression steps
28        16_000_000,          // nested value steps
29        128 * 1_024 * 1_024, // decoded bytes
30        128 * 1_024 * 1_024, // materialized bytes
31        250_000,             // sort entries
32        32_000_000,          // sort comparisons
33        128 * 1_024 * 1_024, // sort temporary bytes
34        100_000,             // group/distinct entries
35        128 * 1_024 * 1_024, // group/distinct state bytes
36        1_000_000,           // cursor steps
37        128 * 1_024 * 1_024, // temporary bytes
38        100_000,             // result rows
39        64 * 1_024 * 1_024,  // result bytes
40        4_500_000_000,       // instruction units
41    ],
42    REQUEST_FAILURE_HEADROOM,
43);
44
45thread_local! {
46    static CURRENT_REQUEST_SCOPE: RefCell<Option<RequestExecutionScope>> =
47        const { RefCell::new(None) };
48}
49
50/// Non-cloneable capability owning one request's aggregate database counters.
51///
52/// Construct this once at request entry and derive every database session used
53/// by the request from it. Sessions retain the counters, so dropping this
54/// value does not reset work already attached to a derived session.
55pub struct RequestExecutionRoot {
56    scope: RequestExecutionScope,
57}
58
59impl RequestExecutionRoot {
60    /// Mint the fixed production request profile.
61    ///
62    /// This constructor is runtime wiring for generated and guarded facade
63    /// request entry. It is intentionally not a budget-policy configuration
64    /// surface.
65    #[doc(hidden)]
66    #[must_use]
67    pub fn __new_runtime_root() -> Self {
68        Self::from_budget(REQUEST_HARD_BUDGET)
69    }
70
71    /// Reuse the active synchronous request scope or mint the production root.
72    ///
73    /// This is runtime wiring for the public scoped-entry helper. Re-entering
74    /// that helper inside one active database segment must retain the existing
75    /// counters instead of creating a budget-reset escape hatch.
76    #[doc(hidden)]
77    #[must_use]
78    pub fn __new_or_current_runtime_root() -> Self {
79        current_request_scope().map_or_else(Self::__new_runtime_root, |scope| Self { scope })
80    }
81
82    /// Make this root current only while one synchronous call tree executes.
83    ///
84    /// The previous scope is restored before this method returns, including
85    /// during host unwinding. The scope is never retained ambiently across an
86    /// async suspension point.
87    #[doc(hidden)]
88    pub fn __with_current_scope<T>(&self, run: impl FnOnce() -> T) -> T {
89        let _guard = CurrentRequestScopeGuard::enter(self.scope());
90        run()
91    }
92
93    /// Whether no root is active or this root owns the active counters.
94    ///
95    /// Generated facade wiring uses this before accepting an explicit root.
96    /// A different active root would reset aggregate accounting inside a
97    /// request and must fail closed.
98    #[doc(hidden)]
99    #[must_use]
100    pub fn __is_compatible_with_current(&self) -> bool {
101        match current_request_scope() {
102            Some(current) => current.same_counters(&self.scope),
103            None => true,
104        }
105    }
106
107    /// Whether this root owns the counters currently installed for this poll.
108    #[doc(hidden)]
109    #[must_use]
110    pub fn __is_current(&self) -> bool {
111        current_request_scope().is_some_and(|current| current.same_counters(&self.scope))
112    }
113
114    #[cfg(test)]
115    #[must_use]
116    pub(in crate::db) fn new_for_tests(budget: HardExecutionBudget) -> Self {
117        Self::from_budget(budget)
118    }
119
120    fn from_budget(budget: HardExecutionBudget) -> Self {
121        Self {
122            scope: RequestExecutionScope {
123                counters: Rc::new(RequestExecutionCounters {
124                    budget,
125                    observed: [const { Cell::new(0) };
126                        DiagnosticExecutionBudgetResource::ALL.len()],
127                }),
128            },
129        }
130    }
131
132    pub(in crate::db) fn scope(&self) -> RequestExecutionScope {
133        self.scope.clone()
134    }
135
136    #[cfg(test)]
137    #[must_use]
138    pub(in crate::db) fn observed(&self, resource: DiagnosticExecutionBudgetResource) -> u64 {
139        self.scope.observed(resource)
140    }
141}
142
143pub(in crate::db) fn current_request_scope() -> Option<RequestExecutionScope> {
144    CURRENT_REQUEST_SCOPE.with(|current| current.borrow().clone())
145}
146
147struct CurrentRequestScopeGuard {
148    previous: Option<RequestExecutionScope>,
149}
150
151impl CurrentRequestScopeGuard {
152    fn enter(scope: RequestExecutionScope) -> Self {
153        let previous = CURRENT_REQUEST_SCOPE.with(|current| current.replace(Some(scope)));
154        Self { previous }
155    }
156}
157
158impl Drop for CurrentRequestScopeGuard {
159    fn drop(&mut self) {
160        CURRENT_REQUEST_SCOPE.with(|current| {
161            current.replace(self.previous.take());
162        });
163    }
164}
165
166/// Shared internal handle retained by every session derived from one root.
167#[derive(Clone)]
168pub(in crate::db) struct RequestExecutionScope {
169    counters: Rc<RequestExecutionCounters>,
170}
171
172impl RequestExecutionScope {
173    fn same_counters(&self, other: &Self) -> bool {
174        Rc::ptr_eq(&self.counters, &other.counters)
175    }
176
177    pub(in crate::db) fn charge(
178        &self,
179        context: HardExecutionContext,
180        resource: DiagnosticExecutionBudgetResource,
181        amount: u64,
182    ) -> Result<(), ExecutionBudgetExceeded> {
183        self.counters.charge(context, resource, amount)
184    }
185
186    /// Return how many complete equal-cost units remain across every named
187    /// request resource without mutating request accounting.
188    pub(in crate::db) fn remaining_budget_units(
189        &self,
190        per_unit: &[(DiagnosticExecutionBudgetResource, u64)],
191    ) -> u64 {
192        per_unit
193            .iter()
194            .filter(|(_resource, amount)| *amount != 0)
195            .map(|(resource, amount)| {
196                let index = resource_index(*resource);
197                self.counters
198                    .budget
199                    .limit(*resource)
200                    .saturating_sub(self.counters.observed[index].get())
201                    / amount
202            })
203            .min()
204            .unwrap_or(u64::MAX)
205    }
206
207    /// Preflight one fixed resource bundle without retaining a rejected or
208    /// speculative charge in the request scope.
209    pub(in crate::db) fn can_charge_budget_bundle(
210        &self,
211        charges: &[(DiagnosticExecutionBudgetResource, u64)],
212    ) -> bool {
213        charges.iter().all(|(resource, amount)| {
214            let index = resource_index(*resource);
215            self.counters.observed[index]
216                .get()
217                .checked_add(*amount)
218                .is_some_and(|observed| observed <= self.counters.budget.limit(*resource))
219        })
220    }
221
222    /// Commit one complete bundle only when every request resource fits.
223    /// `false` leaves every request counter unchanged.
224    pub(in crate::db) fn try_commit_budget_bundle(
225        &self,
226        charges: &[(DiagnosticExecutionBudgetResource, u64)],
227    ) -> bool {
228        if !self.can_charge_budget_bundle(charges) {
229            return false;
230        }
231        for (resource, amount) in charges {
232            let counter = &self.counters.observed[resource_index(*resource)];
233            counter.set(counter.get().saturating_add(*amount));
234        }
235
236        true
237    }
238
239    #[cfg(test)]
240    fn observed(&self, resource: DiagnosticExecutionBudgetResource) -> u64 {
241        self.counters.observed[resource_index(resource)].get()
242    }
243}
244
245struct RequestExecutionCounters {
246    budget: HardExecutionBudget,
247    observed: [Cell<u64>; DiagnosticExecutionBudgetResource::ALL.len()],
248}
249
250impl RequestExecutionCounters {
251    fn charge(
252        &self,
253        context: HardExecutionContext,
254        resource: DiagnosticExecutionBudgetResource,
255        amount: u64,
256    ) -> Result<(), ExecutionBudgetExceeded> {
257        let index = resource_index(resource);
258        let counter = &self.observed[index];
259        let current = counter.get();
260        let (observed, overflowed) = current.overflowing_add(amount);
261        let observed = if overflowed { u64::MAX } else { observed };
262        counter.set(observed);
263        let limit = self.budget.limit(resource);
264        if overflowed || observed > limit {
265            return Err(ExecutionBudgetExceeded::new(
266                resource,
267                limit,
268                observed,
269                context.with_scope(DiagnosticExecutionBudgetScope::Request),
270            ));
271        }
272
273        Ok(())
274    }
275}
276
277#[cfg(test)]
278mod tests {
279    use super::*;
280
281    #[test]
282    fn synchronous_scope_is_installed_then_removed() {
283        assert!(current_request_scope().is_none());
284        let root = RequestExecutionRoot::new_for_tests(REQUEST_HARD_BUDGET);
285
286        root.__with_current_scope(|| {
287            assert!(current_request_scope().is_some());
288        });
289
290        assert!(current_request_scope().is_none());
291    }
292
293    #[test]
294    fn nested_entry_reuses_current_counters() {
295        let resource = DiagnosticExecutionBudgetResource::QueryExecutions;
296        let budget = REQUEST_HARD_BUDGET.with_limit_for_tests(resource, 1);
297        let root = RequestExecutionRoot::new_for_tests(budget);
298        let context = HardExecutionContext::new(
299            DiagnosticExecutionBudgetScope::Execution,
300            icydb_diagnostic_code::DiagnosticExecutionLane::PublicRead,
301            0,
302        );
303
304        root.__with_current_scope(|| {
305            let nested = RequestExecutionRoot::__new_or_current_runtime_root();
306            nested
307                .scope()
308                .charge(context, resource, 1)
309                .expect("first nested charge should fit");
310            let exhausted = root
311                .scope()
312                .charge(context, resource, 1)
313                .expect_err("parent should observe the nested charge");
314
315            assert_eq!(exhausted.scope(), DiagnosticExecutionBudgetScope::Request);
316            assert_eq!(exhausted.observed(), 2);
317        });
318    }
319
320    #[test]
321    fn budget_bundle_preflight_is_non_mutating_and_uses_remaining_capacity() {
322        let entries = DiagnosticExecutionBudgetResource::GroupDistinctEntries;
323        let bytes = DiagnosticExecutionBudgetResource::GroupDistinctStateBytes;
324        let budget = REQUEST_HARD_BUDGET
325            .with_limit_for_tests(entries, 5)
326            .with_limit_for_tests(bytes, 60);
327        let root = RequestExecutionRoot::new_for_tests(budget);
328        let scope = root.scope();
329        let context = HardExecutionContext::new(
330            DiagnosticExecutionBudgetScope::Execution,
331            icydb_diagnostic_code::DiagnosticExecutionLane::PublicRead,
332            0,
333        );
334        scope
335            .charge(context, entries, 2)
336            .expect("initial request charge should fit");
337        scope
338            .charge(context, bytes, 20)
339            .expect("initial request charge should fit");
340
341        assert_eq!(
342            scope.remaining_budget_units(&[(entries, 1), (bytes, 10)]),
343            3
344        );
345        assert!(scope.can_charge_budget_bundle(&[(entries, 3), (bytes, 30)]));
346        assert!(!scope.can_charge_budget_bundle(&[(entries, 4), (bytes, 30)]));
347        assert_eq!(
348            root.observed(entries),
349            2,
350            "preflight must not charge entries"
351        );
352        assert_eq!(root.observed(bytes), 20, "preflight must not charge bytes");
353        assert!(!scope.try_commit_budget_bundle(&[(entries, 4), (bytes, 30)]));
354        assert_eq!(root.observed(entries), 2);
355        assert_eq!(root.observed(bytes), 20);
356        assert!(scope.try_commit_budget_bundle(&[(entries, 3), (bytes, 30)]));
357        assert_eq!(root.observed(entries), 5);
358        assert_eq!(root.observed(bytes), 50);
359    }
360
361    #[test]
362    fn explicit_root_compatibility_rejects_a_different_active_root() {
363        let first = RequestExecutionRoot::new_for_tests(REQUEST_HARD_BUDGET);
364        let second = RequestExecutionRoot::new_for_tests(REQUEST_HARD_BUDGET);
365
366        assert!(first.__is_compatible_with_current());
367        assert!(!first.__is_current());
368        first.__with_current_scope(|| {
369            assert!(first.__is_current());
370            assert!(first.__is_compatible_with_current());
371            assert!(!second.__is_compatible_with_current());
372        });
373        assert!(second.__is_compatible_with_current());
374    }
375}