khive-storage 0.8.0

Storage capability contracts and backend-neutral request context.
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
507
508
509
510
511
512
513
514
515
516
517
518
//! Process-wide open-transaction registry (ADR-091 Plank 0).
//!
//! Every caller-controllable SQL transaction span (`WriterGuard::transaction`,
//! `atomic_unit`'s own registered span, raw `BEGIN IMMEDIATE`/`COMMIT`
//! batch-writer spans, and admitted cached-reader transactions) registers here
//! on open and deregisters via `TxHandle`'s `Drop`. This is observe-only: no
//! enforcement reads the registry in this plank. It exists so the checkpoint
//! task can name which caller, if any, is holding a WAL snapshot open.

use std::collections::HashMap;
use std::ffi::OsString;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{LazyLock, Mutex};
use std::time::{Duration, Instant};

/// Identifier for one registered transaction span.
///
/// Public so consumers of [`oldest`] can detect the oldest entry *changing
/// identity* between observations without a live registration of their own.
/// Equality is the only supported operation; the numeric value carries no
/// meaning beyond "same span" vs. "different span". See
/// `crates/khive-storage/docs/api/tx-registry.md`.
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
pub struct TxId(pub u64);

/// Opaque identity for a file-backed database, keyed on its already-minted
/// canonical path. Minting (resolving aliases to one canonical path) happens
/// at exactly one point — the pool in `khive-db` — every other layer treats
/// this constructor as accepting that output only.
#[derive(Clone, Debug, PartialEq, Eq, Hash)]
pub struct DbIdentity(OsString);

impl DbIdentity {
    pub fn new(canonical_path: impl Into<OsString>) -> Self {
        Self(canonical_path.into())
    }
}

/// Which database, if any, a registered span runs against. Three states,
/// never an option: "no database file" ([`TxOrigin::Memory`]) and "not yet
/// threaded to an origin" ([`TxOrigin::Unscoped`]) are different facts and
/// must never share a representation.
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum TxOrigin {
    /// A file-backed database, identified by its minted [`DbIdentity`].
    Database(DbIdentity),
    /// An in-memory backend: no database file, no WAL, no sidecar —
    /// excluded from WAL-pin attribution entirely.
    Memory,
    /// A registration site not yet threaded to an origin. Observed by the
    /// main backend's attribution view, exactly as every entry was before
    /// origins existed, so scoping can never silently drop a span.
    Unscoped,
}

/// An attribution view over the registry: which origins it observes.
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum TxOriginFilter {
    /// The main backend's view: spans scoped to `main`, plus every
    /// [`TxOrigin::Unscoped`] span (the never-silently-drop fallback).
    Main(DbIdentity),
    /// A secondary backend's view: spans scoped to exactly this database.
    /// Never falls back to unscoped spans — only the main view does.
    Secondary(DbIdentity),
}

impl TxOriginFilter {
    fn matches(&self, origin: &TxOrigin) -> bool {
        match (self, origin) {
            (TxOriginFilter::Main(id), TxOrigin::Database(origin_id)) => origin_id == id,
            (TxOriginFilter::Main(_), TxOrigin::Unscoped) => true,
            (TxOriginFilter::Main(_), TxOrigin::Memory) => false,
            (TxOriginFilter::Secondary(id), TxOrigin::Database(origin_id)) => origin_id == id,
            (TxOriginFilter::Secondary(_), TxOrigin::Unscoped | TxOrigin::Memory) => false,
        }
    }
}

/// Identity, age, label, and origin of the oldest span an attribution view
/// observes. Origin lets the consumer distinguish an evidence-backed
/// `Database(_)` winner from an `Unscoped` fallback winner when writing the
/// attribution-basis field.
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct OldestSpan {
    pub id: TxId,
    pub age: Duration,
    pub label: Option<String>,
    pub origin: TxOrigin,
}

#[derive(Clone, Debug)]
struct TxMeta {
    opened_at: Instant,
    label: Option<String>,
    origin: TxOrigin,
}

static NEXT_ID: AtomicU64 = AtomicU64::new(1);
static REGISTRY: LazyLock<Mutex<HashMap<TxId, TxMeta>>> =
    LazyLock::new(|| Mutex::new(HashMap::new()));

