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, tx-report fan-out, and
3//! high-availability standby takeover.
4
5use std::collections::{BTreeSet, HashMap};
6use std::path::PathBuf;
7use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
8use std::sync::{Arc, Mutex};
9use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
10
11use corium_core::{KeywordInterner, Schema};
12use corium_db::{Db, Idents};
13use corium_log::{LogError, TransactionLog, TxRecord};
14use corium_protocol::codec::{self, CodecError};
15use corium_protocol::pb;
16use corium_protocol::schemaform::{SchemaFormError, schema_from_edn};
17use corium_protocol::txforms::{TxFormError, tx_items_from_edn};
18use corium_query::edn::Edn;
19use corium_store::{RootStore, StoreError, mark_and_sweep_retained};
20use thiserror::Error;
21use tokio::sync::{broadcast, watch};
22use tracing::Instrument;
23
24use crate::backend::{LogBackend, NodeStore, StoreSpec};
25use crate::lease::{self, Lease, LeaseError};
26use crate::metrics::Metrics;
27use crate::{DbRoot, EmbeddedTransactor, TransactError, db_root_name};
28
29/// Expands user database-function invocations in boundary EDN transaction
30/// forms before native conversion. Implemented by `corium-cljrs` (the
31/// sandboxed Clojurust host, ADR-0008) and injected by the process wiring;
32/// the transactor itself stays free of cljrs dependencies.
33pub trait TxFnExpander: Send + Sync {
34    /// Rewrites `forms` with every `[:my/fn arg…]` invocation replaced by
35    /// the function's returned tx-data (recursively).
36    ///
37    /// # Errors
38    /// Returns a display message when a function is missing, rejected by
39    /// the sandbox, fails, or exceeds its budget; the transaction aborts.
40    fn expand(&self, db: &Db, forms: Vec<Edn>) -> Result<Vec<Edn>, String>;
41}
42
43/// Node process configuration.
44#[derive(Clone)]
45pub struct NodeConfig {
46    /// Storage-service backend for blobs and roots (`mem`, `fs`, or Turso).
47    pub store: StoreSpec,
48    /// Data directory holding the filesystem blob/root store (for the `fs`
49    /// backend) and the transaction logs (for every non-`mem` backend).
50    pub data_dir: PathBuf,
51    /// Stable owner identity for lease records.
52    pub owner: String,
53    /// Lease time-to-live in milliseconds.
54    pub lease_ttl_ms: i64,
55    /// How long to wait for a held lease to expire before giving up.
56    pub lease_wait_ms: i64,
57    /// High-availability mode: when another owner holds a database's lease,
58    /// stand by and take over on expiry instead of failing startup, and on
59    /// depose return to standby instead of shutting the process down.
60    pub ha: bool,
61    /// Client endpoint advertised in the lease for peer lease-holder
62    /// rediscovery (e.g. `http://transactor-a:4334`).
63    pub advertise: Option<String>,
64    /// Interval between background index publications.
65    pub index_interval: Duration,
66    /// Interval between heartbeats on subscription streams.
67    pub heartbeat_interval: Duration,
68    /// Interval between scheduled garbage-collection duties; `None` disables it.
69    pub gc_interval: Option<Duration>,
70    /// Minimum age of an unreachable blob before scheduled/manual online GC.
71    pub gc_retention: Duration,
72    /// Optional database-function expander (`:db/fn` support).
73    pub tx_fn_expander: Option<Arc<dyn TxFnExpander>>,
74}
75
76impl std::fmt::Debug for NodeConfig {
77    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
78        f.debug_struct("NodeConfig")
79            .field("store", &self.store)
80            .field("data_dir", &self.data_dir)
81            .field("owner", &self.owner)
82            .field("lease_ttl_ms", &self.lease_ttl_ms)
83            .field("lease_wait_ms", &self.lease_wait_ms)
84            .field("ha", &self.ha)
85            .field("advertise", &self.advertise)
86            .field("index_interval", &self.index_interval)
87            .field("heartbeat_interval", &self.heartbeat_interval)
88            .field("gc_interval", &self.gc_interval)
89            .field("gc_retention", &self.gc_retention)
90            .field("tx_fn_expander", &self.tx_fn_expander.is_some())
91            .finish()
92    }
93}
94
95impl NodeConfig {
96    /// Sensible defaults for a data directory.
97    #[must_use]
98    pub fn new(data_dir: PathBuf) -> Self {
99        Self {
100            store: StoreSpec::Fs,
101            data_dir,
102            owner: format!(
103                "transactor-{}",
104                std::env::var("HOSTNAME").unwrap_or_else(|_| "local".into())
105            ),
106            lease_ttl_ms: 5_000,
107            lease_wait_ms: 15_000,
108            ha: false,
109            advertise: None,
110            index_interval: Duration::from_secs(5),
111            heartbeat_interval: Duration::from_secs(10),
112            gc_interval: Some(Duration::from_secs(60 * 60)),
113            gc_retention: Duration::from_secs(72 * 60 * 60),
114            tx_fn_expander: None,
115        }
116    }
117}
118
119/// Node operation failure.
120#[derive(Debug, Error)]
121pub enum NodeError {
122    /// Named database does not exist.
123    #[error("unknown database {0:?}")]
124    UnknownDb(String),
125    /// Database name is not storable.
126    #[error("invalid database name {0:?}")]
127    InvalidName(String),
128    /// Database root uses a storage format newer than this binary.
129    #[error("storage format {found} is newer than supported format {supported}")]
130    UnsupportedFormat {
131        /// Version found in the root.
132        found: u32,
133        /// Newest version understood by this binary.
134        supported: u32,
135    },
136    /// This node no longer holds the write lease.
137    #[error("deposed: write lease for {0:?} is held elsewhere")]
138    Deposed(String),
139    /// This node is a warm standby for the database; the lease holder
140    /// serves it.
141    #[error("standby for {db:?}: lease held by {owner} at {endpoint:?}")]
142    Standby {
143        /// Database name.
144        db: String,
145        /// Current lease owner id (empty when unknown).
146        owner: String,
147        /// Owner's advertised client endpoint (empty when unadvertised).
148        endpoint: String,
149    },
150    /// Payload failed to decode.
151    #[error(transparent)]
152    Codec(#[from] CodecError),
153    /// Transaction forms failed to convert.
154    #[error(transparent)]
155    TxForm(#[from] TxFormError),
156    /// Schema forms failed to convert.
157    #[error(transparent)]
158    SchemaForm(#[from] SchemaFormError),
159    /// Transaction pipeline failure.
160    #[error(transparent)]
161    Transact(#[from] TransactError),
162    /// Store failure.
163    #[error(transparent)]
164    Store(#[from] StoreError),
165    /// Log failure.
166    #[error(transparent)]
167    Log(#[from] LogError),
168    /// Lease failure.
169    #[error(transparent)]
170    Lease(#[from] LeaseError),
171    /// Malformed request.
172    #[error("bad request: {0}")]
173    BadRequest(String),
174}
175
176struct Naming {
177    schema: Schema,
178    idents: Idents,
179    interner: KeywordInterner,
180}
181
182/// Per-database state hosted by a node.
183pub struct DbState {
184    name: String,
185    transactor: EmbeddedTransactor,
186    log: Arc<dyn TransactionLog>,
187    naming: Mutex<Naming>,
188    commit: tokio::sync::Mutex<()>,
189    broadcast: broadcast::Sender<pb::subscribe_item::Item>,
190    basis: watch::Sender<u64>,
191    index_basis: AtomicU64,
192    held_lease: Mutex<Lease>,
193    deposed: AtomicBool,
194}
195
196impl DbState {
197    /// Database name.
198    #[must_use]
199    pub fn name(&self) -> &str {
200        &self.name
201    }
202
203    /// Current database value.
204    #[must_use]
205    pub fn db(&self) -> Db {
206        self.transactor.db()
207    }
208
209    /// Watch channel following the commit basis.
210    #[must_use]
211    pub fn basis_watch(&self) -> watch::Receiver<u64> {
212        self.basis.subscribe()
213    }
214
215    /// Subscribes to live stream items (reports, index announcements,
216    /// heartbeats).
217    #[must_use]
218    pub fn stream_items(&self) -> broadcast::Receiver<pb::subscribe_item::Item> {
219        self.broadcast.subscribe()
220    }
221
222    /// Basis of the newest published index root.
223    #[must_use]
224    pub fn index_basis(&self) -> u64 {
225        self.index_basis.load(Ordering::Acquire)
226    }
227
228    /// Currently held lease record.
229    #[must_use]
230    pub fn lease(&self) -> Lease {
231        self.held_lease
232            .lock()
233            .unwrap_or_else(std::sync::PoisonError::into_inner)
234            .clone()
235    }
236
237    /// Encoded schema/ident handshake payload plus a consistent basis and
238    /// interner snapshot for backfill encoding.
239    #[must_use]
240    pub fn handshake_snapshot(&self) -> (Vec<u8>, KeywordInterner) {
241        let naming = self
242            .naming
243            .lock()
244            .unwrap_or_else(std::sync::PoisonError::into_inner);
245        (
246            codec::encode_schema(&naming.schema, &naming.idents),
247            naming.interner.clone(),
248        )
249    }
250
251    /// Reads committed records in `[start, end)` from the durable log.
252    ///
253    /// # Errors
254    /// Returns an error when the log cannot be read.
255    pub fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, NodeError> {
256        Ok(self.log.tx_range(start, end)?)
257    }
258
259    /// Verifies this node still owns the write lease (identity check on
260    /// the root record; expiry changes from renewals do not matter).
261    async fn check_lease(&self, store: &dyn RootStore) -> Result<Lease, NodeError> {
262        if self.deposed.load(Ordering::Acquire) {
263            return Err(NodeError::Deposed(self.name.clone()));
264        }
265        let held = self.lease();
266        match lease::verify(store, &self.name, &held).await {
267            Ok(()) => Ok(held),
268            Err(LeaseError::Lost) => {
269                self.deposed.store(true, Ordering::Release);
270                Err(NodeError::Deposed(self.name.clone()))
271            }
272            Err(error) => Err(error.into()),
273        }
274    }
275}
276
277/// A running transactor node hosting every database under one data directory.
278pub struct TransactorNode {
279    config: NodeConfig,
280    store: Arc<NodeStore>,
281    log_backend: LogBackend,
282    dbs: std::sync::RwLock<HashMap<String, Arc<DbState>>>,
283    /// Databases this node is standing by for (HA mode): the lease is held
284    /// elsewhere and the standby poller attempts takeover on expiry.
285    standby: std::sync::RwLock<BTreeSet<String>>,
286    gc_lock: tokio::sync::Mutex<()>,
287    metrics: Metrics,
288    shutdown: watch::Sender<Option<String>>,
289}
290
291fn now_unix_ms() -> i64 {
292    i64::try_from(
293        SystemTime::now()
294            .duration_since(UNIX_EPOCH)
295            .unwrap_or_default()
296            .as_millis(),
297    )
298    .unwrap_or(i64::MAX)
299}
300
301fn valid_db_name(name: &str) -> bool {
302    !name.is_empty()
303        && name.len() <= 128
304        && name
305            .bytes()
306            .all(|byte| byte.is_ascii_alphanumeric() || byte == b'-' || byte == b'_')
307}
308
309fn meta_root_name(db: &str) -> String {
310    format!("meta:{db}")
311}
312
313fn encode_meta(schema: &Schema, idents: &Idents, interner: &KeywordInterner) -> Vec<u8> {
314    let schema_bytes = codec::encode_schema(schema, idents);
315    let naming_bytes = codec::encode_naming(interner);
316    let mut out = Vec::with_capacity(8 + schema_bytes.len() + naming_bytes.len());
317    out.extend_from_slice(&u32::try_from(schema_bytes.len()).unwrap_or(0).to_be_bytes());
318    out.extend_from_slice(&schema_bytes);
319    out.extend_from_slice(&u32::try_from(naming_bytes.len()).unwrap_or(0).to_be_bytes());
320    out.extend_from_slice(&naming_bytes);
321    out
322}
323
324fn decode_meta(bytes: &[u8]) -> Result<(Schema, Idents, KeywordInterner), NodeError> {
325    let take = |input: &mut &[u8]| -> Result<Vec<u8>, NodeError> {
326        let len_bytes = input
327            .get(..4)
328            .ok_or(NodeError::Codec(CodecError::Truncated))?;
329        let len = usize::try_from(u32::from_be_bytes(len_bytes.try_into().unwrap_or_default()))
330            .map_err(|_| NodeError::Codec(CodecError::Length))?;
331        let payload = input
332            .get(4..4 + len)
333            .ok_or(NodeError::Codec(CodecError::Truncated))?
334            .to_vec();
335        *input = &input[4 + len..];
336        Ok(payload)
337    };
338    let mut input = bytes;
339    let schema_bytes = take(&mut input)?;
340    let naming_bytes = take(&mut input)?;
341    let (schema, idents) = codec::decode_schema(&schema_bytes)?;
342    let interner = codec::decode_naming(&naming_bytes)?;
343    Ok((schema, idents, interner))
344}
345
346impl TransactorNode {
347    /// Opens a node over `config.data_dir`, recovering every database found
348    /// there (acquiring its lease, waiting out held leases up to the
349    /// configured bound).
350    ///
351    /// # Errors
352    /// Returns an error when the store cannot be opened or a database cannot
353    /// be recovered.
354    pub async fn open(config: NodeConfig) -> Result<Arc<Self>, NodeError> {
355        let store = Arc::new(NodeStore::open(&config.store, &config.data_dir).await?);
356        let log_backend = LogBackend::for_spec(&config.store, &config.data_dir);
357        let node = Arc::new(Self {
358            config,
359            store,
360            log_backend,
361            dbs: std::sync::RwLock::new(HashMap::new()),
362            standby: std::sync::RwLock::new(BTreeSet::new()),
363            gc_lock: tokio::sync::Mutex::new(()),
364            metrics: Metrics::default(),
365            shutdown: watch::channel(None).0,
366        });
367        let names: Vec<String> = node
368            .store
369            .list_roots("meta:")
370            .await?
371            .into_iter()
372            .filter_map(|root| root.strip_prefix("meta:").map(str::to_owned))
373            .collect();
374        for name in names {
375            match node.open_db(&name).await {
376                Ok(state) => {
377                    node.dbs
378                        .write()
379                        .unwrap_or_else(std::sync::PoisonError::into_inner)
380                        .insert(name, state);
381                }
382                Err(NodeError::Lease(LeaseError::Held { owner, .. })) if node.config.ha => {
383                    tracing::info!(db = %name, %owner, "standing by; lease held elsewhere");
384                    node.standby
385                        .write()
386                        .unwrap_or_else(std::sync::PoisonError::into_inner)
387                        .insert(name);
388                }
389                Err(error) => return Err(error),
390            }
391        }
392        node.spawn_standby_poller();
393        node.spawn_scheduled_gc();
394        Ok(node)
395    }
396
397    /// The node's storage-service backend (blobs + roots).
398    #[must_use]
399    pub fn store(&self) -> &Arc<NodeStore> {
400        &self.store
401    }
402
403    /// Node configuration.
404    #[must_use]
405    pub fn config(&self) -> &NodeConfig {
406        &self.config
407    }
408
409    /// Process observability counters.
410    #[must_use]
411    pub const fn metrics(&self) -> &Metrics {
412        &self.metrics
413    }
414
415    fn spawn_scheduled_gc(self: &Arc<Self>) {
416        let Some(interval) = self.config.gc_interval else {
417            return;
418        };
419        let Ok(runtime) = tokio::runtime::Handle::try_current() else {
420            // Embedded callers may construct an empty catalog before they
421            // enter a runtime. Process wiring opens nodes inside Tokio.
422            return;
423        };
424        let node = Arc::clone(self);
425        runtime.spawn(async move {
426            let mut ticker = tokio::time::interval(interval);
427            ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
428            // `interval` ticks immediately; scheduled duties should wait a full interval.
429            ticker.tick().await;
430            loop {
431                ticker.tick().await;
432                if let Err(error) = node.gc_deleted().await {
433                    tracing::warn!(%error, "scheduled garbage collection failed");
434                }
435            }
436        });
437    }
438
439    /// Watch channel that reports a shutdown reason when the node deposes.
440    #[must_use]
441    pub fn shutdown_watch(&self) -> watch::Receiver<Option<String>> {
442        self.shutdown.subscribe()
443    }
444
445    /// Deposes a hosted database. In HA mode the database returns to
446    /// standby (the poller re-attempts takeover); otherwise the whole
447    /// process shuts down and a supervisor restart re-acquires or waits.
448    fn depose(&self, state: &DbState, reason: &str) {
449        state.deposed.store(true, Ordering::Release);
450        if self.config.ha {
451            tracing::warn!(db = %state.name, reason, "deposed; returning to standby");
452            self.dbs
453                .write()
454                .unwrap_or_else(std::sync::PoisonError::into_inner)
455                .remove(&state.name);
456            self.standby
457                .write()
458                .unwrap_or_else(std::sync::PoisonError::into_inner)
459                .insert(state.name.clone());
460        } else {
461            let _ = self
462                .shutdown
463                .send(Some(format!("database {:?}: {reason}", state.name)));
464        }
465    }
466
467    fn advertised(&self) -> &str {
468        self.config.advertise.as_deref().unwrap_or("")
469    }
470
471    /// Acquires the lease for `name`. In HA mode a held lease surfaces
472    /// immediately (the caller stands by); otherwise startup waits it out
473    /// up to the configured bound.
474    async fn acquire_lease(&self, name: &str) -> Result<Lease, NodeError> {
475        let deadline = now_unix_ms() + self.config.lease_wait_ms;
476        loop {
477            match lease::acquire(
478                self.store.as_ref(),
479                name,
480                &self.config.owner,
481                self.advertised(),
482                self.config.lease_ttl_ms,
483                now_unix_ms(),
484            )
485            .await
486            {
487                Ok(held) => return Ok(held),
488                Err(LeaseError::Held { .. }) if !self.config.ha && now_unix_ms() < deadline => {
489                    tokio::time::sleep(Duration::from_millis(200)).await;
490                }
491                Err(error) => return Err(error.into()),
492            }
493        }
494    }
495
496    async fn open_db(self: &Arc<Self>, name: &str) -> Result<Arc<DbState>, NodeError> {
497        let meta = self
498            .store
499            .get_root(&meta_root_name(name))
500            .await?
501            .ok_or_else(|| NodeError::UnknownDb(name.to_owned()))?;
502        let (schema, idents, interner) = decode_meta(&meta)?;
503        let root_name = db_root_name(name);
504        let current = self
505            .store
506            .get_root(&root_name)
507            .await?
508            .as_deref()
509            .and_then(DbRoot::decode);
510        if let Some(root) = &current
511            && root.format_version > corium_store::FORMAT_VERSION
512        {
513            return Err(NodeError::UnsupportedFormat {
514                found: root.format_version,
515                supported: corium_store::FORMAT_VERSION,
516            });
517        }
518        // Acquisition rewrites the root record under our lease version, so
519        // it doubles as the fence bump: a deposed writer's pending root CAS
520        // now has stale expected bytes and must fail.
521        let held = self.acquire_lease(name).await?;
522        // The log tail replay below happens strictly after the fence, so it
523        // observes every record a previous owner could ever have acked.
524        let log = self.log_backend.open(name, held.version)?;
525        let base = Db::new(schema.clone()).with_naming(idents.clone(), interner.clone());
526        let transactor = EmbeddedTransactor::recover_from(base, Arc::clone(&log))?;
527        let basis_t = transactor.db().basis_t();
528        let index_basis = self
529            .store
530            .get_root(&root_name)
531            .await?
532            .as_deref()
533            .and_then(DbRoot::decode)
534            .map_or(0, |root| root.index_basis_t);
535        let state = Arc::new(DbState {
536            name: name.to_owned(),
537            transactor,
538            log,
539            naming: Mutex::new(Naming {
540                schema,
541                idents,
542                interner,
543            }),
544            commit: tokio::sync::Mutex::new(()),
545            broadcast: broadcast::channel(1024).0,
546            basis: watch::channel(basis_t).0,
547            index_basis: AtomicU64::new(index_basis),
548            held_lease: Mutex::new(held),
549            deposed: AtomicBool::new(false),
550        });
551        self.spawn_maintenance(&state);
552        Ok(state)
553    }
554
555    /// HA standby duty: at the lease-renewal cadence, rediscover databases
556    /// (including ones created on the active after this process started)
557    /// and attempt takeover of any whose lease has lapsed. Takeover is
558    /// ordinary startup — acquire (which fences), replay the log tail,
559    /// serve — per the crash-only design.
560    fn spawn_standby_poller(self: &Arc<Self>) {
561        if !self.config.ha {
562            return;
563        }
564        let Ok(runtime) = tokio::runtime::Handle::try_current() else {
565            return;
566        };
567        let ttl = self.config.lease_ttl_ms;
568        let poll_every = Duration::from_millis(u64::try_from(ttl / 3).unwrap_or(1).max(50));
569        let node = Arc::clone(self);
570        runtime.spawn(async move {
571            let mut ticker = tokio::time::interval(poll_every);
572            ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
573            loop {
574                ticker.tick().await;
575                if let Err(error) = node.standby_scan().await {
576                    tracing::warn!(%error, "standby scan failed");
577                }
578            }
579        });
580    }
581
582    /// One standby pass: refresh the standby set from the catalog and try
583    /// to take over lapsed leases.
584    async fn standby_scan(self: &Arc<Self>) -> Result<(), NodeError> {
585        let names: Vec<String> = self
586            .store
587            .list_roots("meta:")
588            .await?
589            .into_iter()
590            .filter_map(|root| root.strip_prefix("meta:").map(str::to_owned))
591            .collect();
592        {
593            let mut standby = self
594                .standby
595                .write()
596                .unwrap_or_else(std::sync::PoisonError::into_inner);
597            standby.retain(|name| names.contains(name));
598        }
599        for name in names {
600            if self
601                .dbs
602                .read()
603                .unwrap_or_else(std::sync::PoisonError::into_inner)
604                .contains_key(&name)
605            {
606                continue;
607            }
608            match self.open_db(&name).await {
609                Ok(state) => {
610                    tracing::info!(db = %name, owner = %self.config.owner, "standby took over write lease");
611                    self.standby
612                        .write()
613                        .unwrap_or_else(std::sync::PoisonError::into_inner)
614                        .remove(&name);
615                    self.dbs
616                        .write()
617                        .unwrap_or_else(std::sync::PoisonError::into_inner)
618                        .insert(name, state);
619                }
620                Err(NodeError::Lease(LeaseError::Held { .. })) => {
621                    self.standby
622                        .write()
623                        .unwrap_or_else(std::sync::PoisonError::into_inner)
624                        .insert(name);
625                }
626                Err(error) => {
627                    tracing::warn!(db = %name, %error, "standby takeover attempt failed");
628                }
629            }
630        }
631        Ok(())
632    }
633
634    fn spawn_maintenance(self: &Arc<Self>, state: &Arc<DbState>) {
635        let ttl = self.config.lease_ttl_ms;
636        let renew_every = Duration::from_millis(u64::try_from(ttl / 3).unwrap_or(1).max(50));
637        // Lease renewal.
638        let node = Arc::clone(self);
639        let db = Arc::clone(state);
640        tokio::spawn(async move {
641            let mut ticker = tokio::time::interval(renew_every);
642            ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
643            loop {
644                ticker.tick().await;
645                if db.deposed.load(Ordering::Acquire) {
646                    return;
647                }
648                // Serialize the root update and local held-lease update with
649                // transaction lease checks so they cannot observe different
650                // renewal generations and falsely depose this node.
651                let _commit = db.commit.lock().await;
652                let held = db.lease();
653                let name = db.name.clone();
654                let renewed =
655                    lease::renew(node.store.as_ref(), &name, &held, ttl, now_unix_ms()).await;
656                match renewed {
657                    Ok(renewed) => {
658                        *db.held_lease
659                            .lock()
660                            .unwrap_or_else(std::sync::PoisonError::into_inner) = renewed;
661                    }
662                    Err(LeaseError::Lost) => {
663                        node.depose(&db, "write lease lost");
664                        return;
665                    }
666                    Err(_) => {}
667                }
668            }
669        });
670        // Background indexing.
671        let node = Arc::clone(self);
672        let db = Arc::clone(state);
673        tokio::spawn(async move {
674            let mut ticker = tokio::time::interval(node.config.index_interval);
675            ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
676            loop {
677                ticker.tick().await;
678                if db.deposed.load(Ordering::Acquire) {
679                    return;
680                }
681                if db.db().basis_t() <= db.index_basis() {
682                    continue;
683                }
684                let _gc = node.gc_lock.lock().await;
685                let version = db.lease().version;
686                let root_name = db_root_name(&db.name);
687                let started = Instant::now();
688                let published = db
689                    .transactor
690                    .publish_indexes(node.store.as_ref(), &root_name, version)
691                    .await;
692                node.metrics.record_index(started.elapsed());
693                match published {
694                    Ok(root) => {
695                        tracing::debug!(db = %db.name, index_basis_t = root.index_basis_t, "published indexes");
696                        db.index_basis.store(root.index_basis_t, Ordering::Release);
697                        let _ = db.broadcast.send(pb::subscribe_item::Item::IndexBasis(
698                            pb::IndexBasis {
699                                index_basis_t: root.index_basis_t,
700                            },
701                        ));
702                    }
703                    Err(TransactError::Deposed { .. }) => {
704                        node.depose(&db, "database root fenced by a newer lease");
705                        return;
706                    }
707                    Err(_) => {}
708                }
709            }
710        });
711        // Heartbeats.
712        let node = Arc::clone(self);
713        let db = Arc::clone(state);
714        tokio::spawn(async move {
715            let mut ticker = tokio::time::interval(node.config.heartbeat_interval);
716            ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
717            loop {
718                ticker.tick().await;
719                if db.deposed.load(Ordering::Acquire) {
720                    return;
721                }
722                let _ = db
723                    .broadcast
724                    .send(pb::subscribe_item::Item::Heartbeat(pb::Heartbeat {
725                        basis_t: db.db().basis_t(),
726                    }));
727            }
728        });
729    }
730
731    /// Looks up a hosted database.
732    ///
733    /// # Errors
734    /// Returns [`NodeError::Standby`] when this HA node is standing by for
735    /// the database, [`NodeError::UnknownDb`] when absent.
736    pub async fn db_state(&self, name: &str) -> Result<Arc<DbState>, NodeError> {
737        if let Some(state) = self
738            .dbs
739            .read()
740            .unwrap_or_else(std::sync::PoisonError::into_inner)
741            .get(name)
742            .cloned()
743        {
744            return Ok(state);
745        }
746        if self.config.ha
747            && self
748                .standby
749                .read()
750                .unwrap_or_else(std::sync::PoisonError::into_inner)
751                .contains(name)
752        {
753            let root = self
754                .store
755                .get_root(&db_root_name(name))
756                .await?
757                .as_deref()
758                .and_then(DbRoot::decode);
759            return Err(NodeError::Standby {
760                db: name.to_owned(),
761                owner: root.as_ref().map(|r| r.owner.clone()).unwrap_or_default(),
762                endpoint: root.map(|r| r.owner_endpoint).unwrap_or_default(),
763            });
764        }
765        Err(NodeError::UnknownDb(name.to_owned()))
766    }
767
768    /// Databases this node currently stands by for (HA mode).
769    #[must_use]
770    pub fn standby_dbs(&self) -> Vec<String> {
771        self.standby
772            .read()
773            .unwrap_or_else(std::sync::PoisonError::into_inner)
774            .iter()
775            .cloned()
776            .collect()
777    }
778
779    /// Creates a database with the supplied EDN schema forms; returns
780    /// `false` when it already exists.
781    ///
782    /// # Errors
783    /// Returns an error for invalid names/schema or store failures.
784    pub async fn create_db(
785        self: &Arc<Self>,
786        name: &str,
787        schema_edn: &[u8],
788    ) -> Result<bool, NodeError> {
789        if !valid_db_name(name) {
790            return Err(NodeError::InvalidName(name.to_owned()));
791        }
792        if self
793            .dbs
794            .read()
795            .unwrap_or_else(std::sync::PoisonError::into_inner)
796            .contains_key(name)
797        {
798            return Ok(false);
799        }
800        let forms = match codec::decode_edn(schema_edn)? {
801            Edn::Vector(items) | Edn::List(items) => items,
802            Edn::Nil => Vec::new(),
803            other => {
804                return Err(NodeError::BadRequest(format!(
805                    "schema must be a vector of attribute maps, got {other}"
806                )));
807            }
808        };
809        let (schema, idents) = schema_from_edn(&forms)?;
810        let meta = encode_meta(&schema, &idents, &KeywordInterner::default());
811        match self
812            .store
813            .cas_root(&meta_root_name(name), None, &meta)
814            .await
815        {
816            Ok(()) => {}
817            Err(StoreError::CasFailed { .. }) => return Ok(false),
818            Err(error) => return Err(error.into()),
819        }
820        let state = self.open_db(name).await?;
821        self.dbs
822            .write()
823            .unwrap_or_else(std::sync::PoisonError::into_inner)
824            .insert(name.to_owned(), state);
825        Ok(true)
826    }
827
828    /// Deletes a database: unhosts it, releases its lease, and removes its
829    /// roots and log. Blobs remain until [`Self::gc_deleted`].
830    ///
831    /// # Errors
832    /// Returns an error when roots or the log cannot be removed.
833    pub async fn delete_db(&self, name: &str) -> Result<bool, NodeError> {
834        let Some(state) = self
835            .dbs
836            .write()
837            .unwrap_or_else(std::sync::PoisonError::into_inner)
838            .remove(name)
839        else {
840            return Ok(false);
841        };
842        state.deposed.store(true, Ordering::Release);
843        self.standby
844            .write()
845            .unwrap_or_else(std::sync::PoisonError::into_inner)
846            .remove(name);
847        self.store.delete_root(&db_root_name(name)).await?;
848        self.store.delete_root(&meta_root_name(name)).await?;
849        self.log_backend.delete_all(name)?;
850        Ok(true)
851    }
852
853    /// Lists hosted databases.
854    #[must_use]
855    pub fn list_dbs(&self) -> Vec<String> {
856        let mut names: Vec<String> = self
857            .dbs
858            .read()
859            .unwrap_or_else(std::sync::PoisonError::into_inner)
860            .keys()
861            .cloned()
862            .collect();
863        names.sort();
864        names
865    }
866
867    /// Sweeps blobs unreachable from any live database root (including
868    /// everything left behind by deleted databases and superseded indexes).
869    ///
870    /// # Errors
871    /// Returns an error when the store cannot be enumerated or swept.
872    pub async fn gc_deleted(&self) -> Result<u64, NodeError> {
873        self.gc_deleted_with_retention(self.config.gc_retention)
874            .await
875    }
876
877    /// Sweeps unreachable blobs older than the caller-supplied retention.
878    ///
879    /// # Errors
880    /// Returns an error when the store cannot be enumerated or swept.
881    pub async fn gc_deleted_with_retention(&self, retention: Duration) -> Result<u64, NodeError> {
882        let _gc = self.gc_lock.lock().await;
883        let mut live = Vec::new();
884        for root_name in self.store.list_roots("db:").await? {
885            if let Some(root) = self
886                .store
887                .get_root(&root_name)
888                .await?
889                .as_deref()
890                .and_then(DbRoot::decode)
891                && let Some(roots) = root.roots
892            {
893                live.extend(roots);
894            }
895        }
896        let report = mark_and_sweep_retained(
897            self.store.as_ref(),
898            live,
899            |_, _| Ok(Vec::new()),
900            retention,
901            SystemTime::now(),
902        )
903        .await?;
904        self.metrics
905            .record_gc(report.swept as u64, report.retained as u64);
906        tracing::info!(
907            marked = report.marked,
908            swept = report.swept,
909            retained = report.retained,
910            "garbage collection completed"
911        );
912        Ok(report.swept as u64)
913    }
914
915    /// Validates, appends, applies, and reports one transaction supplied as
916    /// composite-encoded EDN transaction forms.
917    ///
918    /// # Errors
919    /// Returns [`NodeError`] for decode/validation failures, lease loss, or
920    /// storage failures.
921    pub async fn transact(
922        &self,
923        name: &str,
924        tx_data: &[u8],
925    ) -> Result<pb::TransactResponse, NodeError> {
926        let started = Instant::now();
927        let result = self
928            .transact_inner(name, tx_data)
929            .instrument(tracing::info_span!("transact", db = name))
930            .await;
931        self.metrics.record_tx(started.elapsed(), result.is_ok());
932        if let Err(error) = &result {
933            tracing::warn!(%error, "transaction failed");
934        }
935        result
936    }
937
938    async fn transact_inner(
939        &self,
940        name: &str,
941        tx_data: &[u8],
942    ) -> Result<pb::TransactResponse, NodeError> {
943        let state = self.db_state(name).await?;
944        let decoded = codec::decode_edn(tx_data)?;
945        let forms = decoded
946            .as_seq()
947            .ok_or_else(|| NodeError::BadRequest("tx-data must be a vector".into()))?
948            .to_vec();
949        let queued = self.metrics.queue_waiter();
950        let commit = state.commit.lock().await;
951        drop(queued);
952        let _commit = commit;
953        state.check_lease(self.store.as_ref()).await?;
954        // Expand user database-function invocations against the
955        // db-in-transaction (the value under the commit lock) before native
956        // conversion. The expander blocks up to its budget deadline, so it
957        // runs off the async workers.
958        let forms = if let Some(expander) = &self.config.tx_fn_expander {
959            let expander = Arc::clone(expander);
960            let db = state.transactor.db();
961            tokio::task::spawn_blocking(move || expander.expand(&db, forms))
962                .await
963                .map_err(|error| NodeError::BadRequest(format!("expander task failed: {error}")))?
964                .map_err(NodeError::BadRequest)?
965        } else {
966            forms
967        };
968        // Convert forms, interning any new keyword values.
969        let (items, naming_changed, idents, interner, schema) = {
970            let mut naming = state
971                .naming
972                .lock()
973                .unwrap_or_else(std::sync::PoisonError::into_inner);
974            let db = state.transactor.db();
975            let before = naming.interner.len();
976            let items = tx_items_from_edn(&db, &mut naming.interner, &forms)?;
977            (
978                items,
979                naming.interner.len() > before,
980                naming.idents.clone(),
981                naming.interner.clone(),
982                naming.schema.clone(),
983            )
984        };
985        if naming_changed {
986            // New keyword names must be durable before the datoms that
987            // reference them; recovery decodes the log against this meta.
988            let meta = encode_meta(&schema, &idents, &interner);
989            loop {
990                let current = self.store.get_root(&meta_root_name(name)).await?;
991                match self
992                    .store
993                    .cas_root(&meta_root_name(name), current.as_deref(), &meta)
994                    .await
995                {
996                    Ok(()) => break,
997                    Err(StoreError::CasFailed { .. }) => {}
998                    Err(error) => return Err(error.into()),
999                }
1000            }
1001            state.transactor.update_naming(idents, interner.clone());
1002        }
1003        let worker = Arc::clone(&state);
1004        let report = tokio::task::spawn_blocking(move || worker.transactor.transact(items))
1005            .await
1006            .map_err(|error| NodeError::BadRequest(format!("transact task failed: {error}")))??;
1007        // Post-append fence: acknowledge only if ownership was intact after
1008        // the record became durable. A takeover that raced the append will
1009        // have replayed the log *after* rewriting the root record, so a
1010        // record we acked here is provably in the successor's replay, and a
1011        // record we refuse here lands in our version's log file where the
1012        // successor's cutoff discards it (see log-and-transactor.md).
1013        if let Err(error) = state.check_lease(self.store.as_ref()).await {
1014            if matches!(error, NodeError::Deposed(_)) {
1015                self.depose(&state, "write lease lost after durable append");
1016            }
1017            // Either way the transaction is not acknowledged; a transient
1018            // store failure here is ambiguous to the caller, exactly like a
1019            // crash between append and reply.
1020            return Err(error);
1021        }
1022        let t = report.db_after.basis_t();
1023        let datoms = codec::encode_datoms(&report.tx.datoms, &interner)?;
1024        let tempids = codec::encode_edn(&Edn::Map(
1025            report
1026                .tx
1027                .tempids
1028                .iter()
1029                .map(|(tempid, eid)| {
1030                    (
1031                        Edn::Str(tempid.clone()),
1032                        Edn::Long(i64::try_from(eid.raw()).unwrap_or(i64::MAX)),
1033                    )
1034                })
1035                .collect(),
1036        ));
1037        let _ = state
1038            .broadcast
1039            .send(pb::subscribe_item::Item::Report(pb::TxReport {
1040                t,
1041                tx_instant: report.tx_instant,
1042                datoms: datoms.clone(),
1043            }));
1044        let _ = state.basis.send(t);
1045        Ok(pb::TransactResponse {
1046            basis_before: report.db_before.basis_t(),
1047            basis_t: t,
1048            tx_instant: report.tx_instant,
1049            tempids,
1050            tx_data: datoms,
1051        })
1052    }
1053
1054    /// Current status for a database.
1055    ///
1056    /// # Errors
1057    /// Returns [`NodeError::UnknownDb`] when absent.
1058    pub async fn status(&self, name: &str) -> Result<pb::StatusResponse, NodeError> {
1059        let state = self.db_state(name).await?;
1060        let db = state.db();
1061        let counts = db.stats();
1062        let held = state.lease();
1063        let metrics = self.metrics.snapshot();
1064        Ok(pb::StatusResponse {
1065            basis_t: db.basis_t(),
1066            index_basis_t: state.index_basis(),
1067            lease_owner: held.owner,
1068            lease_version: held.version,
1069            lease_expires_unix_ms: held.expires_unix_ms,
1070            datom_count: counts.datoms as u64,
1071            entity_count: counts.entities as u64,
1072            attribute_count: counts.attributes as u64,
1073            transaction_count: metrics.tx_total,
1074            transaction_failure_count: metrics.tx_failed,
1075            transaction_queue_depth: metrics.queue_depth,
1076            index_lag: db.basis_t().saturating_sub(state.index_basis()),
1077            indexing_runs: metrics.index_runs,
1078            gc_runs: metrics.gc_runs,
1079            gc_swept_blobs: metrics.gc_swept,
1080            lease_owner_endpoint: held.endpoint,
1081        })
1082    }
1083
1084    /// Releases every held write lease (graceful shutdown): the record is
1085    /// expired in place so a standby's next poll takes over immediately
1086    /// instead of waiting out the TTL. Hosted databases stop accepting
1087    /// work first, so nothing commits after its lease is gone.
1088    pub async fn release_leases(&self) {
1089        let states: Vec<Arc<DbState>> = self
1090            .dbs
1091            .write()
1092            .unwrap_or_else(std::sync::PoisonError::into_inner)
1093            .drain()
1094            .map(|(_, state)| state)
1095            .collect();
1096        for state in states {
1097            state.deposed.store(true, Ordering::Release);
1098            if let Err(error) =
1099                lease::release(self.store.as_ref(), &state.name, &state.lease()).await
1100            {
1101                tracing::warn!(db = %state.name, %error, "lease release failed at shutdown");
1102            }
1103        }
1104    }
1105
1106    /// Waits until the database basis reaches `t`, returning the basis seen.
1107    ///
1108    /// # Errors
1109    /// Returns [`NodeError::UnknownDb`] when absent.
1110    pub async fn sync(&self, name: &str, t: u64) -> Result<u64, NodeError> {
1111        let state = self.db_state(name).await?;
1112        let mut basis = state.basis_watch();
1113        let target = if t == 0 { *basis.borrow() } else { t };
1114        loop {
1115            let current = *basis.borrow();
1116            if current >= target {
1117                return Ok(current);
1118            }
1119            if basis.changed().await.is_err() {
1120                return Ok(*basis.borrow());
1121            }
1122        }
1123    }
1124}