inillucent-base 2.0.2

Checked binary primitives, identifiers, buffers, limits, and the stable error model.
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
//! What one request is allowed to cost, and the one place it is counted.
//!
//! Invariant: **a request that exceeds its budget stops, and the failure names
//! which budget it was.** An MCP server hands a database to an agent, and an
//! agent asks for what it asks for: `SELECT * FROM chunk` against six hundred
//! thousand rows is a reasonable-looking call that materialises gigabytes, and
//! nothing stopped it before this budget existed. `limit_of` in the command
//! surface even turned a negative row limit into zero, which means *unlimited*.
//!
//! ## Why it is a thread-local rather than a parameter
//!
//! The alternative was threading a `&Budget` through every operator, the tree
//! cursors, the sort and the vector search - about forty signatures, most of
//! which would carry it only to hand it on. That is the shape of change that
//! gets a `None` passed at one call site during a later refactor and quietly
//! stops enforcing anything.
//!
//! One connection runs on one thread (`drivers/inillucent-driver-capi`'s header
//! says so, and the engine's own buffer pool assumes it), so a thread-local is
//! exactly the scope a request occupies. [`Guard`] arms it for the length of a
//! call and restores whatever was there before, so a nested call - a virtual
//! table running a query of its own - cannot widen its caller's budget.
//!
//! ## Why cancellation is not a thread-local
//!
//! The flag is an [`Arc<AtomicBool>`] the caller owns, because the whole point
//! of a cancel is that it arrives from somewhere else: another thread, a signal
//! handler, an MCP client's second connection. Everything *else* here is
//! per-request and per-thread; the flag is the one part that is shared, and it
//! is shared deliberately.

use std::cell::RefCell;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use std::time::{Duration, Instant};

use crate::error::{DbError, PrimaryCode};
use crate::DbResult;

/// What ran out.
///
/// Reported by name rather than as one "too big" so that a caller can act on
/// it: a row limit is answered by asking for fewer rows, a deadline by asking
/// for something cheaper, and a cancellation by nothing at all.
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum Exceeded {
    /// More rows were produced than the request was allowed.
    Rows,
    /// More bytes were produced than the request was allowed.
    Bytes,
    /// The request ran longer than it was allowed.
    Time,
    /// Somebody asked for the request to stop.
    Cancelled,
}

impl Exceeded {
    /// Returns the word a structured error carries.
    ///
    /// Stable, because a client branches on it. It is the *budget's* name and
    /// not a sentence, so a translation or a retry rule can key on it.
    pub fn name(self) -> &'static str {
        match self {
            Exceeded::Rows => "rows",
            Exceeded::Bytes => "bytes",
            Exceeded::Time => "time",
            Exceeded::Cancelled => "cancelled",
        }
    }
}

/// What one request may spend.
#[derive(Clone, Debug)]
pub struct Limits {
    /// The most rows an operator may produce, or `None` for no bound.
    pub rows: Option<u64>,
    /// The most bytes of row data it may produce, or `None` for no bound.
    pub bytes: Option<u64>,
    /// How long it may run, or `None` for no bound.
    pub time: Option<Duration>,
}

impl Limits {
    /// Returns limits that bound nothing, which is what a library embedding
    /// the engine gets unless it asks otherwise.
    ///
    /// **Unbounded is the right default for an embedded database and the wrong
    /// one for a served database**, which is why this exists beside
    /// [`Limits::served`] rather than instead of it. An application that has
    /// linked the engine into its own process is not protecting itself from
    /// itself; a server handing an agent a tool is.
    pub fn unbounded() -> Limits {
        Limits {
            rows: None,
            bytes: None,
            time: None,
        }
    }

    /// Returns the budget a served request gets when nobody said otherwise.
    ///
    /// The numbers are chosen to be generous for a question and stingy for a
    /// mistake. Ten million rows is more than any answer an agent reads and far
    /// less than a table scan of a corpus; 256 MiB is the same ceiling
    /// `inillucent-remote` puts on a single protocol message, for the same
    /// reason; sixty seconds is longer than every query in the differential
    /// suite and shorter than a client's patience.
    pub fn served() -> Limits {
        Limits {
            rows: Some(10_000_000),
            bytes: Some(256 * 1024 * 1024),
            time: Some(Duration::from_secs(60)),
        }
    }

    /// Returns these limits with a different deadline.
    ///
    /// @param time - how long a request may run, or `None` for no bound
    pub fn with_time(mut self, time: Option<Duration>) -> Limits {
        self.time = time;
        self
    }