/// RAII handle for a registered transaction span. Deregisters on `Drop` — this is
/// the only deregistration path, so error and panic returns can never leak an
/// entry in the registry.
pub struct TxHandle {
    id: TxId,
}

impl Drop for TxHandle {
    fn drop(&mut self) {
        let mut registry = REGISTRY
            .lock()
            .unwrap_or_else(|poisoned| poisoned.into_inner());
        registry.remove(&self.id);
    }
}

/// Register a new open transaction span with an optional diagnostic label
/// and an explicit origin. Recovers a poisoned lock via `into_inner()`
/// rather than dropping the write — a poisoned registry must keep tracking
/// spans, not go silently blind. See `crates/khive-storage/docs/api/tx-registry.md`.
pub fn register_scoped(label: Option<String>, origin: TxOrigin) -> TxHandle {
    let id = TxId(NEXT_ID.fetch_add(1, Ordering::Relaxed));
    let meta = TxMeta {
        opened_at: Instant::now(),
        label,
        origin,
    };
    let mut registry = REGISTRY
        .lock()
        .unwrap_or_else(|poisoned| poisoned.into_inner());
    registry.insert(id, meta);
    drop(registry);
    TxHandle { id }
}

/// Register a new open transaction span with an optional diagnostic label.
/// Delegates to [`register_scoped`] with [`TxOrigin::Unscoped`].
pub fn register(label: Option<String>) -> TxHandle {
    register_scoped(label, TxOrigin::Unscoped)
}

/// Identity, age, and label of the oldest currently-open registry entry, if
/// any. The [`TxId`] lets callers distinguish "still the same oldest span"
/// from "a different span became oldest" (must re-arm latched escalation
/// state). Recovers a poisoned lock via `into_inner()` (see [`register`])
/// rather than returning `None`, which would read identically to the
/// genuinely-empty case. See `crates/khive-storage/docs/api/tx-registry.md`.
pub fn oldest() -> Option<(TxId, Duration, Option<String>)> {
    let registry = REGISTRY
        .lock()
        .unwrap_or_else(|poisoned| poisoned.into_inner());
    registry
        .iter()
        .min_by_key(|(_, meta)| meta.opened_at)
        .map(|(id, meta)| (*id, meta.opened_at.elapsed(), meta.label.clone()))
}

/// Identity, age, label, and origin of the oldest entry an attribution view
/// observes, if any. See [`TxOriginFilter`] for the main/secondary view
/// semantics and [`oldest`] for the process-wide aggregate this narrows.
/// Recovers a poisoned lock via `into_inner()` (see [`register_scoped`]).
pub fn oldest_for(filter: &TxOriginFilter) -> Option<OldestSpan> {
    let registry = REGISTRY
        .lock()
        .unwrap_or_else(|poisoned| poisoned.into_inner());
    registry
        .iter()
        .filter(|(_, meta)| filter.matches(&meta.origin))
        .min_by_key(|(_, meta)| meta.opened_at)
        .map(|(id, meta)| OldestSpan {
            id: *id,
            age: meta.opened_at.elapsed(),
            label: meta.label.clone(),
            origin: meta.origin.clone(),
        })
}

/// Whether an attribution view currently observes an open span carrying
/// `label`.
///
/// [`oldest_for`] cannot answer this. It reports the oldest span on the
/// backend, so any other span opened earlier on the same backend masks the
/// one the caller is asking about — a caller reasoning about the lifetime of
/// *its own* span reads a different span's presence as its own. Callers that
/// need "is this span still open" need this predicate; callers that need
/// "how long has this backend held a read open" need [`oldest_for`].
///
/// Recovers a poisoned lock via `into_inner()` (see [`register_scoped`]).
pub fn any_open_labeled(filter: &TxOriginFilter, label: &str) -> bool {
    let registry = REGISTRY
        .lock()
        .unwrap_or_else(|poisoned| poisoned.into_inner());
    registry
        .values()
        .any(|meta| filter.matches(&meta.origin) && meta.label.as_deref() == Some(label))
}

