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