Skip to main content

Db

Struct Db 

Source
pub struct Db {
    pub objects: ObjectStore,
    pub id_index: IdIndex,
    pub sorted_indexes: SortedIndexes,
    pub graph: GraphStore,
    pub root: PathBuf,
    pub seq: Atomic<u64>,
    pub startup_ready: Arc<Atomic<bool>>,
    /* private fields */
}

Fields§

§objects: ObjectStore§id_index: IdIndex§sorted_indexes: SortedIndexes§graph: GraphStore§root: PathBuf§seq: Atomic<u64>§startup_ready: Arc<Atomic<bool>>

True once startup is fully ready (MANIFEST loaded or cold scan complete). Warm starts set this true before returning from open(). Cold starts set this true in the background thread when scan completes. Writes are held with 503 until this is true; reads always proceed.

Implementations§

Source§

impl Db

Source

pub fn in_memory() -> Db

Create a pure in-memory database — no disk I/O, no migration, instant startup. Perfect for tests, hot-cache layers, and ephemeral sessions. All data is lost when the Db is dropped.

Source

pub fn open(db_root: &Path, dek: Option<Dek>) -> Result<Db, Error>

Open (or create) a database. Runs v1→v2 migration automatically if log.aof is present.

Source

pub fn start_cold_scan(self_arc: Arc<Db>)

Call this from Manager::open_all() after Arc::new(db). Spawns the cold scan background thread with stable heap addresses. No-op if startup is already complete (warm start).

Source

pub fn rebuild_id_index(&self) -> Result<usize, Error>

Rebuild the id index from the object store, synchronously.

Every object carries its own coll, id and seq, so the id index is fully derivable: for each (coll, id) the highest seq wins. Use this to recover a database whose id-index WAL never reached disk — the objects are intact and verify, but list()/get() return nothing.

Idempotent, and safe on a healthy store (it rewrites the same winners). Returns the number of entries written. Flushes before returning.

Source

pub fn repair(&self) -> Result<usize, Error>

Full repair: rebuild the seq index and the id index from objects, even on a WARM store, then flush.

[start_cold_scan] deliberately no-ops when startup is already complete, which meant the documented repair path (“idempotent — a no-op on a warm store, a full self-heal on a stale MANIFEST”) could never repair a database that had a valid MANIFEST and a damaged id index. This is the forcing entry point; start_cold_scan keeps its O(1) warm-boot contract.

Source

pub fn put( &self, coll: &str, id: &str, data: Value, caused_by: Vec<String>, valid_from: Option<String>, valid_to: Option<String>, ) -> Result<Node, Error>

Write a document. Returns the new node with its content hash set.

Source

pub fn put_batch( &self, ops: Vec<(String, String, Value, Vec<String>, Option<String>, Option<String>)>, ) -> Result<Vec<Node>, Error>

Batch put: write N documents in parallel, preserving monotonic seq ordering. Pre-allocates N seq numbers atomically, then parallelises object writes and id-index updates via Rayon. Each op is independent — safe to parallelise. Returns nodes in input order with assigned seq numbers.

Source

pub fn try_flush_all(&self) -> Result<(), Error>

Flush both the id-index WAL and MANIFEST, REPORTING failure.

This is the durability boundary: until it returns Ok(()), writes that put() acknowledged may not be on disk. Callers that must not lose data — anything about to take a destructive or externally-visible action on the strength of a persisted record — should use this, not [flush_all].

Every stage is attempted even if an earlier one fails (a MANIFEST flush is still worth doing when one index leaf failed), and the first error is returned. Failed id-index entries stay in the WAL for retry.

Source

pub fn flush_all(&self)

Flush both the id-index WAL and MANIFEST. Used on graceful shutdown.

Errors are logged, not returned — kept for back-compat and for the ticker/Drop paths that have nowhere to propagate. Prefer [try_flush_all] whenever the outcome matters.

Source

pub fn compact(&self) -> Result<CompactStats, Error>