/// Age and label of every currently-open registry entry. Recovers a poisoned
/// lock via `into_inner()` (see [`register`]) instead of returning an empty
/// `Vec` — see `crates/khive-storage/docs/api/tx-registry.md`.
pub fn snapshot() -> Vec<(Duration, Option<String>)> {
    let registry = REGISTRY
        .lock()
        .unwrap_or_else(|poisoned| poisoned.into_inner());
    registry
        .values()
        .map(|meta| (meta.opened_at.elapsed(), meta.label.clone()))
        .collect()
}

#[cfg(test)]
mod tests {
    use super::*;
    use std::panic::{self, AssertUnwindSafe};

    // The registry is a process-wide singleton, shared across every test in this
    // binary (cargo runs `#[test]`s in parallel threads of the same process). A
    // module-local lock serializes these tests so one test's entries can't be
    // observed as another's "oldest" or leak into another's snapshot assertion.
    //
    // Acquired with the same poison recovery the registry itself uses. A failing
    // assertion panics while holding this lock, so under `.unwrap()` every later
    // test in the module dies of `PoisonError` instead of running: one real
    // failure presents as a dozen, and the reader has to work out which one is
    // the signal. Measured — a single-predicate mutation reddened 13 of 15 tests,
    // 10 of them purely as poison cascade.
    static TEST_LOCK: Mutex<()> = Mutex::new(());

    #[test]
    fn register_reports_oldest_with_label() {
        let _guard = TEST_LOCK
            .lock()
            .unwrap_or_else(|poisoned| poisoned.into_inner());
        let handle = register(Some("test_span".to_string()));
        let (_, _, label) = oldest().expect("expected an open entry");
        assert_eq!(label.as_deref(), Some("test_span"));
        drop(handle);
    }

    #[test]
    fn drop_deregisters() {
        let _guard = TEST_LOCK
            .lock()
            .unwrap_or_else(|poisoned| poisoned.into_inner());
        let handle = register(Some("drop_me".to_string()));
        let id_present_before = snapshot()
            .iter()
            .any(|(_, label)| label.as_deref() == Some("drop_me"));
        assert!(id_present_before);
        drop(handle);
        let id_present_after = snapshot()
            .iter()
            .any(|(_, label)| label.as_deref() == Some("drop_me"));
        assert!(!id_present_after);
    }

    #[test]
    fn oldest_is_genuinely_oldest() {
        let _guard = TEST_LOCK
            .lock()
            .unwrap_or_else(|poisoned| poisoned.into_inner());
        let first = register(Some("first".to_string()));
        std::thread::sleep(Duration::from_millis(5));
        let second = register(Some("second".to_string()));
        let (first_id, _, label) = oldest().expect("expected an open entry");
        assert_eq!(label.as_deref(), Some("first"));
        drop(first);
        let (second_id, _, label) = oldest().expect("expected an open entry");
        assert_eq!(label.as_deref(), Some("second"));
        assert_ne!(
            first_id, second_id,
            "distinct registrations must carry distinct TxIds"
        );
        drop(second);
    }

    #[test]
    fn snapshot_contains_all_open_entries() {
        let _guard = TEST_LOCK
            .lock()
            .unwrap_or_else(|poisoned| poisoned.into_inner());
        let a = register(Some("snap_a".to_string()));
        let b = register(Some("snap_b".to_string()));
        let labels: Vec<Option<String>> = snapshot().into_iter().map(|(_, label)| label).collect();
        assert!(labels.contains(&Some("snap_a".to_string())));
        assert!(labels.contains(&Some("snap_b".to_string())));
        drop(a);
        drop(b);
    }

