Skip to main content

corium_transactor/
node.rs

1//! The transactor as a process: multi-database state, durable naming,
2//! lease acquisition/renewal, background indexing, and tx-report fan-out.
3
4use std::collections::HashMap;
5use std::path::PathBuf;
6use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
7use std::sync::{Arc, Mutex};
8use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
9
10use corium_core::{KeywordInterner, Schema};
11use corium_db::{Db, Idents};
12use corium_log::{FileLog, LogError, TransactionLog, TxRecord};
13use corium_protocol::codec::{self, CodecError};
14use corium_protocol::pb;
15use corium_protocol::schemaform::{SchemaFormError, schema_from_edn};
16use corium_protocol::txforms::{TxFormError, tx_items_from_edn};
17use corium_query::edn::Edn;
18use corium_store::{FsStore, RootStore, StoreError, mark_and_sweep_retained};
19use thiserror::Error;
20use tokio::sync::{broadcast, watch};
21use tracing::Instrument;
22
23use crate::lease::{self, Lease, LeaseError};
24use crate::metrics::Metrics;
25use crate::{DbRoot, EmbeddedTransactor, TransactError, db_root_name, publish_root};
26
27/// Expands user database-function invocations in boundary EDN transaction
28/// forms before native conversion. Implemented by `corium-cljrs` (the
29/// sandboxed Clojurust host, ADR-0008) and injected by the process wiring;
30/// the transactor itself stays free of cljrs dependencies.
31pub trait TxFnExpander: Send + Sync {
32    /// Rewrites `forms` with every `[:my/fn arg…]` invocation replaced by
33    /// the function's returned tx-data (recursively).
34    ///
35    /// # Errors
36    /// Returns a display message when a function is missing, rejected by
37    /// the sandbox, fails, or exceeds its budget; the transaction aborts.
38    fn expand(&self, db: &Db, forms: Vec<Edn>) -> Result<Vec<Edn>, String>;
39}
40
41/// Node process configuration.
42#[derive(Clone)]
43pub struct NodeConfig {
44    /// Data directory holding the blob/root store and transaction logs.
45    pub data_dir: PathBuf,
46    /// Stable owner identity for lease records.
47    pub owner: String,
48    /// Lease time-to-live in milliseconds.
49    pub lease_ttl_ms: i64,
50    /// How long to wait for a held lease to expire before giving up.
51    pub lease_wait_ms: i64,
52    /// Interval between background index publications.
53    pub index_interval: Duration,
54    /// Interval between heartbeats on subscription streams.
55    pub heartbeat_interval: Duration,
56    /// Interval between scheduled garbage-collection duties; `None` disables it.
57    pub gc_interval: Option<Duration>,
58    /// Minimum age of an unreachable blob before scheduled/manual online GC.
59    pub gc_retention: Duration,
60    /// Optional database-function expander (`:db/fn` support).
61    pub tx_fn_expander: Option<Arc<dyn TxFnExpander>>,
62}
63
64impl std::fmt::Debug for NodeConfig {
65    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
66        f.debug_struct("NodeConfig")
67            .field("data_dir", &self.data_dir)
68            .field("owner", &self.owner)
69            .field("lease_ttl_ms", &self.lease_ttl_ms)
70            .field("lease_wait_ms", &self.lease_wait_ms)
71            .field("index_interval", &self.index_interval)
72            .field("heartbeat_interval", &self.heartbeat_interval)
73            .field("gc_interval", &self.gc_interval)
74            .field("gc_retention", &self.gc_retention)
75            .field("tx_fn_expander", &self.tx_fn_expander.is_some())
76            .finish()
77    }
78}
79
80impl NodeConfig {
81    /// Sensible defaults for a data directory.
82    #[must_use]
83    pub fn new(data_dir: PathBuf) -> Self {
84        Self {
85            data_dir,
86            owner: format!(
87                "transactor-{}",
88                std::env::var("HOSTNAME").unwrap_or_else(|_| "local".into())
89            ),
90            lease_ttl_ms: 5_000,
91            lease_wait_ms: 15_000,
92            index_interval: Duration::from_secs(5),
93            heartbeat_interval: Duration::from_secs(10),
94            gc_interval: Some(Duration::from_secs(60 * 60)),
95            gc_retention: Duration::from_secs(72 * 60 * 60),
96            tx_fn_expander: None,
97        }
98    }
99}
100
101/// Node operation failure.
102#[derive(Debug, Error)]
103pub enum NodeError {
104    /// Named database does not exist.
105    #[error("unknown database {0:?}")]
106    UnknownDb(String),
107    /// Database name is not storable.
108    #[error("invalid database name {0:?}")]
109    InvalidName(String),
110    /// Database root uses a storage format newer than this binary.
111    #[error("storage format {found} is newer than supported format {supported}")]
112    UnsupportedFormat {
113        /// Version found in the root.
114        found: u32,
115        /// Newest version understood by this binary.
116        supported: u32,
117    },
118    /// This node no longer holds the write lease.
119    #[error("deposed: write lease for {0:?} is held elsewhere")]
120    Deposed(String),
121    /// Payload failed to decode.
122    #[error(transparent)]
123    Codec(#[from] CodecError),
124    /// Transaction forms failed to convert.
125    #[error(transparent)]
126    TxForm(#[from] TxFormError),
127    /// Schema forms failed to convert.
128    #[error(transparent)]
129    SchemaForm(#[from] SchemaFormError),
130    /// Transaction pipeline failure.
131    #[error(transparent)]
132    Transact(#[from] TransactError),
133    /// Store failure.
134    #[error(transparent)]
135    Store(#[from] StoreError),
136    /// Log failure.
137    #[error(transparent)]
138    Log(#[from] LogError),
139    /// Lease failure.
140    #[error(transparent)]
141    Lease(#[from] LeaseError),
142    /// Malformed request.
143    #[error("bad request: {0}")]
144    BadRequest(String),
145}
146
147struct Naming {
148    schema: Schema,
149    idents: Idents,
150    interner: KeywordInterner,
151}
152
153/// Per-database state hosted by a node.
154pub struct DbState {
155    name: String,
156    transactor: EmbeddedTransactor,
157    log: Arc<FileLog>,
158    naming: Mutex<Naming>,
159    commit: tokio::sync::Mutex<()>,
160    broadcast: broadcast::Sender<pb::subscribe_item::Item>,
161    basis: watch::Sender<u64>,
162    index_basis: AtomicU64,
163    held_lease: Mutex<Lease>,
164    deposed: AtomicBool,
165}
166
167impl DbState {
168    /// Database name.
169    #[must_use]
170    pub fn name(&self) -> &str {
171        &self.name
172    }
173
174    /// Current database value.
175    #[must_use]
176    pub fn db(&self) -> Db {
177        self.transactor.db()
178    }
179
180    /// Watch channel following the commit basis.
181    #[must_use]
182    pub fn basis_watch(&self) -> watch::Receiver<u64> {
183        self.basis.subscribe()
184    }
185
186    /// Subscribes to live stream items (reports, index announcements,
187    /// heartbeats).
188    #[must_use]
189    pub fn stream_items(&self) -> broadcast::Receiver<pb::subscribe_item::Item> {
190        self.broadcast.subscribe()
191    }
192
193    /// Basis of the newest published index root.
194    #[must_use]
195    pub fn index_basis(&self) -> u64 {
196        self.index_basis.load(Ordering::Acquire)
197    }
198
199    /// Currently held lease record.
200    #[must_use]
201    pub fn lease(&self) -> Lease {
202        self.held_lease
203            .lock()
204            .unwrap_or_else(std::sync::PoisonError::into_inner)
205            .clone()
206    }
207
208    /// Encoded schema/ident handshake payload plus a consistent basis and
209    /// interner snapshot for backfill encoding.
210    #[must_use]
211    pub fn handshake_snapshot(&self) -> (Vec<u8>, KeywordInterner) {
212        let naming = self
213            .naming
214            .lock()
215            .unwrap_or_else(std::sync::PoisonError::into_inner);
216        (
217            codec::encode_schema(&naming.schema, &naming.idents),
218            naming.interner.clone(),
219        )
220    }
221
222    /// Reads committed records in `[start, end)` from the durable log.
223    ///
224    /// # Errors
225    /// Returns an error when the log cannot be read.
226    pub fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, NodeError> {
227        Ok(self.log.tx_range(start, end)?)
228    }
229
230    fn check_lease(&self, store: &dyn RootStore) -> Result<Lease, NodeError> {
231        if self.deposed.load(Ordering::Acquire) {
232            return Err(NodeError::Deposed(self.name.clone()));
233        }
234        let held = self.lease();
235        let stored = store.get_root(&lease::lease_root(&self.name))?;
236        if stored.as_deref() == Some(held.encode().as_slice()) {
237            return Ok(held);
238        }
239        // Renewal publishes the new expiry before updating `held_lease`.
240        // A transaction can land in that tiny interval; matching owner and
241        // fence version is still the same lease, so adopt the newer record.
242        if let Some(stored) = stored.as_deref().and_then(Lease::decode) {
243            if stored.owner == held.owner && stored.version == held.version {
244                *self
245                    .held_lease
246                    .lock()
247                    .unwrap_or_else(std::sync::PoisonError::into_inner) = stored.clone();
248                return Ok(stored);
249            }
250        }
251        self.deposed.store(true, Ordering::Release);
252        Err(NodeError::Deposed(self.name.clone()))
253    }
254}
255
256/// A running transactor node hosting every database under one data directory.
257pub struct TransactorNode {
258    config: NodeConfig,
259    store: Arc<FsStore>,
260    dbs: std::sync::RwLock<HashMap<String, Arc<DbState>>>,
261    gc_lock: tokio::sync::Mutex<()>,
262    metrics: Metrics,
263    shutdown: watch::Sender<Option<String>>,
264}
265
266fn now_unix_ms() -> i64 {
267    i64::try_from(
268        SystemTime::now()
269            .duration_since(UNIX_EPOCH)
270            .unwrap_or_default()
271            .as_millis(),
272    )
273    .unwrap_or(i64::MAX)
274}
275
276fn valid_db_name(name: &str) -> bool {
277    !name.is_empty()
278        && name.len() <= 128
279        && name
280            .bytes()
281            .all(|byte| byte.is_ascii_alphanumeric() || byte == b'-' || byte == b'_')
282}
283
284fn meta_root_name(db: &str) -> String {
285    format!("meta:{db}")
286}
287
288fn encode_meta(schema: &Schema, idents: &Idents, interner: &KeywordInterner) -> Vec<u8> {
289    let schema_bytes = codec::encode_schema(schema, idents);
290    let naming_bytes = codec::encode_naming(interner);
291    let mut out = Vec::with_capacity(8 + schema_bytes.len() + naming_bytes.len());
292    out.extend_from_slice(&u32::try_from(schema_bytes.len()).unwrap_or(0).to_be_bytes());
293    out.extend_from_slice(&schema_bytes);
294    out.extend_from_slice(&u32::try_from(naming_bytes.len()).unwrap_or(0).to_be_bytes());
295    out.extend_from_slice(&naming_bytes);
296    out
297}
298
299fn decode_meta(bytes: &[u8]) -> Result<(Schema, Idents, KeywordInterner), NodeError> {
300    let take = |input: &mut &[u8]| -> Result<Vec<u8>, NodeError> {
301        let len_bytes = input
302            .get(..4)
303            .ok_or(NodeError::Codec(CodecError::Truncated))?;
304        let len = usize::try_from(u32::from_be_bytes(len_bytes.try_into().unwrap_or_default()))
305            .map_err(|_| NodeError::Codec(CodecError::Length))?;
306        let payload = input
307            .get(4..4 + len)
308            .ok_or(NodeError::Codec(CodecError::Truncated))?
309            .to_vec();
310        *input = &input[4 + len..];
311        Ok(payload)
312    };
313    let mut input = bytes;
314    let schema_bytes = take(&mut input)?;
315    let naming_bytes = take(&mut input)?;
316    let (schema, idents) = codec::decode_schema(&schema_bytes)?;
317    let interner = codec::decode_naming(&naming_bytes)?;
318    Ok((schema, idents, interner))
319}
320
321impl TransactorNode {
322    /// Opens a node over `config.data_dir`, recovering every database found
323    /// there (acquiring its lease, waiting out held leases up to the
324    /// configured bound).
325    ///
326    /// # Errors
327    /// Returns an error when the store cannot be opened or a database cannot
328    /// be recovered.
329    pub fn open(config: NodeConfig) -> Result<Arc<Self>, NodeError> {
330        let store = Arc::new(FsStore::open(config.data_dir.join("store"))?);
331        let node = Arc::new(Self {
332            config,
333            store,
334            dbs: std::sync::RwLock::new(HashMap::new()),
335            gc_lock: tokio::sync::Mutex::new(()),
336            metrics: Metrics::default(),
337            shutdown: watch::channel(None).0,
338        });
339        let names: Vec<String> = node
340            .store
341            .list_roots("meta:")?
342            .into_iter()
343            .filter_map(|root| root.strip_prefix("meta:").map(str::to_owned))
344            .collect();
345        for name in names {
346            let state = node.open_db(&name)?;
347            node.dbs
348                .write()
349                .unwrap_or_else(std::sync::PoisonError::into_inner)
350                .insert(name, state);
351        }
352        node.spawn_scheduled_gc();
353        Ok(node)
354    }
355
356    /// The node's data-directory store.
357    #[must_use]
358    pub fn store(&self) -> &Arc<FsStore> {
359        &self.store
360    }
361
362    /// Node configuration.
363    #[must_use]
364    pub fn config(&self) -> &NodeConfig {
365        &self.config
366    }
367
368    /// Process observability counters.
369    #[must_use]
370    pub const fn metrics(&self) -> &Metrics {
371        &self.metrics
372    }
373
374    fn spawn_scheduled_gc(self: &Arc<Self>) {
375        let Some(interval) = self.config.gc_interval else {
376            return;
377        };
378        let Ok(runtime) = tokio::runtime::Handle::try_current() else {
379            // Embedded callers may construct an empty catalog before they
380            // enter a runtime. Process wiring opens nodes inside Tokio.
381            return;
382        };
383        let node = Arc::clone(self);
384        runtime.spawn(async move {
385            let mut ticker = tokio::time::interval(interval);
386            ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
387            // `interval` ticks immediately; scheduled duties should wait a full interval.
388            ticker.tick().await;
389            loop {
390                ticker.tick().await;
391                if let Err(error) = node.gc_deleted().await {
392                    tracing::warn!(%error, "scheduled garbage collection failed");
393                }
394            }
395        });
396    }
397
398    /// Watch channel that reports a shutdown reason when the node deposes.
399    #[must_use]
400    pub fn shutdown_watch(&self) -> watch::Receiver<Option<String>> {
401        self.shutdown.subscribe()
402    }
403
404    fn depose(&self, state: &DbState, reason: &str) {
405        state.deposed.store(true, Ordering::Release);
406        let _ = self
407            .shutdown
408            .send(Some(format!("database {:?}: {reason}", state.name)));
409    }
410
411    fn log_path(&self, name: &str) -> PathBuf {
412        self.config
413            .data_dir
414            .join("logs")
415            .join(format!("{name}.log"))
416    }
417
418    fn acquire_with_wait(&self, name: &str) -> Result<Lease, NodeError> {
419        let deadline = now_unix_ms() + self.config.lease_wait_ms;
420        loop {
421            match lease::acquire(
422                self.store.as_ref(),
423                name,
424                &self.config.owner,
425                self.config.lease_ttl_ms,
426                now_unix_ms(),
427            ) {
428                Ok(held) => return Ok(held),
429                Err(LeaseError::Held { .. }) if now_unix_ms() < deadline => {
430                    std::thread::sleep(Duration::from_millis(200));
431                }
432                Err(error) => return Err(error.into()),
433            }
434        }
435    }
436
437    fn open_db(self: &Arc<Self>, name: &str) -> Result<Arc<DbState>, NodeError> {
438        let meta = self
439            .store
440            .get_root(&meta_root_name(name))?
441            .ok_or_else(|| NodeError::UnknownDb(name.to_owned()))?;
442        let (schema, idents, interner) = decode_meta(&meta)?;
443        let root_name = db_root_name(name);
444        let current = self
445            .store
446            .get_root(&root_name)?
447            .as_deref()
448            .and_then(DbRoot::decode);
449        if let Some(root) = &current {
450            if root.format_version > corium_store::FORMAT_VERSION {
451                return Err(NodeError::UnsupportedFormat {
452                    found: root.format_version,
453                    supported: corium_store::FORMAT_VERSION,
454                });
455            }
456        }
457        let held = self.acquire_with_wait(name)?;
458        // Fence bump: ensure the db root carries our lease version so any
459        // deposed writer's pending CAS fails and observes the new version.
460        if current
461            .as_ref()
462            .is_none_or(|root| root.lease_version < held.version)
463        {
464            publish_root(
465                self.store.as_ref(),
466                &root_name,
467                &DbRoot {
468                    format_version: corium_store::FORMAT_VERSION,
469                    lease_version: held.version,
470                    index_basis_t: current.as_ref().map_or(0, |root| root.index_basis_t),
471                    roots: current.and_then(|root| root.roots),
472                },
473            )?;
474        }
475        let log = Arc::new(FileLog::open(self.log_path(name))?);
476        let base = Db::new(schema.clone()).with_naming(idents.clone(), interner.clone());
477        let transactor = EmbeddedTransactor::recover_from(base, Arc::clone(&log) as _)?;
478        let basis_t = transactor.db().basis_t();
479        let index_basis = self
480            .store
481            .get_root(&root_name)?
482            .as_deref()
483            .and_then(DbRoot::decode)
484            .map_or(0, |root| root.index_basis_t);
485        let state = Arc::new(DbState {
486            name: name.to_owned(),
487            transactor,
488            log,
489            naming: Mutex::new(Naming {
490                schema,
491                idents,
492                interner,
493            }),
494            commit: tokio::sync::Mutex::new(()),
495            broadcast: broadcast::channel(1024).0,
496            basis: watch::channel(basis_t).0,
497            index_basis: AtomicU64::new(index_basis),
498            held_lease: Mutex::new(held),
499            deposed: AtomicBool::new(false),
500        });
501        self.spawn_maintenance(&state);
502        Ok(state)
503    }
504
505    fn spawn_maintenance(self: &Arc<Self>, state: &Arc<DbState>) {
506        let ttl = self.config.lease_ttl_ms;
507        let renew_every = Duration::from_millis(u64::try_from(ttl / 3).unwrap_or(1).max(50));
508        // Lease renewal.
509        let node = Arc::clone(self);
510        let db = Arc::clone(state);
511        tokio::spawn(async move {
512            let mut ticker = tokio::time::interval(renew_every);
513            ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
514            loop {
515                ticker.tick().await;
516                if db.deposed.load(Ordering::Acquire) {
517                    return;
518                }
519                let held = db.lease();
520                let store = Arc::clone(&node.store);
521                let name = db.name.clone();
522                let renewed = tokio::task::spawn_blocking(move || {
523                    lease::renew(store.as_ref(), &name, &held, ttl, now_unix_ms())
524                })
525                .await;
526                match renewed {
527                    Ok(Ok(renewed)) => {
528                        *db.held_lease
529                            .lock()
530                            .unwrap_or_else(std::sync::PoisonError::into_inner) = renewed;
531                    }
532                    Ok(Err(LeaseError::Lost)) => {
533                        node.depose(&db, "write lease lost");
534                        return;
535                    }
536                    Ok(Err(_)) | Err(_) => {}
537                }
538            }
539        });
540        // Background indexing.
541        let node = Arc::clone(self);
542        let db = Arc::clone(state);
543        tokio::spawn(async move {
544            let mut ticker = tokio::time::interval(node.config.index_interval);
545            ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
546            loop {
547                ticker.tick().await;
548                if db.deposed.load(Ordering::Acquire) {
549                    return;
550                }
551                if db.db().basis_t() <= db.index_basis() {
552                    continue;
553                }
554                let _gc = node.gc_lock.lock().await;
555                let store = Arc::clone(&node.store);
556                let version = db.lease().version;
557                let root_name = db_root_name(&db.name);
558                let worker = Arc::clone(&db);
559                let started = Instant::now();
560                let published = tokio::task::spawn_blocking(move || {
561                    worker
562                        .transactor
563                        .publish_indexes(store.as_ref(), &root_name, version)
564                })
565                .await;
566                node.metrics.record_index(started.elapsed());
567                match published {
568                    Ok(Ok(root)) => {
569                        tracing::debug!(db = %db.name, index_basis_t = root.index_basis_t, "published indexes");
570                        db.index_basis.store(root.index_basis_t, Ordering::Release);
571                        let _ = db.broadcast.send(pb::subscribe_item::Item::IndexBasis(
572                            pb::IndexBasis {
573                                index_basis_t: root.index_basis_t,
574                            },
575                        ));
576                    }
577                    Ok(Err(TransactError::Deposed { .. })) => {
578                        node.depose(&db, "database root fenced by a newer lease");
579                        return;
580                    }
581                    Ok(Err(_)) | Err(_) => {}
582                }
583            }
584        });
585        // Heartbeats.
586        let node = Arc::clone(self);
587        let db = Arc::clone(state);
588        tokio::spawn(async move {
589            let mut ticker = tokio::time::interval(node.config.heartbeat_interval);
590            ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
591            loop {
592                ticker.tick().await;
593                if db.deposed.load(Ordering::Acquire) {
594                    return;
595                }
596                let _ = db
597                    .broadcast
598                    .send(pb::subscribe_item::Item::Heartbeat(pb::Heartbeat {
599                        basis_t: db.db().basis_t(),
600                    }));
601            }
602        });
603    }
604
605    /// Looks up a hosted database.
606    ///
607    /// # Errors
608    /// Returns [`NodeError::UnknownDb`] when absent.
609    pub fn db_state(&self, name: &str) -> Result<Arc<DbState>, NodeError> {
610        self.dbs
611            .read()
612            .unwrap_or_else(std::sync::PoisonError::into_inner)
613            .get(name)
614            .cloned()
615            .ok_or_else(|| NodeError::UnknownDb(name.to_owned()))
616    }
617
618    /// Creates a database with the supplied EDN schema forms; returns
619    /// `false` when it already exists.
620    ///
621    /// # Errors
622    /// Returns an error for invalid names/schema or store failures.
623    pub fn create_db(self: &Arc<Self>, name: &str, schema_edn: &[u8]) -> Result<bool, NodeError> {
624        if !valid_db_name(name) {
625            return Err(NodeError::InvalidName(name.to_owned()));
626        }
627        if self
628            .dbs
629            .read()
630            .unwrap_or_else(std::sync::PoisonError::into_inner)
631            .contains_key(name)
632        {
633            return Ok(false);
634        }
635        let forms = match codec::decode_edn(schema_edn)? {
636            Edn::Vector(items) | Edn::List(items) => items,
637            Edn::Nil => Vec::new(),
638            other => {
639                return Err(NodeError::BadRequest(format!(
640                    "schema must be a vector of attribute maps, got {other}"
641                )));
642            }
643        };
644        let (schema, idents) = schema_from_edn(&forms)?;
645        let meta = encode_meta(&schema, &idents, &KeywordInterner::default());
646        match self.store.cas_root(&meta_root_name(name), None, &meta) {
647            Ok(()) => {}
648            Err(StoreError::CasFailed { .. }) => return Ok(false),
649            Err(error) => return Err(error.into()),
650        }
651        let state = self.open_db(name)?;
652        self.dbs
653            .write()
654            .unwrap_or_else(std::sync::PoisonError::into_inner)
655            .insert(name.to_owned(), state);
656        Ok(true)
657    }
658
659    /// Deletes a database: unhosts it, releases its lease, and removes its
660    /// roots and log. Blobs remain until [`Self::gc_deleted`].
661    ///
662    /// # Errors
663    /// Returns an error when roots or the log cannot be removed.
664    pub fn delete_db(&self, name: &str) -> Result<bool, NodeError> {
665        let Some(state) = self
666            .dbs
667            .write()
668            .unwrap_or_else(std::sync::PoisonError::into_inner)
669            .remove(name)
670        else {
671            return Ok(false);
672        };
673        state.deposed.store(true, Ordering::Release);
674        let _ = lease::release(self.store.as_ref(), name, &state.lease());
675        self.store.delete_root(&db_root_name(name))?;
676        self.store.delete_root(&meta_root_name(name))?;
677        self.store.delete_root(&lease::lease_root(name))?;
678        match std::fs::remove_file(self.log_path(name)) {
679            Ok(()) => {}
680            Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
681            Err(error) => return Err(NodeError::Store(StoreError::Io(error))),
682        }
683        Ok(true)
684    }
685
686    /// Lists hosted databases.
687    #[must_use]
688    pub fn list_dbs(&self) -> Vec<String> {
689        let mut names: Vec<String> = self
690            .dbs
691            .read()
692            .unwrap_or_else(std::sync::PoisonError::into_inner)
693            .keys()
694            .cloned()
695            .collect();
696        names.sort();
697        names
698    }
699
700    /// Sweeps blobs unreachable from any live database root (including
701    /// everything left behind by deleted databases and superseded indexes).
702    ///
703    /// # Errors
704    /// Returns an error when the store cannot be enumerated or swept.
705    pub async fn gc_deleted(&self) -> Result<u64, NodeError> {
706        self.gc_deleted_with_retention(self.config.gc_retention)
707            .await
708    }
709
710    /// Sweeps unreachable blobs older than the caller-supplied retention.
711    ///
712    /// # Errors
713    /// Returns an error when the store cannot be enumerated or swept.
714    pub async fn gc_deleted_with_retention(&self, retention: Duration) -> Result<u64, NodeError> {
715        let _gc = self.gc_lock.lock().await;
716        let store = Arc::clone(&self.store);
717        let report =
718            tokio::task::spawn_blocking(move || -> Result<corium_store::GcReport, NodeError> {
719                let mut live = Vec::new();
720                for root_name in store.list_roots("db:")? {
721                    if let Some(root) = store
722                        .get_root(&root_name)?
723                        .as_deref()
724                        .and_then(DbRoot::decode)
725                    {
726                        if let Some(roots) = root.roots {
727                            live.extend(roots);
728                        }
729                    }
730                }
731                Ok(mark_and_sweep_retained(
732                    store.as_ref(),
733                    live,
734                    |_, _| Ok(Vec::new()),
735                    retention,
736                    SystemTime::now(),
737                )?)
738            })
739            .await
740            .map_err(|error| NodeError::BadRequest(format!("gc task failed: {error}")))??;
741        self.metrics
742            .record_gc(report.swept as u64, report.retained as u64);
743        tracing::info!(
744            marked = report.marked,
745            swept = report.swept,
746            retained = report.retained,
747            "garbage collection completed"
748        );
749        Ok(report.swept as u64)
750    }
751
752    /// Validates, appends, applies, and reports one transaction supplied as
753    /// composite-encoded EDN transaction forms.
754    ///
755    /// # Errors
756    /// Returns [`NodeError`] for decode/validation failures, lease loss, or
757    /// storage failures.
758    pub async fn transact(
759        &self,
760        name: &str,
761        tx_data: &[u8],
762    ) -> Result<pb::TransactResponse, NodeError> {
763        let started = Instant::now();
764        let result = self
765            .transact_inner(name, tx_data)
766            .instrument(tracing::info_span!("transact", db = name))
767            .await;
768        self.metrics.record_tx(started.elapsed(), result.is_ok());
769        if let Err(error) = &result {
770            tracing::warn!(%error, "transaction failed");
771        }
772        result
773    }
774
775    async fn transact_inner(
776        &self,
777        name: &str,
778        tx_data: &[u8],
779    ) -> Result<pb::TransactResponse, NodeError> {
780        let state = self.db_state(name)?;
781        let decoded = codec::decode_edn(tx_data)?;
782        let forms = decoded
783            .as_seq()
784            .ok_or_else(|| NodeError::BadRequest("tx-data must be a vector".into()))?
785            .to_vec();
786        let queued = self.metrics.queue_waiter();
787        let commit = state.commit.lock().await;
788        drop(queued);
789        let _commit = commit;
790        state.check_lease(self.store.as_ref())?;
791        // Expand user database-function invocations against the
792        // db-in-transaction (the value under the commit lock) before native
793        // conversion. The expander blocks up to its budget deadline, so it
794        // runs off the async workers.
795        let forms = if let Some(expander) = &self.config.tx_fn_expander {
796            let expander = Arc::clone(expander);
797            let db = state.transactor.db();
798            tokio::task::spawn_blocking(move || expander.expand(&db, forms))
799                .await
800                .map_err(|error| NodeError::BadRequest(format!("expander task failed: {error}")))?
801                .map_err(NodeError::BadRequest)?
802        } else {
803            forms
804        };
805        // Convert forms, interning any new keyword values.
806        let (items, naming_changed, idents, interner, schema) = {
807            let mut naming = state
808                .naming
809                .lock()
810                .unwrap_or_else(std::sync::PoisonError::into_inner);
811            let db = state.transactor.db();
812            let before = naming.interner.len();
813            let items = tx_items_from_edn(&db, &mut naming.interner, &forms)?;
814            (
815                items,
816                naming.interner.len() > before,
817                naming.idents.clone(),
818                naming.interner.clone(),
819                naming.schema.clone(),
820            )
821        };
822        if naming_changed {
823            // New keyword names must be durable before the datoms that
824            // reference them; recovery decodes the log against this meta.
825            let meta = encode_meta(&schema, &idents, &interner);
826            loop {
827                let current = self.store.get_root(&meta_root_name(name))?;
828                match self
829                    .store
830                    .cas_root(&meta_root_name(name), current.as_deref(), &meta)
831                {
832                    Ok(()) => break,
833                    Err(StoreError::CasFailed { .. }) => {}
834                    Err(error) => return Err(error.into()),
835                }
836            }
837            state.transactor.update_naming(idents, interner.clone());
838        }
839        let worker = Arc::clone(&state);
840        let report = tokio::task::spawn_blocking(move || worker.transactor.transact(items))
841            .await
842            .map_err(|error| NodeError::BadRequest(format!("transact task failed: {error}")))??;
843        let t = report.db_after.basis_t();
844        let datoms = codec::encode_datoms(&report.tx.datoms, &interner)?;
845        let tempids = codec::encode_edn(&Edn::Map(
846            report
847                .tx
848                .tempids
849                .iter()
850                .map(|(tempid, eid)| {
851                    (
852                        Edn::Str(tempid.clone()),
853                        Edn::Long(i64::try_from(eid.raw()).unwrap_or(i64::MAX)),
854                    )
855                })
856                .collect(),
857        ));
858        let _ = state
859            .broadcast
860            .send(pb::subscribe_item::Item::Report(pb::TxReport {
861                t,
862                tx_instant: report.tx_instant,
863                datoms: datoms.clone(),
864            }));
865        let _ = state.basis.send(t);
866        Ok(pb::TransactResponse {
867            basis_before: report.db_before.basis_t(),
868            basis_t: t,
869            tx_instant: report.tx_instant,
870            tempids,
871            tx_data: datoms,
872        })
873    }
874
875    /// Current status for a database.
876    ///
877    /// # Errors
878    /// Returns [`NodeError::UnknownDb`] when absent.
879    pub fn status(&self, name: &str) -> Result<pb::StatusResponse, NodeError> {
880        let state = self.db_state(name)?;
881        let db = state.db();
882        let counts = db.stats();
883        let held = state.lease();
884        let metrics = self.metrics.snapshot();
885        Ok(pb::StatusResponse {
886            basis_t: db.basis_t(),
887            index_basis_t: state.index_basis(),
888            lease_owner: held.owner,
889            lease_version: held.version,
890            lease_expires_unix_ms: held.expires_unix_ms,
891            datom_count: counts.datoms as u64,
892            entity_count: counts.entities as u64,
893            attribute_count: counts.attributes as u64,
894            transaction_count: metrics.tx_total,
895            transaction_failure_count: metrics.tx_failed,
896            transaction_queue_depth: metrics.queue_depth,
897            index_lag: db.basis_t().saturating_sub(state.index_basis()),
898            indexing_runs: metrics.index_runs,
899            gc_runs: metrics.gc_runs,
900            gc_swept_blobs: metrics.gc_swept,
901        })
902    }
903
904    /// Waits until the database basis reaches `t`, returning the basis seen.
905    ///
906    /// # Errors
907    /// Returns [`NodeError::UnknownDb`] when absent.
908    pub async fn sync(&self, name: &str, t: u64) -> Result<u64, NodeError> {
909        let state = self.db_state(name)?;
910        let mut basis = state.basis_watch();
911        let target = if t == 0 { *basis.borrow() } else { t };
912        loop {
913            let current = *basis.borrow();
914            if current >= target {
915                return Ok(current);
916            }
917            if basis.changed().await.is_err() {
918                return Ok(*basis.borrow());
919            }
920        }
921    }
922}