Compact the v3 packed object store: keep the CURRENT version of every document (from the id-index) and reclaim everything else. No-op unless running with the v3 segment substrate (--dag-v3 / NEDB_DAG_V3).

This is a PRUNING operation: superseded/historical object versions are dropped, so AS OF / TRACE over pruned versions is discarded — that is what reclaims the space. Flushes first so all data is durable on disk before the old segments are deleted.

Source

pub fn flush_manifest_if_dirty(&self)

Flush MANIFEST to disk if dirty. No-op for in-memory databases.

Source

pub fn try_flush_manifest(&self) -> Result<(), Error>

Atomically persist current seq+head to MANIFEST, reporting failure. No-op (Ok) for in-memory databases.

A silently failed MANIFEST write is not data loss — the startup self-heal rescans — but it IS a warm-boot regression and, on a full disk, the first symptom that persistence is failing. Callers deserve to know.

Source

pub fn flush_manifest(&self)

Atomically persist current seq+head to MANIFEST. No-op for in-memory databases. Errors are logged; prefer [try_flush_manifest] when the outcome matters.

Source

pub fn embedded_flush_interval_ms() -> Option<u64>

Start a background thread that flushes both the id-index WAL and MANIFEST every interval_ms milliseconds. Call this after Arc::new(db) — the Arc keeps Db alive for the thread’s lifetime. Flush cadence for EMBEDDED durable handles (the napi and pyo3 open() paths).

nedbd has always run the manifest ticker at 1 s, so a server flushes the id-index WAL and MANIFEST every second and a hard kill loses at most a second of acknowledged writes. The embedded bindings did not start a ticker at all: their WAL was flushed only by the exit hooks (SIGINT/SIGTERM/atexit) — so an embedded app killed with SIGKILL, OOM-killed, or cut by power lost EVERY write since open, with no bound. Found by CHALK / Sports-Rater on 2026-09-04 (acknowledged fan writes gone after kill -9). Since 2.8.5 the bindings start the ticker on durable open with this cadence — parity with nedbd.

NEDB_FLUSH_MS overrides: an integer of milliseconds (min 50), or 0 / off to disable (only for hosts that own their own flush cadence). Unset → 1000.

Source

pub fn start_manifest_ticker(self_arc: Arc<Db>, interval_ms: u64)

Spawn the background flush ticker.

The ticker holds a Weak<Db> and exits the first time the upgrade fails — i.e. as soon as the last real owner drops the database. The caller must therefore keep its own Arc alive for as long as it wants ticking; every current caller already does (nedbd stores it in its database map, the napi and pyo3 handles own theirs).

It used to hold a strong Arc inside an unconditional loop, which meant the thread never exited and the Db was never dropped. Three consequences, all of them live since 2.8.5:

  • The exclusive data-dir LOCK taken in Db::open was never released, so reopening the same path in the same process failed with “locked by another process (pid N)” where N was the caller’s own pid.
  • Every open() leaked a thread and the entire Db — indexes, caches, segment handles — for the lifetime of the process.
  • Drop for Db (flush-on-close) could never fire for embedded users, exactly as its own doc comment warned: it “only fires once every owning handle is gone”, and an immortal thread always held one.

nedbd’s drop_db was hit by the same thing: removing a database from the map did not free it, and an orphaned ticker went on fsyncing it.

The Arc is upgraded inside the loop and dropped before the next sleep, so the ticker never extends the database’s life across a tick. No final flush is needed here — the owner’s Drop does it.

Source

pub fn head(&self) -> String

Return the current Merkle head string. O(1) — read from cache.

Source

pub fn delete(&self, coll: &str, id: &str) -> Result<bool, Error>

Delete a document — writes a tombstone node and removes the id from the index. The object history is preserved in the DAG; only the live id pointer is cleared.

Source

pub fn get(&self, coll: &str, id: &str) -> Option<Node>

Get the current version of a document by id.

Source

pub fn get_by_hash(&self, hash: &str) -> Option<Node>

Get a specific version of a document by object hash.

Source