    /// Returns these limits with a different row bound.
    ///
    /// @param rows - the most rows a request may produce
    pub fn with_rows(mut self, rows: Option<u64>) -> Limits {
        self.rows = rows;
        self
    }

    /// Returns these limits with a different byte bound.
    ///
    /// @param bytes - the most bytes a request may produce
    pub fn with_bytes(mut self, bytes: Option<u64>) -> Limits {
        self.bytes = bytes;
        self
    }
}

impl Default for Limits {
    /// Returns unbounded limits.
    fn default() -> Limits {
        Limits::unbounded()
    }
}

/// One request's budget: what it may spend, and what it has spent.
#[derive(Debug)]
struct Spending {
    /// What it may spend.
    limits: Limits,
    /// When it must stop, computed once from `limits.time`.
    deadline: Option<Instant>,
    /// The flag a cancel from another thread sets.
    cancel: Arc<AtomicBool>,
    /// Rows produced so far.
    rows: u64,
    /// Bytes of row data produced so far.
    bytes: u64,
}

thread_local! {
    /// The budget this thread's request is running under, if any.
    static ACTIVE: RefCell<Option<Spending>> = const { RefCell::new(None) };
}

/// Arms a budget for the length of its lifetime.
///
/// Dropping it restores whatever budget was in force before, which is what
/// makes a nested call - a virtual table's own query, a trigger's statement -
/// run inside its caller's budget rather than beside it.
#[derive(Debug)]
pub struct Guard {
    /// What was armed before, restored on drop.
    previous: Option<Spending>,
}

impl Drop for Guard {
    /// Restores the budget that was in force before this one.
    fn drop(&mut self) {
        let previous = self.previous.take();
        ACTIVE.with(|held| {
            *held.borrow_mut() = previous;
        });
    }
}

/// Arms a budget on this thread until the returned guard is dropped.
///
/// @param limits - what the request may spend
/// @param cancel - the flag another thread sets to stop it
pub fn arm(limits: Limits, cancel: Arc<AtomicBool>) -> Guard {
    // A cancel that arrived before the request started belongs to the request
    // that has finished, not to this one. Clearing it here is what makes the
    // flag safe to reuse across calls on one connection.
    cancel.store(false, Ordering::Relaxed);
    arm_as_it_stands(limits, cancel)
}

/// Arms a budget without clearing the cancellation flag first.
///
/// **For a caller that has already decided this request is cancelled
/// (task-1932, H11).** The clear in [`arm`] is right for a surface that arms
/// once per call and cannot know what arrived in between; it is wrong for one
/// that reads its input on a second thread, because there is a window between
/// taking a request off the queue and arming it, and a cancellation that lands
/// inside that window is wiped by the clear. `inillucent-mcp` clears the flag
/// and publishes which request is running under one lock, so by the time this
/// is called the flag means "this request was cancelled" and clearing it would
/// throw that away.
///
/// @param limits - what the request may spend
/// @param cancel - the flag another thread sets to stop it
pub fn arm_as_it_stands(limits: Limits, cancel: Arc<AtomicBool>) -> Guard {
    let spending = Spending {
        // `checked_add` because the crate denies wrapping arithmetic and an
        // `Instant` plus a caller's `Duration` is a caller's number. A window
        // so large it overflows the clock is the same thing as no deadline.
        deadline: limits
            .time
            .and_then(|window| Instant::now().checked_add(window)),
        limits,
        cancel,
        rows: 0,
        bytes: 0,
    };
    let previous = ACTIVE.with(|held| held.borrow_mut().replace(spending));
    Guard { previous }
}

/// Reports whether a budget is armed on this thread.
pub fn armed() -> bool {
    ACTIVE.with(|held| held.borrow().is_some())
}

/// Refuses when the request has been cancelled or has run out of time.
///
/// **Cheap enough to call in a loop**, which is the property that decides where
/// it can go: an atomic load and, when there is a deadline, one `Instant::now`.
/// Callers that run a tight inner loop check every batch rather than every row.
pub fn check() -> DbResult<()> {
    ACTIVE.with(|held| {
        let Some(spending) = held
            .borrow()
            .as_ref()
            .map(|spending| (spending.cancel.load(Ordering::Relaxed), spending.deadline))
        else {
            return Ok(());
        };
        let (cancelled, deadline) = spending;
        if cancelled {
            return Err(exceeded(Exceeded::Cancelled, 0, 0));
        }
        if let Some(deadline) = deadline {
            if Instant::now() >= deadline {
                return Err(exceeded(Exceeded::Time, 0, 0));
            }
        }
        Ok(())
    })
}