    #[test]
    fn oldest_for_partitions_by_origin() {
        let _guard = TEST_LOCK
            .lock()
            .unwrap_or_else(|poisoned| poisoned.into_inner());
        let main_id = DbIdentity::new("main.db");
        let secondary_id = DbIdentity::new("secondary.db");
        let main_view = TxOriginFilter::Main(main_id.clone());
        let secondary_view = TxOriginFilter::Secondary(secondary_id.clone());

        let main_handle = register_scoped(
            Some("main_span".to_string()),
            TxOrigin::Database(main_id.clone()),
        );
        let secondary_handle = register_scoped(
            Some("secondary_span".to_string()),
            TxOrigin::Database(secondary_id.clone()),
        );

        let main_oldest = oldest_for(&main_view).expect("main span visible in main view");
        assert_eq!(main_oldest.label.as_deref(), Some("main_span"));
        assert_eq!(main_oldest.origin, TxOrigin::Database(main_id));

        let secondary_oldest =
            oldest_for(&secondary_view).expect("secondary span visible in secondary view");
        assert_eq!(secondary_oldest.label.as_deref(), Some("secondary_span"));

        drop(main_handle);
        drop(secondary_handle);
    }

    #[test]
    fn register_delegates_to_unscoped_origin() {
        let _guard = TEST_LOCK
            .lock()
            .unwrap_or_else(|poisoned| poisoned.into_inner());
        let main_id = DbIdentity::new("main.db");
        let main_view = TxOriginFilter::Main(main_id);
        let handle = register(Some("legacy_span".to_string()));

        let oldest = oldest_for(&main_view).expect("register() delegates to Unscoped");
        assert_eq!(oldest.origin, TxOrigin::Unscoped);

        drop(handle);
    }

    #[test]
    fn unscoped_visible_in_main_view_and_oldest() {
        let _guard = TEST_LOCK
            .lock()
            .unwrap_or_else(|poisoned| poisoned.into_inner());
        let main_id = DbIdentity::new("main.db");
        let main_view = TxOriginFilter::Main(main_id);
        let handle = register_scoped(Some("unscoped_span".to_string()), TxOrigin::Unscoped);

        let via_filter = oldest_for(&main_view).expect("unscoped span falls back to main view");
        assert_eq!(via_filter.label.as_deref(), Some("unscoped_span"));
        assert_eq!(via_filter.origin, TxOrigin::Unscoped);

        let via_aggregate = oldest().expect("unscoped span visible in aggregate oldest()");
        assert_eq!(via_aggregate.2.as_deref(), Some("unscoped_span"));

        drop(handle);
    }

    #[test]
    fn memory_absent_from_attribution_views_but_present_in_oldest() {
        let _guard = TEST_LOCK
            .lock()
            .unwrap_or_else(|poisoned| poisoned.into_inner());
        let main_id = DbIdentity::new("main.db");
        let secondary_id = DbIdentity::new("secondary.db");
        let main_view = TxOriginFilter::Main(main_id);
        let secondary_view = TxOriginFilter::Secondary(secondary_id);
        let handle = register_scoped(Some("memory_span".to_string()), TxOrigin::Memory);

        assert!(oldest_for(&main_view).is_none());
        assert!(oldest_for(&secondary_view).is_none());

        let via_aggregate = oldest().expect("memory span visible in aggregate oldest()");
        assert_eq!(via_aggregate.2.as_deref(), Some("memory_span"));

        drop(handle);
    }

    #[test]
    fn database_entries_never_leak_across_views() {
        let _guard = TEST_LOCK
            .lock()
            .unwrap_or_else(|poisoned| poisoned.into_inner());
        let main_id = DbIdentity::new("main.db");
        let secondary_id = DbIdentity::new("secondary.db");
        let main_view = TxOriginFilter::Main(main_id);
        let secondary_view = TxOriginFilter::Secondary(secondary_id.clone());
        let handle = register_scoped(
            Some("secondary_span".to_string()),
            TxOrigin::Database(secondary_id),
        );

        assert!(oldest_for(&main_view).is_none());
        let via_secondary =
            oldest_for(&secondary_view).expect("secondary span visible in its own view");
        assert_eq!(via_secondary.label.as_deref(), Some("secondary_span"));

        drop(handle);
    }

    #[test]
    fn panic_inside_scope_still_deregisters() {
        let _guard = TEST_LOCK
            .lock()
            .unwrap_or_else(|poisoned| poisoned.into_inner());
        let result = panic::catch_unwind(AssertUnwindSafe(|| {
            let _handle = register(Some("panics".to_string()));
            panic!("boom");
        }));
        assert!(result.is_err());
        let still_present = snapshot()
            .iter()
            .any(|(_, label)| label.as_deref() == Some("panics"));
        assert!(!still_present);
    }