pub fn get_as_of(&self, coll: &str, id: &str, target_seq: u64) -> Option<Node>

Get a document AS OF a specific sequence number. Walks the version chain (prev links) backward until seq <= target.

Source

pub fn list(&self, coll: &str) -> Vec<Node>

List all documents in a collection, returning current versions.

Source

pub fn range_scan( &self, coll: &str, field: &str, low: Option<&Value>, high: Option<&Value>, low_incl: bool, high_incl: bool, ) -> Option<Vec<Node>>

Candidate nodes whose field falls in the given range, via the sorted index. None when no index covers (coll, field) — the caller must then fall back to a scan.

Returns CURRENT versions only (the index drops a superseded hash on overwrite), so this must not be used to serve an AS OF query.

Source

pub fn index_lookup( &self, coll: &str, field: &str, values: &[Value], ) -> Option<Vec<Node>>

Candidate nodes whose field equals any of values — the indexed path for = and for IN (...). None when no index covers the field.

Source

pub fn range_cardinality( &self, coll: &str, field: &str, low: Option<&Value>, high: Option<&Value>, low_incl: bool, high_incl: bool, ) -> Option<usize>

How many rows an indexed range covers, without reading any of them. None when no index covers the field.

Source

pub fn has_sorted_index(&self, coll: &str, field: &str) -> bool

True when a sorted index covers (coll, field).

Source

pub fn order_by_asc(&self, coll: &str, field: &str, limit: usize) -> Vec<Node>

ORDER BY field ASC LIMIT n — uses sorted index if available, else falls back to full scan.

Source

pub fn order_by_desc(&self, coll: &str, field: &str, limit: usize) -> Vec<Node>

ORDER BY field DESC LIMIT n

Source

pub fn trace(&self, hash: &str, reverse: bool, limit: usize) -> Vec<Node>

TRACE caused_by — walk causal graph from a node.

Source

pub fn verify(&self) -> (usize, Vec<String>)

Verify tamper-evidence of all objects.

Source

pub fn create_sorted_index(&self, coll: &str, field: &str)

Create a sorted index for a (coll, field) pair.

Source

pub fn get_hash_by_seq(&self, seq: u64) -> Option<String>

Resolve a sequence number to its content hash (v1 compatibility). Only covers nodes written in the current process session + cold-scan nodes.

Source

pub fn tip(&self) -> Option<Node>

The tip — the most recently written node (highest seq), or None if the database is empty. O(1): self.seq is the next-to-assign counter, so the latest write sits at seq - 1; we resolve it through the same seq_index → object-store path a normal read uses, so the returned Node is byte-identical to one fetched by id or hash (it carries its own seq, hash, causal links, and valid-time). This is the cheap “give me the latest write” primitive — the head of the log, not an aggregate.

Source

pub fn tip_collection(&self, coll: &str) -> Option<Node>

The collection-local tip — the most recent write into coll (highest seq in that collection), or None if the collection has no writes. O(1): resolves through coll_tip_hash, a dedicated per-collection map kept current on every write (update_head), restored from MANIFEST on warm boot, and rebuilt by the cold scan — durable across restarts by construction, same contract as tip() for the global head. Conceptually a different index than the global tip() (global head vs collection head), kept as a separate method so each is explicit — parity with the Python reference’s tip(coll). Lets a consumer resume one chain (e.g. blocks / tx / utxo) without pulling global tip and filtering.

Source

pub fn since(&self, after_seq: u64, limit: usize) -> SinceBatch

Changefeed page: up to limit nodes written AFTER after_seq (EXCLUSIVE), ascending by seq, wrapped in a SinceBatch cursor envelope. after_seq is the cursor you last applied (a prior tip() seq or to_seq). limit bounds the page — 0 means DEFAULT_SINCE_LIMIT, so the engine primitive can never materialize an unbounded batch even when embedders call it directly (the safety is here, not only in the HTTP layer). Drain by paging while has_more, advancing your cursor to to_seq, then hand off to the live subscribe edge. The append-only log IS the changefeed, so this is an O(page) walk; unresolved seqs (outside seq_index coverage — see scan_status()) are skipped rather than faked.

