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};
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, 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, 1_000_000, 100_000, 64 * 1_024 * 1_024, 4_500_000_000, ],
51 REQUEST_FAILURE_HEADROOM,
52);
53
54thread_local! {
55 static CURRENT_REQUEST_SCOPE: RefCell<Option<RequestExecutionScope>> =
56 const { RefCell::new(None) };
57}
58
59pub struct RequestExecutionRoot {
65 scope: RequestExecutionScope,
66}
67
68impl RequestExecutionRoot {
69 #[doc(hidden)]
75 #[must_use]
76 pub fn __new_runtime_root() -> Self {
77 Self::from_budget(REQUEST_HARD_BUDGET)
78 }
79
80 #[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 #[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 #[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 #[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#[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 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 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 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}