1use 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, 1_024, 256, 250_000, 250_000, 128 * 1_024 * 1_024, 16_000_000, 16_000_000, 128 * 1_024 * 1_024, 128 * 1_024 * 1_024, 250_000, 32_000_000, 128 * 1_024 * 1_024, 100_000, 128 * 1_024 * 1_024, 1_000_000, 128 * 1_024 * 1_024, 100_000, 64 * 1_024 * 1_024, 4_500_000_000, ],
42 REQUEST_FAILURE_HEADROOM,
43);
44
45thread_local! {
46 static CURRENT_REQUEST_SCOPE: RefCell<Option<RequestExecutionScope>> =
47 const { RefCell::new(None) };
48}
49
50pub struct RequestExecutionRoot {
56 scope: RequestExecutionScope,
57}
58
59impl RequestExecutionRoot {
60 #[doc(hidden)]
66 #[must_use]
67 pub fn __new_runtime_root() -> Self {
68 Self::from_budget(REQUEST_HARD_BUDGET)
69 }
70
71 #[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 #[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 #[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 #[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#[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 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 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 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}