Source

pub fn scan_status(&self) -> ScanStatus

Replication readiness — see ScanStatus. scan_complete gates safe historical catch-up: a consumer pulling an old cursor right after a cold start must wait for it, or since() may hand back a partial page that looks like “caught up”. Computes the indexed range by scanning the in-memory seq index (O(index)) — intended for periodic status polls, not the per-write hot path.

Add an explicit named relation edge between two documents. Add an explicit named relation between two “coll:id” nodes. Relations stored as links documents — NQL-queryable, time-travelable, consistent with the PyO3 binding which uses the same links convention.

Remove a named relation (deletes the links document).

Source

pub fn neighbors(&self, frm: &str, rel: &str) -> Vec<Node>

Get neighbor nodes via a named relation. Queries links — consistent with the PyO3 binding.

Source§

impl Db

Source

pub fn install_exit_flush(self_arc: Arc<Db>)

Flush this durable database’s buffered state on SIGINT/SIGTERM (Ctrl+C, kill, orchestrator shutdown) — the flush-on-close contract extended to hard exits that never run Drop.

Call once, after the database is wrapped in an Arc (the registry holds a Weak, so this never keeps the Db alive). Idempotent; safe to call from multiple databases. A no-op for in-memory (:memory:) databases.

let db = Arc::new(Db::open(std::path::Path::new("/data/mydb"), None)?);
Db::install_exit_flush(Arc::clone(&db));   // durable across Ctrl+C / SIGTERM

Trait Implementations§

Source§

impl Drop for Db

Source§

fn drop(&mut self)

Flush buffered state when the database is closed so a write-then-drop sequence is durable without an explicit flush_all().

IdIndex::set only stages updates in the in-memory WAL write_buf; disk persistence happens in flush_write_buf(), normally driven by the manifest ticker. A short-lived Db (a library user’s { let db = Db::open(p)?; db.put(..)?; } block, or a test) has no ticker, so without this its writes would be silently lost on reopen. Flushing on drop mirrors the flush-on-close contract of other embedded stores (sled, RocksDB).

In production this is a harmless safety net, not the primary durability path: the manifest ticker thread holds an Arc<Db> for the process lifetime, so Drop only fires once every owning handle is gone. No-op for in-memory databases (flush_all short-circuits on :memory:).

Source§

fn pin_drop(self: Pin<&mut Self>)

🔬This is a nightly-only experimental API. (pin_ergonomics)
Execute the destructor for this type, but different to Drop::drop, it requires self to be pinned. Read more

Auto Trait Implementations§

§

impl !Freeze for Db

§

impl !RefUnwindSafe for Db

§

impl !UnwindSafe for Db

§

impl Send for Db

§

impl Sync for Db

§

impl Unpin for Db

§

impl UnsafeUnpin for Db

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<A, B, T> HttpServerConnExec<A, B> for T
where B: Body,

Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self>

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self>

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> IntoEither for T

Source§

fn into_either(self, into_left: bool) -> Either<Self, Self>

Converts self into a Left variant of Either<Self, Self> if into_left is true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
where F: FnOnce(&Self) -> bool,

Converts self into a Left variant of Either<Self, Self> if into_left(&self) returns true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

impl<T> Pointable for T

Source§

const ALIGN: usize

The alignment of pointer.
Source§

type Init = T

The type for initializers.
Source§

unsafe fn init(init: <T as Pointable>::Init) -> usize

Initializes a with the given initializer. Read more
Source§

unsafe fn deref<'a>(ptr: usize) -> &'a T

Dereferences the given pointer. Read more
Source§

unsafe fn deref_mut<'a>(ptr: usize) -> &'a mut T

Mutably dereferences the given pointer. Read more
Source§

unsafe fn drop(ptr: usize)

Drops the object pointed to by the given pointer. Read more
Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = !

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, !>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self>
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self>

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more