#[cfg(test)]
mod maintenance_tests;
use std::cell::{Cell, RefCell};
#[cfg(test)]
use crate::db::QueryError;
use crate::db::executor::{EntityAuthority, RuntimeGroupedRow, SharedPreparedExecutionPlan};
use crate::db::session::RequestExecutionScope;
use crate::{error::InternalError, value::Value};
use icydb_diagnostic_code::{
DiagnosticExecutionBudgetResource, DiagnosticExecutionBudgetScope, DiagnosticExecutionLane,
};
const RESOURCE_COUNT: usize = DiagnosticExecutionBudgetResource::ALL.len();
const INSTRUCTION_WATERMARK_CHARGE_INTERVAL: u16 = 64;
const INSTRUCTION_WATERMARK_LARGE_CHARGE: u64 = 1_024 * 1_024;
const READ_FAILURE_HEADROOM: HardExecutionFailureHeadroom =
HardExecutionFailureHeadroom::new(500_000_000, 64 * 1_024);
pub(in crate::db) const MUTATION_EXECUTION_INSTRUCTION_LIMIT: u64 = 30_000_000_000;
pub(in crate::db) const MUTATION_EXECUTION_INSTRUCTION_FAILURE_RESERVE: u64 = 5_000_000_000;
const MUTATION_FAILURE_HEADROOM: HardExecutionFailureHeadroom =
HardExecutionFailureHeadroom::new(MUTATION_EXECUTION_INSTRUCTION_FAILURE_RESERVE, 64 * 1_024);
static READ_HARD_BUDGET: HardExecutionBudget = HardExecutionBudget::new(
[
1, 2_000_000, 64, 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, ],
READ_FAILURE_HEADROOM,
);
static MUTATION_HARD_BUDGET: HardExecutionBudget = HardExecutionBudget::new(
[
1, 2_000_000, 64, 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, MUTATION_EXECUTION_INSTRUCTION_LIMIT, ],
MUTATION_FAILURE_HEADROOM,
);
pub(in crate::db) const MUTATION_EXECUTION_BUDGET_POLICY_IDENTITY: u32 =
hard_execution_budget_policy_identity(&MUTATION_HARD_BUDGET);
const fn hard_execution_budget_policy_identity(budget: &HardExecutionBudget) -> u32 {
let mut identity = 0x811c_9dc5_u32;
let mut index = 0;
while index < RESOURCE_COUNT {
identity = fold_execution_budget_policy_u64(identity, budget.limits[index]);
index += 1;
}
identity =
fold_execution_budget_policy_u64(identity, budget.failure_headroom.instruction_units);
fold_execution_budget_policy_u64(identity, budget.failure_headroom.response_bytes)
}
const fn fold_execution_budget_policy_u64(mut identity: u32, value: u64) -> u32 {
let bytes = value.to_le_bytes();
identity ^= u32::from_le_bytes([bytes[0], bytes[1], bytes[2], bytes[3]]);
identity = identity.wrapping_mul(0x0100_0193);
identity ^= u32::from_le_bytes([bytes[4], bytes[5], bytes[6], bytes[7]]);
identity = identity.wrapping_mul(0x0100_0193);
identity
}
#[cfg(test)]
pub(in crate::db::executor) const fn read_hard_budget_limit_for_tests(
resource: DiagnosticExecutionBudgetResource,
) -> u64 {
READ_HARD_BUDGET.limit(resource)
}
std::thread_local! {
static ACTIVE_EXECUTION_BUDGET: RefCell<Option<HardExecutionBudgetTracker>> =
const { RefCell::new(None) };
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(in crate::db) struct HardExecutionFailureHeadroom {
instruction_units: u64,
response_bytes: u64,
}
impl HardExecutionFailureHeadroom {
#[must_use]
pub(in crate::db) const fn new(instruction_units: u64, response_bytes: u64) -> Self {
Self {
instruction_units,
response_bytes,
}
}
#[cfg(test)]
#[must_use]
pub(in crate::db) const fn instruction_units(self) -> u64 {
self.instruction_units
}
#[cfg(test)]
#[must_use]
pub(in crate::db) const fn response_bytes(self) -> u64 {
self.response_bytes
}
const fn is_reserved(self) -> bool {
self.instruction_units != 0 && self.response_bytes != 0
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(in crate::db) struct HardExecutionBudget {
limits: [u64; RESOURCE_COUNT],
failure_headroom: HardExecutionFailureHeadroom,
}
impl HardExecutionBudget {
#[must_use]
pub(in crate::db) const fn new(
limits: [u64; RESOURCE_COUNT],
failure_headroom: HardExecutionFailureHeadroom,
) -> Self {
Self {
limits,
failure_headroom,
}
}
#[must_use]
pub(in crate::db) const fn limit(&self, resource: DiagnosticExecutionBudgetResource) -> u64 {
self.limits[resource_index(resource)]
}
pub(in crate::db) fn can_charge_budget_bundle(
&self,
charges: &[(DiagnosticExecutionBudgetResource, u64)],
observed: impl Fn(DiagnosticExecutionBudgetResource) -> u64,
) -> bool {
let Some(totals) = budget_bundle_totals(charges) else {
return false;
};
charges.iter().all(|(resource, _amount)| {
observed(*resource)
.checked_add(totals[resource_index(*resource)])
.is_some_and(|next| next <= self.limit(*resource))
})
}
pub(in crate::db) fn remaining_budget_units(
&self,
per_unit: &[(DiagnosticExecutionBudgetResource, u64)],
observed: impl Fn(DiagnosticExecutionBudgetResource) -> u64,
) -> u64 {
let Some(totals) = budget_bundle_totals(per_unit) else {
return 0;
};
DiagnosticExecutionBudgetResource::ALL
.into_iter()
.zip(totals)
.filter(|(_resource, amount)| *amount != 0)
.map(|(resource, amount)| {
self.limit(resource).saturating_sub(observed(resource)) / amount
})
.min()
.unwrap_or(u64::MAX)
}
#[cfg(test)]
#[must_use]
pub(in crate::db) const fn failure_headroom(&self) -> HardExecutionFailureHeadroom {
self.failure_headroom
}
#[cfg(test)]
#[must_use]
pub(in crate::db) const fn uniform_for_tests(
limit: u64,
failure_headroom: HardExecutionFailureHeadroom,
) -> Self {
Self::new([limit; RESOURCE_COUNT], failure_headroom)
}
#[cfg(test)]
#[must_use]
pub(in crate::db) const fn with_limit_for_tests(
mut self,
resource: DiagnosticExecutionBudgetResource,
limit: u64,
) -> Self {
self.limits[resource_index(resource)] = limit;
self
}
}
fn budget_bundle_totals(
charges: &[(DiagnosticExecutionBudgetResource, u64)],
) -> Option<[u64; RESOURCE_COUNT]> {
let mut totals = [0_u64; RESOURCE_COUNT];
for (resource, amount) in charges {
let total = &mut totals[resource_index(*resource)];
*total = total.checked_add(*amount)?;
}
Some(totals)
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(in crate::db) struct HardExecutionContext {
scope: DiagnosticExecutionBudgetScope,
lane: DiagnosticExecutionLane,
normalized_shape_fingerprint_prefix: u64,
}
impl HardExecutionContext {
#[must_use]
pub(in crate::db) const fn new(
scope: DiagnosticExecutionBudgetScope,
lane: DiagnosticExecutionLane,
normalized_shape_fingerprint_prefix: u64,
) -> Self {
Self {
scope,
lane,
normalized_shape_fingerprint_prefix,
}
}
#[must_use]
pub(in crate::db) const fn with_scope(self, scope: DiagnosticExecutionBudgetScope) -> Self {
Self { scope, ..self }
}
#[must_use]
pub(in crate::db) const fn lane(self) -> DiagnosticExecutionLane {
self.lane
}
}
pub(in crate::db::executor) fn prepared_read_execution_context(
plan: &SharedPreparedExecutionPlan,
lane: DiagnosticExecutionLane,
) -> HardExecutionContext {
HardExecutionContext::new(
DiagnosticExecutionBudgetScope::Execution,
lane,
plan.execution_shape_fingerprint_prefix(),
)
}
pub(in crate::db::executor) fn read_shape_fingerprint_prefix(
authority: &EntityAuthority,
logical: &crate::db::query::plan::AccessPlannedQuery,
) -> Result<u64, crate::error::InternalError> {
let fingerprint = authority.accepted_schema_fingerprint();
let mut prefix = u64::from_be_bytes([
fingerprint[0],
fingerprint[1],
fingerprint[2],
fingerprint[3],
fingerprint[4],
fingerprint[5],
fingerprint[6],
fingerprint[7],
]) ^ authority.entity_tag().value().rotate_left(17);
let scalar = logical.scalar_plan();
prefix ^= u64::from(logical.has_residual_filter_predicate()?).rotate_left(7);
prefix ^= u64::from(scalar.distinct).rotate_left(11);
prefix ^=
usize_as_u64(scalar.order.as_ref().map_or(0, |order| order.fields.len())).rotate_left(23);
prefix ^= usize_as_u64(
logical
.scalar_projection_plan()
.map_or(0, <[crate::db::query::plan::expr::CompiledExpr]>::len),
)
.rotate_left(31);
prefix ^= usize_as_u64(logical.grouped_aggregate_execution_specs().map_or(
0,
<[crate::db::query::plan::GroupedAggregateExecutionSpec]>::len,
))
.rotate_left(41);
Ok(prefix)
}
pub(in crate::db) const fn direct_read_execution_context(
authority: &EntityAuthority,
lane: DiagnosticExecutionLane,
shape_domain: u64,
) -> HardExecutionContext {
let fingerprint = authority.accepted_schema_fingerprint();
let prefix = u64::from_be_bytes([
fingerprint[0],
fingerprint[1],
fingerprint[2],
fingerprint[3],
fingerprint[4],
fingerprint[5],
fingerprint[6],
fingerprint[7],
]) ^ authority.entity_tag().value().rotate_left(17)
^ shape_domain;
HardExecutionContext::new(DiagnosticExecutionBudgetScope::Execution, lane, prefix)
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(in crate::db) struct ExecutionBudgetExceeded {
resource: DiagnosticExecutionBudgetResource,
limit: u64,
observed: u64,
context: HardExecutionContext,
}
impl ExecutionBudgetExceeded {
#[must_use]
pub(in crate::db) const fn new(
resource: DiagnosticExecutionBudgetResource,
limit: u64,
observed: u64,
context: HardExecutionContext,
) -> Self {
Self {
resource,
limit,
observed,
context,
}
}
#[must_use]
pub(in crate::db) const fn resource(self) -> DiagnosticExecutionBudgetResource {
self.resource
}
#[must_use]
pub(in crate::db) const fn limit(self) -> u64 {
self.limit
}
#[must_use]
pub(in crate::db) const fn observed(self) -> u64 {
self.observed
}
#[must_use]
pub(in crate::db) const fn scope(self) -> DiagnosticExecutionBudgetScope {
self.context.scope
}
#[must_use]
pub(in crate::db) const fn lane(self) -> DiagnosticExecutionLane {
self.context.lane
}
#[must_use]
pub(in crate::db) const fn normalized_shape_fingerprint_prefix(self) -> u64 {
self.context.normalized_shape_fingerprint_prefix
}
}
impl From<ExecutionBudgetExceeded> for InternalError {
fn from(exhausted: ExecutionBudgetExceeded) -> Self {
Self::execution_budget_exceeded(
exhausted.resource(),
exhausted.limit(),
exhausted.observed(),
exhausted.scope(),
exhausted.lane(),
exhausted.normalized_shape_fingerprint_prefix(),
)
}
}
enum HardExecutionBudgetAuthority {
Static(&'static HardExecutionBudget),
#[cfg(test)]
TestOwned(Box<HardExecutionBudget>),
}
impl HardExecutionBudgetAuthority {
const fn budget(&self) -> &HardExecutionBudget {
match self {
Self::Static(budget) => budget,
#[cfg(test)]
Self::TestOwned(budget) => budget,
}
}
}
pub(in crate::db) struct HardExecutionBudgetTracker {
budget: HardExecutionBudgetAuthority,
context: HardExecutionContext,
request_scope: Option<RequestExecutionScope>,
observed: [u64; RESOURCE_COUNT],
last_instruction_counter: Option<u64>,
charges_since_instruction_watermark: u16,
}
impl HardExecutionBudgetTracker {
#[must_use]
pub(in crate::db) fn new(
budget: &'static HardExecutionBudget,
context: HardExecutionContext,
) -> Self {
debug_assert!(budget.failure_headroom.is_reserved());
Self {
budget: HardExecutionBudgetAuthority::Static(budget),
context,
request_scope: None,
observed: [0; RESOURCE_COUNT],
last_instruction_counter: None,
charges_since_instruction_watermark: 0,
}
}
#[must_use]
pub(in crate::db) fn new_with_request_scope(
budget: &'static HardExecutionBudget,
context: HardExecutionContext,
request_scope: &RequestExecutionScope,
) -> Self {
let mut tracker = Self::new(budget, context);
tracker.request_scope = Some(request_scope.clone());
tracker
}
#[cfg(test)]
#[must_use]
pub(in crate::db) fn new_for_tests(
budget: HardExecutionBudget,
context: HardExecutionContext,
) -> Self {
debug_assert!(budget.failure_headroom.is_reserved());
Self {
budget: HardExecutionBudgetAuthority::TestOwned(Box::new(budget)),
context,
request_scope: None,
observed: [0; RESOURCE_COUNT],
last_instruction_counter: None,
charges_since_instruction_watermark: 0,
}
}
pub(in crate::db) fn precharge(
&mut self,
resource: DiagnosticExecutionBudgetResource,
amount: u64,
) -> Result<(), ExecutionBudgetExceeded> {
self.charge_raw(resource, amount)
}
pub(in crate::db) fn charge_periodic(
&mut self,
resource: DiagnosticExecutionBudgetResource,
amount: u64,
) -> Result<(), ExecutionBudgetExceeded> {
self.charge(resource, amount)
}
#[cfg(test)]
#[must_use]
pub(in crate::db) const fn observed(&self, resource: DiagnosticExecutionBudgetResource) -> u64 {
self.observed[resource_index(resource)]
}
pub(in crate::db) fn check_instruction_watermark(
&mut self,
) -> Result<(), ExecutionBudgetExceeded> {
let current = crate::runtime::local_instruction_counter();
let delta = self
.last_instruction_counter
.map_or(0, |previous| current.saturating_sub(previous));
self.last_instruction_counter = Some(current);
self.charges_since_instruction_watermark = 0;
self.charge_raw(DiagnosticExecutionBudgetResource::InstructionUnits, delta)
}
pub(in crate::db) fn finish_instruction_watermark(
&mut self,
) -> Result<(), ExecutionBudgetExceeded> {
if self.last_instruction_counter.is_none() {
return Ok(());
}
self.check_instruction_watermark()
}
fn remaining_budget_units(&self, per_unit: &[(DiagnosticExecutionBudgetResource, u64)]) -> u64 {
let execution_remaining = self
.budget
.budget()
.remaining_budget_units(per_unit, |resource| self.observed[resource_index(resource)]);
let request_remaining = self
.request_scope
.as_ref()
.map_or(u64::MAX, |scope| scope.remaining_budget_units(per_unit));
execution_remaining.min(request_remaining)
}
fn can_charge_budget_bundle(
&self,
charges: &[(DiagnosticExecutionBudgetResource, u64)],
) -> bool {
self.budget
.budget()
.can_charge_budget_bundle(charges, |resource| self.observed[resource_index(resource)])
&& self
.request_scope
.as_ref()
.is_none_or(|scope| scope.can_charge_budget_bundle(charges))
}
fn try_charge_budget_bundle(
&mut self,
charges: &[(DiagnosticExecutionBudgetResource, u64)],
) -> Result<bool, ExecutionBudgetExceeded> {
if !self.can_charge_budget_bundle(charges) {
return Ok(false);
}
self.check_instruction_watermark()?;
if !self.can_charge_budget_bundle(charges) {
return Ok(false);
}
if let Some(scope) = self.request_scope.as_ref()
&& !scope.try_commit_budget_bundle(charges)
{
return Ok(false);
}
for (resource, amount) in charges {
let observed = &mut self.observed[resource_index(*resource)];
*observed = observed.saturating_add(*amount);
}
self.charges_since_instruction_watermark = self
.charges_since_instruction_watermark
.saturating_add(u16::try_from(charges.len()).unwrap_or(u16::MAX));
Ok(true)
}
#[cfg(test)]
#[must_use]
pub(in crate::db) const fn failure_headroom(&self) -> HardExecutionFailureHeadroom {
self.budget.budget().failure_headroom()
}
fn charge(
&mut self,
resource: DiagnosticExecutionBudgetResource,
amount: u64,
) -> Result<(), ExecutionBudgetExceeded> {
if !matches!(
resource,
DiagnosticExecutionBudgetResource::InstructionUnits
) {
self.charges_since_instruction_watermark =
self.charges_since_instruction_watermark.saturating_add(1);
if self.charges_since_instruction_watermark >= INSTRUCTION_WATERMARK_CHARGE_INTERVAL
|| amount >= INSTRUCTION_WATERMARK_LARGE_CHARGE
{
self.check_instruction_watermark()?;
}
}
self.charge_raw(resource, amount)
}
fn charge_raw(
&mut self,
resource: DiagnosticExecutionBudgetResource,
amount: u64,
) -> Result<(), ExecutionBudgetExceeded> {
let index = resource_index(resource);
let current = self.observed[index];
let (observed, overflowed) = current.overflowing_add(amount);
let observed = if overflowed { u64::MAX } else { observed };
self.observed[index] = observed;
let limit = self.budget.budget().limit(resource);
let execution_result = if overflowed || observed > limit {
Err(ExecutionBudgetExceeded::new(
resource,
limit,
observed,
self.context,
))
} else {
Ok(())
};
let request_result = self
.request_scope
.as_ref()
.map_or(Ok(()), |scope| scope.charge(self.context, resource, amount));
execution_result?;
request_result
}
}
pub(in crate::db::executor) fn with_read_execution_budget<T>(
request_scope: &RequestExecutionScope,
context: HardExecutionContext,
run: impl FnOnce() -> Result<T, InternalError>,
) -> Result<T, InternalError> {
with_execution_budget(
HardExecutionBudgetTracker::new_with_request_scope(
&READ_HARD_BUDGET,
context,
request_scope,
),
run,
std::convert::identity,
ExecutionBudgetFinish::Automatic,
)
}
pub(in crate::db) fn with_mutation_execution_budget<T, E>(
context: HardExecutionContext,
run: impl FnOnce() -> Result<T, E>,
map_internal: fn(InternalError) -> E,
) -> Result<T, E> {
let mut tracker = HardExecutionBudgetTracker::new(&MUTATION_HARD_BUDGET, context);
tracker
.check_instruction_watermark()
.map_err(InternalError::from)
.map_err(map_internal)?;
with_execution_budget(tracker, run, map_internal, ExecutionBudgetFinish::Explicit)
}
pub(in crate::db) fn charge_current_execution_budget(
resource: DiagnosticExecutionBudgetResource,
amount: u64,
) -> Result<(), InternalError> {
if amount == 0 {
return Ok(());
}
ACTIVE_EXECUTION_BUDGET.with(|budget| {
let mut budget = budget
.try_borrow_mut()
.map_err(|_| InternalError::query_executor_invariant())?;
let Some(budget) = budget.as_mut() else {
return Ok(());
};
budget
.charge_periodic(resource, amount)
.map_err(InternalError::from)
})
}
pub(in crate::db) struct MaintenanceConstructionBudget {
tracker: RefCell<HardExecutionBudgetTracker>,
exhausted: Cell<Option<ExecutionBudgetExceeded>>,
}
impl MaintenanceConstructionBudget {
#[cfg(test)]
pub(in crate::db) fn with_limit_for_tests(
resource: DiagnosticExecutionBudgetResource,
limit: u64,
) -> Self {
Self {
exhausted: Cell::new(None),
tracker: RefCell::new(HardExecutionBudgetTracker::new_for_tests(
MUTATION_HARD_BUDGET.with_limit_for_tests(resource, limit),
HardExecutionContext::new(
DiagnosticExecutionBudgetScope::Execution,
DiagnosticExecutionLane::Mutation,
0,
),
)),
}
}
pub(in crate::db) fn new() -> Self {
Self {
exhausted: Cell::new(None),
tracker: RefCell::new(HardExecutionBudgetTracker::new(
&MUTATION_HARD_BUDGET,
HardExecutionContext::new(
DiagnosticExecutionBudgetScope::Execution,
DiagnosticExecutionLane::Mutation,
0,
),
)),
}
}
pub(in crate::db) fn run<T, E>(
&self,
run: impl FnOnce(&Self) -> Result<T, E>,
map_error: fn(InternalError) -> E,
) -> Result<T, E> {
self.with_tracker(HardExecutionBudgetTracker::check_instruction_watermark)
.map_err(map_error)?;
let result = run(self);
self.with_tracker(HardExecutionBudgetTracker::check_instruction_watermark)
.map_err(map_error)?;
result
}
fn with_tracker(
&self,
charge: impl FnOnce(&mut HardExecutionBudgetTracker) -> Result<(), ExecutionBudgetExceeded>,
) -> Result<(), InternalError> {
if let Some(error) = self.exhausted.get() {
return Err(error.into());
}
let mut tracker = self
.tracker
.try_borrow_mut()
.map_err(|_| InternalError::query_executor_invariant())?;
charge(&mut tracker).map_err(|error| {
self.exhausted.set(Some(error));
InternalError::from(error)
})
}
}
impl crate::db::query::construction::ConstructionBudget for MaintenanceConstructionBudget {
fn charge(
&self,
resource: DiagnosticExecutionBudgetResource,
amount: u64,
) -> Result<(), InternalError> {
self.with_tracker(|tracker| tracker.charge_periodic(resource, amount))
}
}
pub(in crate::db) struct ExecutionConstructionBudget;
impl crate::db::query::construction::ConstructionBudget for ExecutionConstructionBudget {
fn charge(
&self,
resource: DiagnosticExecutionBudgetResource,
amount: u64,
) -> Result<(), InternalError> {
ACTIVE_EXECUTION_BUDGET.with(|budget| {
let mut budget = budget
.try_borrow_mut()
.map_err(|_| InternalError::query_executor_invariant())?;
let budget = budget
.as_mut()
.ok_or_else(InternalError::query_executor_invariant)?;
budget
.charge_periodic(resource, amount)
.map_err(InternalError::from)
})
}
}
pub(in crate::db) fn current_execution_remaining_budget_units(
per_unit: &[(DiagnosticExecutionBudgetResource, u64)],
) -> Result<u64, InternalError> {
ACTIVE_EXECUTION_BUDGET.with(|budget| {
let budget = budget
.try_borrow()
.map_err(|_| InternalError::query_executor_invariant())?;
let budget = budget
.as_ref()
.ok_or_else(InternalError::query_executor_invariant)?;
Ok(budget.remaining_budget_units(per_unit))
})
}
pub(in crate::db) fn try_charge_current_execution_budget_bundle(
charges: &[(DiagnosticExecutionBudgetResource, u64)],
) -> Result<bool, InternalError> {
ACTIVE_EXECUTION_BUDGET.with(|budget| {
let mut budget = budget
.try_borrow_mut()
.map_err(|_| InternalError::query_executor_invariant())?;
let budget = budget
.as_mut()
.ok_or_else(InternalError::query_executor_invariant)?;
budget
.try_charge_budget_bundle(charges)
.map_err(InternalError::from)
})
}
pub(in crate::db::executor) fn charge_current_execution_budget_pair(
first: (DiagnosticExecutionBudgetResource, u64),
second: (DiagnosticExecutionBudgetResource, u64),
) -> Result<(), InternalError> {
if first.1 == 0 && second.1 == 0 {
return Ok(());
}
ACTIVE_EXECUTION_BUDGET.with(|budget| {
let mut budget = budget
.try_borrow_mut()
.map_err(|_| InternalError::query_executor_invariant())?;
let Some(budget) = budget.as_mut() else {
return Ok(());
};
if first.1 != 0 {
budget
.charge_periodic(first.0, first.1)
.map_err(InternalError::from)?;
}
if second.1 != 0 {
budget
.charge_periodic(second.0, second.1)
.map_err(InternalError::from)?;
}
Ok(())
})
}
pub(in crate::db::executor) fn charge_decoded_row(
row_bytes: usize,
nested_steps: usize,
) -> Result<(), InternalError> {
charge_current_execution_budget_pair(
(
DiagnosticExecutionBudgetResource::DecodedBytes,
usize_as_u64(row_bytes),
),
(
DiagnosticExecutionBudgetResource::NestedValueSteps,
usize_as_u64(nested_steps),
),
)
}
macro_rules! charge_materialized_data_row {
($row:expr) => {{
let row_bytes = u64::try_from($row.len()).unwrap_or(u64::MAX);
charge_current_execution_budget(
DiagnosticExecutionBudgetResource::StoredBytesRead,
row_bytes,
)?;
charge_current_execution_budget(
DiagnosticExecutionBudgetResource::MaterializedBytes,
row_bytes,
)
}};
}
pub(in crate::db::executor) use charge_materialized_data_row;
pub(in crate::db::executor) fn current_execution_budget_exceeded(
resource: DiagnosticExecutionBudgetResource,
limit: u64,
observed: u64,
) -> InternalError {
ACTIVE_EXECUTION_BUDGET.with(|budget| {
let Ok(budget) = budget.try_borrow() else {
return InternalError::query_executor_invariant();
};
let Some(budget) = budget.as_ref() else {
return InternalError::query_executor_invariant();
};
ExecutionBudgetExceeded::new(resource, limit, observed, budget.context).into()
})
}
#[derive(Clone, Copy)]
pub(in crate::db::executor) struct ExecutionBudgetUsage {
observed: [u64; RESOURCE_COUNT],
}
impl ExecutionBudgetUsage {
#[must_use]
pub(in crate::db::executor) const fn observed(
self,
resource: DiagnosticExecutionBudgetResource,
) -> u64 {
self.observed[resource_index(resource)]
}
}
pub(in crate::db::executor) fn current_execution_budget_usage()
-> Result<ExecutionBudgetUsage, InternalError> {
ACTIVE_EXECUTION_BUDGET.with(|budget| {
let budget = budget
.try_borrow()
.map_err(|_| InternalError::query_executor_invariant())?;
let budget = budget
.as_ref()
.ok_or_else(InternalError::query_executor_invariant)?;
Ok(ExecutionBudgetUsage {
observed: budget.observed,
})
})
}
pub(in crate::db::executor) fn charge_sort_work<R>(entries: usize) -> Result<(), InternalError> {
let comparisons_per_entry = if entries <= 1 {
0
} else {
usize::BITS.saturating_sub(entries.saturating_sub(1).leading_zeros())
};
let comparisons = entries.saturating_mul(comparisons_per_entry as usize);
charge_current_execution_budget(
DiagnosticExecutionBudgetResource::SortEntries,
usize_as_u64(entries),
)?;
charge_current_execution_budget(
DiagnosticExecutionBudgetResource::SortComparisons,
usize_as_u64(comparisons),
)?;
charge_current_execution_budget(
DiagnosticExecutionBudgetResource::SortTemporaryBytes,
usize_as_u64(entries.saturating_mul(std::mem::size_of::<R>())),
)
}
pub(in crate::db) fn finish_current_execution_instruction_watermark() -> Result<(), InternalError> {
ACTIVE_EXECUTION_BUDGET.with(|budget| {
let mut budget = budget
.try_borrow_mut()
.map_err(|_| InternalError::query_executor_invariant())?;
let Some(budget) = budget.as_mut() else {
return Ok(());
};
budget
.finish_instruction_watermark()
.map_err(InternalError::from)
})
}
pub(in crate::db::executor) fn charge_runtime_value_rows(
rows: &[Vec<Value>],
) -> Result<(), InternalError> {
let (bytes, nested_steps) = rows.iter().flatten().fold((0_u64, 0_u64), |total, value| {
let value_work = runtime_value_work(value);
(
total.0.saturating_add(value_work.0),
total.1.saturating_add(value_work.1),
)
});
charge_current_execution_budget(
DiagnosticExecutionBudgetResource::ResultRows,
usize_as_u64(rows.len()),
)?;
charge_current_execution_budget(
DiagnosticExecutionBudgetResource::NestedValueSteps,
nested_steps,
)?;
charge_current_execution_budget(DiagnosticExecutionBudgetResource::MaterializedBytes, bytes)?;
charge_current_execution_budget(DiagnosticExecutionBudgetResource::ResultBytes, bytes)
}
pub(in crate::db::executor) fn charge_runtime_grouped_rows(
rows: &[RuntimeGroupedRow],
) -> Result<(), InternalError> {
let (bytes, nested_steps) = rows
.iter()
.flat_map(|row| row.group_key().iter().chain(row.aggregate_values()))
.fold((0_u64, 0_u64), |total, value| {
let value_work = runtime_value_work(value);
(
total.0.saturating_add(value_work.0),
total.1.saturating_add(value_work.1),
)
});
charge_current_execution_budget(
DiagnosticExecutionBudgetResource::ResultRows,
usize_as_u64(rows.len()),
)?;
charge_current_execution_budget(
DiagnosticExecutionBudgetResource::NestedValueSteps,
nested_steps,
)?;
charge_current_execution_budget(DiagnosticExecutionBudgetResource::MaterializedBytes, bytes)?;
charge_current_execution_budget(DiagnosticExecutionBudgetResource::ResultBytes, bytes)
}
pub(in crate::db::executor) const RUNTIME_VALUE_NODE_OVERHEAD_BYTES: u64 = 32;
pub(in crate::db::executor) fn runtime_value_work(value: &Value) -> (u64, u64) {
const VALUE_OVERHEAD: u64 = RUNTIME_VALUE_NODE_OVERHEAD_BYTES;
match value {
Value::Blob(value) => (VALUE_OVERHEAD.saturating_add(usize_as_u64(value.len())), 1),
Value::Text(value) => (VALUE_OVERHEAD.saturating_add(usize_as_u64(value.len())), 1),
Value::IntBig(value) => (VALUE_OVERHEAD.saturating_add(value.leb128_len()), 1),
Value::NatBig(value) => (VALUE_OVERHEAD.saturating_add(value.leb128_len()), 1),
Value::Principal(value) => (
VALUE_OVERHEAD.saturating_add(usize_as_u64(value.as_slice().len())),
1,
),
Value::List(values) => values.iter().fold((VALUE_OVERHEAD, 1_u64), |total, value| {
let value_work = runtime_value_work(value);
(
total.0.saturating_add(value_work.0),
total.1.saturating_add(value_work.1),
)
}),
Value::Map(entries) => {
entries
.iter()
.fold((VALUE_OVERHEAD, 1_u64), |total, (key, value)| {
let key_work = runtime_value_work(key);
let value_work = runtime_value_work(value);
(
total
.0
.saturating_add(key_work.0)
.saturating_add(value_work.0),
total
.1
.saturating_add(key_work.1)
.saturating_add(value_work.1),
)
})
}
Value::Enum(value) => value.payload().map_or((VALUE_OVERHEAD, 1), |payload| {
let payload_work = runtime_value_work(payload);
(
VALUE_OVERHEAD.saturating_add(payload_work.0),
1_u64.saturating_add(payload_work.1),
)
}),
Value::Account(_)
| Value::Bool(_)
| Value::Date(_)
| Value::Decimal(_)
| Value::Duration(_)
| Value::Float32(_)
| Value::Float64(_)
| Value::Int64(_)
| Value::Int128(_)
| Value::Null
| Value::Subaccount(_)
| Value::Timestamp(_)
| Value::Nat64(_)
| Value::Nat128(_)
| Value::Ulid(_)
| Value::Unit => (VALUE_OVERHEAD, 1),
Value::U256(_) => (VALUE_OVERHEAD, 32),
}
}
#[derive(Clone, Copy)]
enum ExecutionBudgetFinish {
Automatic,
Explicit,
}
fn with_execution_budget<T, E>(
mut tracker: HardExecutionBudgetTracker,
run: impl FnOnce() -> Result<T, E>,
map_internal: fn(InternalError) -> E,
finish: ExecutionBudgetFinish,
) -> Result<T, E> {
tracker
.precharge(DiagnosticExecutionBudgetResource::QueryExecutions, 1)
.map_err(InternalError::from)
.map_err(map_internal)?;
let installed = ACTIVE_EXECUTION_BUDGET
.with(|budget| {
let mut budget = budget
.try_borrow_mut()
.map_err(|_| InternalError::query_executor_invariant())?;
if budget.is_some() {
return Ok(false);
}
*budget = Some(tracker);
Ok::<bool, InternalError>(true)
})
.map_err(map_internal)?;
if !installed {
return run();
}
let result = run();
let final_budget_result = match finish {
ExecutionBudgetFinish::Automatic => finish_current_execution_instruction_watermark(),
ExecutionBudgetFinish::Explicit => Ok(()),
};
let removed = ACTIVE_EXECUTION_BUDGET.with(|budget| {
budget
.try_borrow_mut()
.map_err(|_| InternalError::query_executor_invariant())?
.take()
.ok_or_else(InternalError::query_executor_invariant)
});
let removed = removed.map_err(map_internal)?;
let _ = removed;
final_budget_result.map_err(map_internal)?;
result
}
fn usize_as_u64(value: usize) -> u64 {
u64::try_from(value).unwrap_or(u64::MAX)
}
#[cfg(test)]
pub(in crate::db) fn with_query_execution_budget_for_tests<T>(
budget: HardExecutionBudget,
context: HardExecutionContext,
run: impl FnOnce() -> Result<T, QueryError>,
) -> Result<T, QueryError> {
with_execution_budget(
HardExecutionBudgetTracker::new_for_tests(budget, context),
run,
QueryError::execute,
ExecutionBudgetFinish::Automatic,
)
}
#[cfg(all(test, feature = "sql"))]
pub(in crate::db) fn with_execution_budget_for_tests<T, E>(
budget: HardExecutionBudget,
context: HardExecutionContext,
run: impl FnOnce() -> Result<T, E>,
map_internal: fn(InternalError) -> E,
) -> Result<T, E> {
with_execution_budget(
HardExecutionBudgetTracker::new_for_tests(budget, context),
run,
map_internal,
ExecutionBudgetFinish::Automatic,
)
}
pub(in crate::db) const fn resource_index(resource: DiagnosticExecutionBudgetResource) -> usize {
match resource {
DiagnosticExecutionBudgetResource::QueryExecutions => 0,
DiagnosticExecutionBudgetResource::PlanningSteps => 1,
DiagnosticExecutionBudgetResource::PlanCompilations => 2,
DiagnosticExecutionBudgetResource::KeyIndexEntriesVisited => 3,
DiagnosticExecutionBudgetResource::RowsVisited => 4,
DiagnosticExecutionBudgetResource::StoredBytesRead => 5,
DiagnosticExecutionBudgetResource::PredicateExpressionSteps => 6,
DiagnosticExecutionBudgetResource::NestedValueSteps => 7,
DiagnosticExecutionBudgetResource::DecodedBytes => 8,
DiagnosticExecutionBudgetResource::MaterializedBytes => 9,
DiagnosticExecutionBudgetResource::SortEntries => 10,
DiagnosticExecutionBudgetResource::SortComparisons => 11,
DiagnosticExecutionBudgetResource::SortTemporaryBytes => 12,
DiagnosticExecutionBudgetResource::GroupDistinctEntries => 13,
DiagnosticExecutionBudgetResource::GroupDistinctStateBytes => 14,
DiagnosticExecutionBudgetResource::CursorSteps => 15,
DiagnosticExecutionBudgetResource::TemporaryBytes => 16,
DiagnosticExecutionBudgetResource::ResultRows => 17,
DiagnosticExecutionBudgetResource::ResultBytes => 18,
DiagnosticExecutionBudgetResource::InstructionUnits => 19,
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::db::{RequestExecutionRoot, data::RawRow};
use icydb_diagnostic_code::{DiagnosticDetail, DiagnosticFactTag, RuntimeBoundaryCode};
const TEST_HEADROOM: HardExecutionFailureHeadroom = HardExecutionFailureHeadroom::new(500, 256);
#[test]
fn bigint_runtime_work_preserves_encoded_size_and_nested_node_charges() {
for integer in [-8193_i32, -8192, -65, -64, -1, 0, 63, 64, 8191, 8192] {
let signed = crate::types::IntBig::from(integer);
let unsigned = crate::types::NatBig::from(integer.unsigned_abs());
let payload_bytes = (signed.to_leb128().len() + unsigned.to_leb128().len()) as u64;
let value = Value::List(vec![Value::IntBig(signed), Value::NatBig(unsigned)]);
assert_eq!(
runtime_value_work(&value),
(3 * RUNTIME_VALUE_NODE_OVERHEAD_BYTES + payload_bytes, 3)
);
}
}
#[test]
fn construction_requires_active_authority_and_shares_execution_request_counters() {
use crate::db::query::construction::ConstructionBudget;
let resource = DiagnosticExecutionBudgetResource::TemporaryBytes;
let construction = &ExecutionConstructionBudget as &dyn ConstructionBudget;
assert!(construction.charge(resource, 0).is_err());
let root = RequestExecutionRoot::new_for_tests(
HardExecutionBudget::uniform_for_tests(u64::MAX, TEST_HEADROOM)
.with_limit_for_tests(resource, 3),
);
with_execution_budget(
HardExecutionBudgetTracker::new_with_request_scope(
&BUNDLE_EXECUTION_BUDGET,
TEST_CONTEXT,
&root.scope(),
),
|| {
construction.copy_text("abc")?;
assert_eq!(current_execution_budget_usage()?.observed(resource), 3);
assert_eq!(root.observed(resource), 3);
let err = construction.copy_text("x").unwrap_err();
assert!(
err.diagnostic_facts()
.contains(&(DiagnosticFactTag::BudgetResource, resource.raw()))
);
Ok::<_, InternalError>(())
},
std::convert::identity,
ExecutionBudgetFinish::Automatic,
)
.unwrap();
assert!(construction.copy_text("x").is_err());
}
const TEST_CONTEXT: HardExecutionContext = HardExecutionContext::new(
DiagnosticExecutionBudgetScope::Execution,
DiagnosticExecutionLane::PublicRead,
0x0102_0304_0506_0708,
);
static PAIR_FIRST_FAILURE_BUDGET: HardExecutionBudget =
HardExecutionBudget::uniform_for_tests(u64::MAX, TEST_HEADROOM)
.with_limit_for_tests(DiagnosticExecutionBudgetResource::KeyIndexEntriesVisited, 0);
static BUNDLE_EXECUTION_BUDGET: HardExecutionBudget =
HardExecutionBudget::uniform_for_tests(u64::MAX, TEST_HEADROOM)
.with_limit_for_tests(DiagnosticExecutionBudgetResource::GroupDistinctEntries, 5)
.with_limit_for_tests(
DiagnosticExecutionBudgetResource::GroupDistinctStateBytes,
50,
);
#[test]
fn every_resource_charges_monotonically_and_retains_rejected_work() {
let budget = HardExecutionBudget::new([1; RESOURCE_COUNT], TEST_HEADROOM);
for resource in DiagnosticExecutionBudgetResource::ALL {
let mut tracker = HardExecutionBudgetTracker::new_for_tests(budget, TEST_CONTEXT);
tracker
.precharge(resource, 1)
.expect("work at the hard ceiling should be admitted");
let exhausted = tracker
.charge_periodic(resource, 1)
.expect_err("work above the hard ceiling should reject");
assert_eq!(exhausted.resource(), resource);
assert_eq!(exhausted.limit(), 1);
assert_eq!(exhausted.observed(), 2);
assert_eq!(tracker.observed(resource), 2);
}
}
#[test]
fn mutation_execution_policy_identity_covers_limits_and_failure_reserves() {
assert_eq!(MUTATION_EXECUTION_BUDGET_POLICY_IDENTITY, 0x35cc_94d9);
assert_eq!(
MUTATION_HARD_BUDGET.limit(DiagnosticExecutionBudgetResource::InstructionUnits),
MUTATION_EXECUTION_INSTRUCTION_LIMIT,
);
assert_eq!(
MUTATION_HARD_BUDGET.failure_headroom().instruction_units(),
MUTATION_EXECUTION_INSTRUCTION_FAILURE_RESERVE,
);
for resource in DiagnosticExecutionBudgetResource::ALL {
let changed = MUTATION_HARD_BUDGET.with_limit_for_tests(
resource,
MUTATION_HARD_BUDGET.limit(resource).wrapping_add(1),
);
assert_ne!(
hard_execution_budget_policy_identity(&changed),
MUTATION_EXECUTION_BUDGET_POLICY_IDENTITY,
"resource {resource:?} must participate in mutation execution policy identity",
);
}
for headroom in [
HardExecutionFailureHeadroom::new(
MUTATION_EXECUTION_INSTRUCTION_FAILURE_RESERVE + 1,
MUTATION_FAILURE_HEADROOM.response_bytes,
),
HardExecutionFailureHeadroom::new(
MUTATION_EXECUTION_INSTRUCTION_FAILURE_RESERVE,
MUTATION_FAILURE_HEADROOM.response_bytes + 1,
),
] {
let changed = HardExecutionBudget::new(MUTATION_HARD_BUDGET.limits, headroom);
assert_ne!(
hard_execution_budget_policy_identity(&changed),
MUTATION_EXECUTION_BUDGET_POLICY_IDENTITY,
);
}
}
#[test]
fn arithmetic_overflow_is_exhaustion_and_never_refunds_usage() {
let budget = HardExecutionBudget::new([u64::MAX; RESOURCE_COUNT], TEST_HEADROOM);
let resource = DiagnosticExecutionBudgetResource::PlanningSteps;
let mut tracker = HardExecutionBudgetTracker::new_for_tests(budget, TEST_CONTEXT);
tracker
.precharge(resource, u64::MAX)
.expect("the representable ceiling should be admitted");
let exhausted = tracker
.charge_periodic(resource, 1)
.expect_err("counter overflow must reject");
assert_eq!(exhausted.observed(), u64::MAX);
assert_eq!(tracker.observed(resource), u64::MAX);
}
#[test]
fn exhaustion_maps_to_complete_typed_diagnostic_facts() {
let budget = HardExecutionBudget::new([0; RESOURCE_COUNT], TEST_HEADROOM);
let mut tracker = HardExecutionBudgetTracker::new_for_tests(budget, TEST_CONTEXT);
let exhausted = tracker
.precharge(DiagnosticExecutionBudgetResource::QueryExecutions, 1)
.expect_err("zero query allowance should reject");
let error = InternalError::from(exhausted);
assert!(matches!(
error.diagnostic().detail(),
Some(DiagnosticDetail::RuntimeBoundary {
boundary: RuntimeBoundaryCode::ExecutionBudgetExceeded,
})
));
assert_eq!(
error.diagnostic_facts(),
vec![
(DiagnosticFactTag::BudgetResource, 1),
(DiagnosticFactTag::Limit, 0),
(DiagnosticFactTag::Actual, 1),
(DiagnosticFactTag::ExecutionBudgetScope, 1),
(DiagnosticFactTag::ExecutionLane, 1),
(
DiagnosticFactTag::QueryShapeFingerprintPrefix,
0x0102_0304_0506_0708,
),
],
);
assert_eq!(tracker.failure_headroom(), TEST_HEADROOM);
assert_eq!(TEST_HEADROOM.instruction_units(), 500);
assert_eq!(TEST_HEADROOM.response_bytes(), 256);
}
#[test]
fn paired_budget_charges_preserve_sequential_failure_order() {
let root = RequestExecutionRoot::new_for_tests(HardExecutionBudget::uniform_for_tests(
u64::MAX,
TEST_HEADROOM,
));
let error = with_execution_budget(
HardExecutionBudgetTracker::new_with_request_scope(
&PAIR_FIRST_FAILURE_BUDGET,
TEST_CONTEXT,
&root.scope(),
),
|| {
charge_current_execution_budget_pair(
(DiagnosticExecutionBudgetResource::KeyIndexEntriesVisited, 1),
(DiagnosticExecutionBudgetResource::CursorSteps, 1),
)
},
std::convert::identity,
ExecutionBudgetFinish::Automatic,
)
.expect_err("the first paired charge should retain its ordinary hard limit");
assert!(matches!(
error.diagnostic().detail(),
Some(DiagnosticDetail::RuntimeBoundary {
boundary: RuntimeBoundaryCode::ExecutionBudgetExceeded,
})
));
assert_eq!(
root.observed(DiagnosticExecutionBudgetResource::KeyIndexEntriesVisited),
1,
);
assert_eq!(
root.observed(DiagnosticExecutionBudgetResource::CursorSteps),
0,
"the second charge must not run after the first fails",
);
}
#[test]
fn decoded_row_charges_preserve_exact_limits_and_failure_accounting() {
use DiagnosticExecutionBudgetResource::{DecodedBytes, NestedValueSteps};
const DECODE_BUDGET: HardExecutionBudget =
HardExecutionBudget::uniform_for_tests(u64::MAX, TEST_HEADROOM)
.with_limit_for_tests(DecodedBytes, 7)
.with_limit_for_tests(NestedValueSteps, 3);
static DECODE_BUDGETS: [HardExecutionBudget; 3] = [
DECODE_BUDGET,
DECODE_BUDGET.with_limit_for_tests(DecodedBytes, 6),
DECODE_BUDGET.with_limit_for_tests(NestedValueSteps, 2),
];
for (budget, rejected_resource, observed_steps) in [
(&DECODE_BUDGETS[0], None, 3),
(&DECODE_BUDGETS[1], Some(DecodedBytes), 0),
(&DECODE_BUDGETS[2], Some(NestedValueSteps), 3),
] {
let root = RequestExecutionRoot::new_for_tests(HardExecutionBudget::uniform_for_tests(
u64::MAX,
TEST_HEADROOM,
));
let result = with_execution_budget(
HardExecutionBudgetTracker::new_with_request_scope(
budget,
TEST_CONTEXT,
&root.scope(),
),
|| charge_decoded_row(7, 3),
std::convert::identity,
ExecutionBudgetFinish::Automatic,
);
if let Some(resource) = rejected_resource {
let error = result.expect_err("decoding above either ceiling must reject");
assert!(
error
.diagnostic_facts()
.contains(&(DiagnosticFactTag::BudgetResource, resource.raw())),
);
} else {
result.expect("decoding at both ceilings must succeed");
}
assert_eq!(root.observed(DecodedBytes), 7);
assert_eq!(root.observed(NestedValueSteps), observed_steps);
}
}
#[test]
fn materialized_data_row_charge_owns_byte_resources() {
let row = RawRow::try_new(vec![0; 7]).expect("bounded test row");
let budget = HardExecutionBudget::uniform_for_tests(u64::MAX, TEST_HEADROOM);
with_execution_budget(
HardExecutionBudgetTracker::new_for_tests(budget, TEST_CONTEXT),
|| {
charge_materialized_data_row!(&row)?;
let usage = current_execution_budget_usage()?;
assert_eq!(
usage.observed(DiagnosticExecutionBudgetResource::StoredBytesRead),
7,
);
assert_eq!(
usage.observed(DiagnosticExecutionBudgetResource::MaterializedBytes),
7,
);
Ok::<_, InternalError>(())
},
std::convert::identity,
ExecutionBudgetFinish::Automatic,
)
.expect("materialized row accounting should complete");
}
#[test]
fn semantic_budget_bundle_is_atomic_across_execution_and_request() {
let entries = DiagnosticExecutionBudgetResource::GroupDistinctEntries;
let bytes = DiagnosticExecutionBudgetResource::GroupDistinctStateBytes;
let request_budget = HardExecutionBudget::uniform_for_tests(u64::MAX, TEST_HEADROOM)
.with_limit_for_tests(entries, 3)
.with_limit_for_tests(bytes, 30);
let root = RequestExecutionRoot::new_for_tests(request_budget);
with_execution_budget(
HardExecutionBudgetTracker::new_with_request_scope(
&BUNDLE_EXECUTION_BUDGET,
TEST_CONTEXT,
&root.scope(),
),
|| {
charge_current_execution_budget(entries, 1)?;
charge_current_execution_budget(bytes, 10)?;
let before = current_execution_budget_usage()?;
assert_eq!(
current_execution_remaining_budget_units(&[(entries, 1), (entries, 1)])?,
1
);
assert!(!try_charge_current_execution_budget_bundle(&[
(entries, 2),
(bytes, 10),
(entries, 1),
])?);
let duplicate_rejected = current_execution_budget_usage()?;
assert_eq!(
duplicate_rejected.observed(entries),
before.observed(entries)
);
assert_eq!(duplicate_rejected.observed(bytes), before.observed(bytes));
assert_eq!(root.observed(entries), 1);
assert_eq!(root.observed(bytes), 10);
assert!(!try_charge_current_execution_budget_bundle(&[
(entries, 3),
(bytes, 10),
])?);
let rejected = current_execution_budget_usage()?;
assert_eq!(rejected.observed(entries), before.observed(entries));
assert_eq!(rejected.observed(bytes), before.observed(bytes));
assert_eq!(root.observed(entries), 1);
assert_eq!(root.observed(bytes), 10);
assert!(try_charge_current_execution_budget_bundle(&[
(entries, 2),
(bytes, 20),
])?);
let committed = current_execution_budget_usage()?;
assert_eq!(committed.observed(entries), 3);
assert_eq!(committed.observed(bytes), 30);
assert_eq!(root.observed(entries), 3);
assert_eq!(root.observed(bytes), 30);
Ok::<_, InternalError>(())
},
std::convert::identity,
ExecutionBudgetFinish::Automatic,
)
.expect("atomic bundle proof should complete");
}
#[test]
fn repeated_budget_bundle_resources_are_atomic_for_each_limiting_scope() {
let entries = DiagnosticExecutionBudgetResource::GroupDistinctEntries;
let bytes = DiagnosticExecutionBudgetResource::GroupDistinctStateBytes;
for (execution_limit, request_limit) in [(10, 100), (100, 10), (10, 10)] {
let root = RequestExecutionRoot::new_for_tests(
HardExecutionBudget::uniform_for_tests(u64::MAX, TEST_HEADROOM)
.with_limit_for_tests(entries, request_limit),
);
let mut tracker = HardExecutionBudgetTracker::new_for_tests(
HardExecutionBudget::uniform_for_tests(u64::MAX, TEST_HEADROOM)
.with_limit_for_tests(entries, execution_limit),
TEST_CONTEXT,
);
tracker.request_scope = Some(root.scope());
tracker.precharge(entries, 2).unwrap();
tracker.precharge(bytes, 1).unwrap();
let before_execution = tracker.observed;
let before_request = root.request_budget();
for bundle in [
[(entries, 5), (bytes, 1), (entries, 4)],
[(entries, 4), (entries, 5), (bytes, 1)],
] {
assert!(!tracker.try_charge_budget_bundle(&bundle).unwrap());
assert_eq!(tracker.observed, before_execution);
assert_eq!(root.request_budget(), before_request);
}
assert_eq!(
tracker.remaining_budget_units(&[(entries, 2), (entries, 2)]),
2
);
assert!(
tracker
.try_charge_budget_bundle(&[
(entries, 5),
(bytes, 1),
(entries, 0),
(entries, 3),
])
.unwrap()
);
assert_eq!(tracker.observed(entries), 10);
assert_eq!(root.observed(entries), 10);
assert_eq!(tracker.observed(bytes), 2);
assert_eq!(root.observed(bytes), 2);
let committed_request = root.request_budget();
let committed_execution = tracker.observed;
assert!(
!tracker
.try_charge_budget_bundle(&[(entries, 1), (bytes, 1)])
.unwrap()
);
assert_eq!(root.request_budget(), committed_request);
assert_eq!(tracker.observed, committed_execution);
}
}
#[test]
fn budget_bundle_instruction_exhaustion_preserves_semantic_counters() {
let instructions = DiagnosticExecutionBudgetResource::InstructionUnits;
let entries = DiagnosticExecutionBudgetResource::GroupDistinctEntries;
for (execution_limit, request_limit, expected_scope) in [
(0, u64::MAX, DiagnosticExecutionBudgetScope::Execution),
(u64::MAX, 0, DiagnosticExecutionBudgetScope::Request),
] {
let root = RequestExecutionRoot::new_for_tests(
HardExecutionBudget::uniform_for_tests(u64::MAX, TEST_HEADROOM)
.with_limit_for_tests(instructions, request_limit),
);
let mut tracker = HardExecutionBudgetTracker::new_for_tests(
HardExecutionBudget::uniform_for_tests(u64::MAX, TEST_HEADROOM)
.with_limit_for_tests(instructions, execution_limit),
TEST_CONTEXT,
);
tracker.request_scope = Some(root.scope());
tracker.precharge(instructions, 1).unwrap_err();
let before_request = root.request_budget();
let before_execution = tracker.observed;
let failure = tracker
.try_charge_budget_bundle(&[(entries, 1), (entries, 1)])
.unwrap_err();
assert_eq!(failure.resource(), instructions);
assert_eq!(failure.scope(), expected_scope);
assert_eq!(failure.observed(), 1);
assert_eq!(root.request_budget(), before_request);
assert_eq!(tracker.observed, before_execution);
}
}
#[test]
fn budget_bundle_overflow_is_rejected_without_saturating_into_admission() {
let resource = DiagnosticExecutionBudgetResource::ResultBytes;
let unlimited = HardExecutionBudget::uniform_for_tests(u64::MAX, TEST_HEADROOM);
let root = RequestExecutionRoot::new_for_tests(unlimited);
let mut tracker = HardExecutionBudgetTracker::new_for_tests(unlimited, TEST_CONTEXT);
let scope = root.scope();
let overflow = [(resource, u64::MAX), (resource, 1)];
assert!(!tracker.try_charge_budget_bundle(&overflow).unwrap());
assert_eq!(tracker.observed(resource), 0);
tracker.request_scope = Some(root.scope());
assert!(!scope.can_charge_budget_bundle(&overflow));
assert!(!scope.try_commit_budget_bundle(&overflow));
assert!(!tracker.try_charge_budget_bundle(&overflow).unwrap());
assert_eq!(scope.remaining_budget_units(&overflow), 0);
assert_eq!(tracker.remaining_budget_units(&overflow), 0);
assert_eq!(root.observed(resource), 0);
assert_eq!(tracker.observed(resource), 0);
assert_eq!(scope.remaining_budget_units(&[]), u64::MAX);
assert_eq!(tracker.remaining_budget_units(&[(resource, 0)]), u64::MAX);
assert!(scope.try_commit_budget_bundle(&[]));
assert!(tracker.try_charge_budget_bundle(&[(resource, 0)]).unwrap());
assert!(
tracker
.try_charge_budget_bundle(&[(resource, u64::MAX - 1), (resource, 1)])
.unwrap()
);
let exhausted = root.request_budget();
assert!(!tracker.try_charge_budget_bundle(&[(resource, 1)]).unwrap());
assert_eq!(root.request_budget(), exhausted);
assert_eq!(tracker.observed(resource), u64::MAX);
}
#[test]
fn derived_execution_trackers_cannot_reset_the_shared_request_scope() {
let request_budget = HardExecutionBudget::uniform_for_tests(u64::MAX, TEST_HEADROOM)
.with_limit_for_tests(DiagnosticExecutionBudgetResource::QueryExecutions, 2);
let root = RequestExecutionRoot::new_for_tests(request_budget);
let scope = root.scope();
for _ in 0..2 {
HardExecutionBudgetTracker::new_with_request_scope(
&READ_HARD_BUDGET,
TEST_CONTEXT,
&scope,
)
.precharge(DiagnosticExecutionBudgetResource::QueryExecutions, 1)
.expect("work at the aggregate request ceiling should admit");
}
let exhausted = HardExecutionBudgetTracker::new_with_request_scope(
&READ_HARD_BUDGET,
TEST_CONTEXT,
&root.scope(),
)
.precharge(DiagnosticExecutionBudgetResource::QueryExecutions, 1)
.expect_err("a fresh derived scope handle must not reset request counters");
assert_eq!(exhausted.scope(), DiagnosticExecutionBudgetScope::Request);
assert_eq!(exhausted.observed(), 3);
assert_eq!(
root.observed(DiagnosticExecutionBudgetResource::QueryExecutions),
3,
);
}
#[test]
fn failures_retries_and_nested_executions_remain_charged() {
let request_budget = HardExecutionBudget::uniform_for_tests(u64::MAX, TEST_HEADROOM);
let root = RequestExecutionRoot::new_for_tests(request_budget);
let scope = root.scope();
let failed = with_execution_budget(
HardExecutionBudgetTracker::new_with_request_scope(
&READ_HARD_BUDGET,
TEST_CONTEXT,
&scope,
),
|| Err::<(), _>(InternalError::query_executor_invariant()),
std::convert::identity,
ExecutionBudgetFinish::Automatic,
);
assert!(failed.is_err());
let retried = with_execution_budget(
HardExecutionBudgetTracker::new_with_request_scope(
&READ_HARD_BUDGET,
TEST_CONTEXT,
&scope,
),
|| {
with_execution_budget(
HardExecutionBudgetTracker::new_with_request_scope(
&READ_HARD_BUDGET,
TEST_CONTEXT,
&scope,
),
|| Ok::<_, InternalError>(()),
std::convert::identity,
ExecutionBudgetFinish::Automatic,
)
},
std::convert::identity,
ExecutionBudgetFinish::Automatic,
);
assert!(retried.is_ok());
assert_eq!(
root.observed(DiagnosticExecutionBudgetResource::QueryExecutions),
3,
"the failed attempt, retry, and nested execution all stay charged",
);
}
#[test]
fn planning_and_compilation_charges_share_the_request_scope() {
let request_budget = HardExecutionBudget::uniform_for_tests(u64::MAX, TEST_HEADROOM)
.with_limit_for_tests(DiagnosticExecutionBudgetResource::PlanningSteps, 2)
.with_limit_for_tests(DiagnosticExecutionBudgetResource::PlanCompilations, 1);
let root = RequestExecutionRoot::new_for_tests(request_budget);
let first_scope = root.scope();
first_scope
.charge(
TEST_CONTEXT,
DiagnosticExecutionBudgetResource::PlanningSteps,
1,
)
.expect("the first planning operation should be admitted");
first_scope
.charge(
TEST_CONTEXT,
DiagnosticExecutionBudgetResource::PlanCompilations,
1,
)
.expect("the first compilation should be admitted");
let second_scope = root.scope();
second_scope
.charge(
TEST_CONTEXT,
DiagnosticExecutionBudgetResource::PlanningSteps,
1,
)
.expect("planning at the aggregate ceiling should be admitted");
let exhausted = second_scope
.charge(
TEST_CONTEXT,
DiagnosticExecutionBudgetResource::PlanCompilations,
1,
)
.expect_err("a derived scope must retain the earlier compilation charge");
assert_eq!(exhausted.scope(), DiagnosticExecutionBudgetScope::Request);
assert_eq!(exhausted.limit(), 1);
assert_eq!(exhausted.observed(), 2);
assert_eq!(
root.observed(DiagnosticExecutionBudgetResource::PlanningSteps),
2,
);
assert_eq!(
root.observed(DiagnosticExecutionBudgetResource::PlanCompilations),
2,
);
}
#[test]
fn collection_scale_n_plus_one_work_fails_at_the_aggregate_request_boundary() {
let root = RequestExecutionRoot::__new_runtime_root();
let scope = root.scope();
let mut rejected = None;
for _ in 0..257 {
let charge = HardExecutionBudgetTracker::new_with_request_scope(
&READ_HARD_BUDGET,
HardExecutionContext::new(
DiagnosticExecutionBudgetScope::Execution,
DiagnosticExecutionLane::TrustedRead,
0x746f_6b6f_2d6e_2b31,
),
&scope,
)
.precharge(DiagnosticExecutionBudgetResource::QueryExecutions, 1);
if let Err(exhausted) = charge {
rejected = Some(exhausted);
break;
}
}
let exhausted = rejected.expect("the 257th individually bounded query should reject");
assert_eq!(exhausted.scope(), DiagnosticExecutionBudgetScope::Request);
assert_eq!(exhausted.lane(), DiagnosticExecutionLane::TrustedRead);
assert_eq!(exhausted.limit(), 256);
assert_eq!(exhausted.observed(), 257);
}
}