    #[test]
    fn any_open_labeled_matches_only_the_named_label() {
        let _guard = TEST_LOCK
            .lock()
            .unwrap_or_else(|poisoned| poisoned.into_inner());
        let main_id = DbIdentity::new("main.db");
        let main_view = TxOriginFilter::Main(main_id.clone());
        let handle = register_scoped(
            Some("wanted".to_string()),
            TxOrigin::Database(main_id.clone()),
        );

        assert!(any_open_labeled(&main_view, "wanted"));
        assert!(!any_open_labeled(&main_view, "other"));
        assert!(!any_open_labeled(&main_view, "wante"));

        drop(handle);
        assert!(!any_open_labeled(&main_view, "wanted"));
    }

    /// The whole reason this predicate exists: an unrelated span open on the
    /// same backend satisfies [`oldest_for`] while saying nothing about the
    /// span the caller named. A caller asking "is *my* span still open" must
    /// not read a different span's presence as its own.
    #[test]
    fn any_open_labeled_ignores_unrelated_spans_that_oldest_for_reports() {
        let _guard = TEST_LOCK
            .lock()
            .unwrap_or_else(|poisoned| poisoned.into_inner());
        let main_id = DbIdentity::new("main.db");
        let main_view = TxOriginFilter::Main(main_id.clone());
        let unrelated = register_scoped(
            Some("someone_elses_write".to_string()),
            TxOrigin::Database(main_id.clone()),
        );

        assert!(
            oldest_for(&main_view).is_some(),
            "aggregate must observe the unrelated span — otherwise this test \
             proves nothing about the discrimination"
        );
        assert!(!any_open_labeled(&main_view, "mine"));

        drop(unrelated);
    }

    #[test]
    fn any_open_labeled_survives_duplicate_labels() {
        let _guard = TEST_LOCK
            .lock()
            .unwrap_or_else(|poisoned| poisoned.into_inner());
        let main_id = DbIdentity::new("main.db");
        let main_view = TxOriginFilter::Main(main_id.clone());
        let first = register_scoped(Some("dup".to_string()), TxOrigin::Database(main_id.clone()));
        let second = register_scoped(Some("dup".to_string()), TxOrigin::Database(main_id.clone()));

        assert!(any_open_labeled(&main_view, "dup"));
        drop(first);
        assert!(
            any_open_labeled(&main_view, "dup"),
            "one of two same-labeled spans closing must not read as all of them closing"
        );
        drop(second);
        assert!(!any_open_labeled(&main_view, "dup"));
    }

    #[test]
    fn any_open_labeled_never_matches_an_unlabeled_span() {
        let _guard = TEST_LOCK
            .lock()
            .unwrap_or_else(|poisoned| poisoned.into_inner());
        let main_id = DbIdentity::new("main.db");
        let main_view = TxOriginFilter::Main(main_id.clone());
        let unlabeled = register_scoped(None, TxOrigin::Database(main_id.clone()));

        assert!(
            oldest_for(&main_view).is_some(),
            "the unlabeled span must be registered for this test to discriminate"
        );
        assert!(!any_open_labeled(&main_view, "anything"));

        drop(unlabeled);
    }

    #[test]
    fn any_open_labeled_partitions_by_origin() {
        let _guard = TEST_LOCK
            .lock()
            .unwrap_or_else(|poisoned| poisoned.into_inner());
        let main_id = DbIdentity::new("main.db");
        let secondary_id = DbIdentity::new("secondary.db");
        let main_view = TxOriginFilter::Main(main_id);
        let secondary_view = TxOriginFilter::Secondary(secondary_id.clone());
        let handle = register_scoped(
            Some("secondary_only".to_string()),
            TxOrigin::Database(secondary_id),
        );

        assert!(any_open_labeled(&secondary_view, "secondary_only"));
        assert!(!any_open_labeled(&main_view, "secondary_only"));

        drop(handle);
    }
}