/// Counts rows and their bytes against the budget, refusing when either runs out.
///
/// @param rows - how many rows were produced
/// @param bytes - roughly how many bytes they hold
pub fn spend(rows: u64, bytes: u64) -> DbResult<()> {
    ACTIVE.with(|held| {
        let mut borrowed = held.borrow_mut();
        let Some(spending) = borrowed.as_mut() else {
            return Ok(());
        };
        spending.rows = spending.rows.saturating_add(rows);
        spending.bytes = spending.bytes.saturating_add(bytes);
        if let Some(most) = spending.limits.rows {
            if spending.rows > most {
                return Err(exceeded(Exceeded::Rows, spending.rows, most));
            }
        }
        if let Some(most) = spending.limits.bytes {
            if spending.bytes > most {
                return Err(exceeded(Exceeded::Bytes, spending.bytes, most));
            }
        }
        Ok(())
    })
}

/// Counts bytes a statement had to hold on to, refusing when the byte budget
/// runs out.
///
/// **The row count is deliberately not touched (task-1932, H6).** A request's
/// row budget bounds what the caller is handed - `Limits::served()` says 10,000
/// rows, and `Rows::total` has to keep meaning "rows in the answer" for that
/// number to be worth anything. What a statement *materialises* on the way to
/// that answer is a different quantity and is bounded by the byte budget:
/// a hash join's build side, a group table, a `DISTINCT` set, a window's
/// partition buffer and a recursive CTE's accumulated answer are all rows that
/// occupy memory and never reach the caller.
///
/// Until this existed the only `spend` in the engine was `Collect::push`, the
/// result sink. A join whose build side is a hundred million rows and whose
/// output is one row was bounded by nothing at all: the 256 MiB cap counted the
/// one row it handed back.
///
/// The cancellation and deadline checks are folded in because every caller
/// wants both and one thread-local borrow is cheaper than two.
///
/// @param bytes - roughly how many bytes the statement is now holding
pub fn materialise(bytes: u64) -> DbResult<()> {
    ACTIVE.with(|held| {
        let mut borrowed = held.borrow_mut();
        let Some(spending) = borrowed.as_mut() else {
            return Ok(());
        };
        if spending.cancel.load(Ordering::Relaxed) {
            return Err(exceeded(Exceeded::Cancelled, 0, 0));
        }
        if let Some(deadline) = spending.deadline {
            if Instant::now() >= deadline {
                return Err(exceeded(Exceeded::Time, 0, 0));
            }
        }
        spending.bytes = spending.bytes.saturating_add(bytes);
        if let Some(most) = spending.limits.bytes {
            if spending.bytes > most {
                return Err(exceeded(Exceeded::Bytes, spending.bytes, most));
            }
        }
        Ok(())
    })
}

/// Returns what the armed request has spent so far, as `(rows, bytes)`.
///
/// Zero when nothing is armed, which is the same answer a request that has
/// produced nothing gives - and the caller that asks is reporting rather than
/// deciding.
pub fn spent() -> (u64, u64) {
    ACTIVE.with(|held| {
        held.borrow()
            .as_ref()
            .map(|spending| (spending.rows, spending.bytes))
            .unwrap_or((0, 0))
    })
}

/// Builds the failure a spent budget produces.
///
/// **`Interrupt`, and the same code for all four.** A caller distinguishes them
/// by the sentence, which names the budget; what they share is the property a
/// caller has to branch on first, which is that the connection is still usable
/// and the statement did not run. A row limit reported as `TooBig` and a
/// deadline reported as `Interrupt` would make that one question two.
///
/// @param what - which budget ran out
/// @param used - what the request had spent
/// @param most - what it was allowed
fn exceeded(what: Exceeded, used: u64, most: u64) -> DbError {
    let said = match what {
        Exceeded::Cancelled => "this request was cancelled.".to_string(),
        Exceeded::Time => "this request ran past the time it was allowed.".to_string(),
        Exceeded::Rows => format!(
            "this request produced {used} rows, past the {most} it was allowed. Ask for fewer \
             with a WHERE clause or a LIMIT."
        ),
        Exceeded::Bytes => format!(
            "this request produced {used} bytes of row data, past the {most} it was allowed. Ask \
             for fewer columns or fewer rows."
        ),
    };
    DbError::primary(PrimaryCode::Interrupt)
        .with_message(said.clone())
        // The budget's name, machine-readable, so a client can retry a
        // deadline and not retry a cancellation.
        .with_detail(format!("budget={} {said}", what.name()))
}

/// Returns which budget a failure names, when it is a budget failure.
///
/// The reader for the tag `exceeded` writes. It reads the detail rather than
/// the message because the message is prose a person sees and may be reworded.
///
/// @param error - the failure to classify
pub fn exceeded_kind(error: &DbError) -> Option<Exceeded> {
    let detail = error.detail()?;
    let rest = detail.strip_prefix("budget=")?;
    let name = rest.split_whitespace().next()?;
    match name {
        "rows" => Some(Exceeded::Rows),
        "bytes" => Some(Exceeded::Bytes),
        "time" => Some(Exceeded::Time),
        "cancelled" => Some(Exceeded::Cancelled),
        _ => None,
    }
}

#[cfg(test)]
mod tests {
    use super::*;

    /// With nothing armed, nothing is spent and nothing is refused: the engine
    /// is not a sandbox unless somebody asked for one.
    #[test]
    fn an_unarmed_thread_spends_nothing() {
        assert!(!armed());
        assert!(check().is_ok());
        assert!(spend(u64::MAX, u64::MAX).is_ok());
        assert_eq!(spent(), (0, 0));
    }

    /// A row budget refuses once it is past, and names itself.
    #[test]
    fn a_row_budget_refuses_and_says_which_one() {
        let _guard = arm(
            Limits::unbounded().with_rows(Some(10)),
            Arc::new(AtomicBool::new(false)),
        );
        assert!(spend(9, 0).is_ok());
        let refused = spend(2, 0).expect_err("the eleventh row is past ten");
        assert_eq!(exceeded_kind(&refused), Some(Exceeded::Rows));
        assert!(
            refused.message().contains("11 rows"),
            "{}",
            refused.message()
        );
    }

    /// A byte budget refuses independently of the row budget.
    #[test]
    fn a_byte_budget_refuses_on_its_own() {
        let _guard = arm(
            Limits::unbounded().with_bytes(Some(100)),
            Arc::new(AtomicBool::new(false)),
        );
        assert!(spend(1_000_000, 99).is_ok());
        let refused = spend(0, 2).expect_err("101 bytes is past 100");
        assert_eq!(exceeded_kind(&refused), Some(Exceeded::Bytes));
    }

    /// A deadline that has passed refuses.
    #[test]
    fn a_passed_deadline_refuses() {
        let _guard = arm(
            Limits::unbounded().with_time(Some(Duration::from_millis(0))),
            Arc::new(AtomicBool::new(false)),
        );
        std::thread::sleep(Duration::from_millis(2));
        let refused = check().expect_err("the deadline has passed");
        assert_eq!(exceeded_kind(&refused), Some(Exceeded::Time));
    }

    /// A cancel set from another thread stops the request.
    #[test]
    fn a_cancel_from_another_thread_stops_it() {
        let flag = Arc::new(AtomicBool::new(false));
        let _guard = arm(Limits::unbounded(), Arc::clone(&flag));
        assert!(check().is_ok());
        let other = Arc::clone(&flag);
        std::thread::spawn(move || other.store(true, Ordering::Relaxed))
            .join()
            .expect("the setter runs");
        let refused = check().expect_err("a cancelled request stops");
        assert_eq!(exceeded_kind(&refused), Some(Exceeded::Cancelled));
    }

    /// A cancel left over from the previous request does not stop this one.
    #[test]
    fn a_stale_cancel_does_not_stop_the_next_request() {
        let flag = Arc::new(AtomicBool::new(true));
        let _guard = arm(Limits::unbounded(), Arc::clone(&flag));
        assert!(check().is_ok(), "arming clears a flag from the last call");
    }

    /// A nested budget cannot widen the one it is inside.
    ///
    /// The property that makes this safe to arm around a whole statement: a
    /// virtual table that ran a query of its own would otherwise be able to
    /// hand itself an unbounded one.
    #[test]
    fn a_nested_budget_is_restored_when_it_ends() {
        let outer = arm(
            Limits::unbounded().with_rows(Some(5)),
            Arc::new(AtomicBool::new(false)),
        );
        assert!(spend(4, 0).is_ok());
        {
            let _inner = arm(Limits::unbounded(), Arc::new(AtomicBool::new(false)));
            assert!(spend(1_000, 0).is_ok(), "the inner budget is its own");
        }
        let refused = spend(2, 0).expect_err("the outer budget still counts its own four rows");
        assert_eq!(exceeded_kind(&refused), Some(Exceeded::Rows));
        drop(outer);
        assert!(!armed());
    }

    /// A failure that is not a budget failure is not read as one.
    #[test]
    fn an_ordinary_failure_names_no_budget() {
        assert_eq!(
            exceeded_kind(&crate::error::refusal("something else")),
            None
        );
    }
}