1use 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
29pub trait TxFnExpander: Send + Sync {
34 fn expand(&self, db: &Db, forms: Vec<Edn>) -> Result<Vec<Edn>, String>;
41}
42
43#[derive(Clone)]
45pub struct NodeConfig {
46 pub store: StoreSpec,
48 pub data_dir: PathBuf,
51 pub owner: String,
53 pub lease_ttl_ms: i64,
55 pub lease_wait_ms: i64,
57 pub ha: bool,
61 pub advertise: Option<String>,
64 pub index_interval: Duration,
66 pub heartbeat_interval: Duration,
68 pub gc_interval: Option<Duration>,
70 pub gc_retention: Duration,
72 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 #[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#[derive(Debug, Error)]
121pub enum NodeError {
122 #[error("unknown database {0:?}")]
124 UnknownDb(String),
125 #[error("invalid database name {0:?}")]
127 InvalidName(String),
128 #[error("storage format {found} is newer than supported format {supported}")]
130 UnsupportedFormat {
131 found: u32,
133 supported: u32,
135 },
136 #[error("deposed: write lease for {0:?} is held elsewhere")]
138 Deposed(String),
139 #[error("standby for {db:?}: lease held by {owner} at {endpoint:?}")]
142 Standby {
143 db: String,
145 owner: String,
147 endpoint: String,
149 },
150 #[error(transparent)]
152 Codec(#[from] CodecError),
153 #[error(transparent)]
155 TxForm(#[from] TxFormError),
156 #[error(transparent)]
158 SchemaForm(#[from] SchemaFormError),
159 #[error(transparent)]
161 Transact(#[from] TransactError),
162 #[error(transparent)]
164 Store(#[from] StoreError),
165 #[error(transparent)]
167 Log(#[from] LogError),
168 #[error(transparent)]
170 Lease(#[from] LeaseError),
171 #[error("bad request: {0}")]
173 BadRequest(String),
174}
175
176struct Naming {
177 schema: Schema,
178 idents: Idents,
179 interner: KeywordInterner,
180}
181
182pub 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 #[must_use]
199 pub fn name(&self) -> &str {
200 &self.name
201 }
202
203 #[must_use]
205 pub fn db(&self) -> Db {
206 self.transactor.db()
207 }
208
209 #[must_use]
211 pub fn basis_watch(&self) -> watch::Receiver<u64> {
212 self.basis.subscribe()
213 }
214
215 #[must_use]
218 pub fn stream_items(&self) -> broadcast::Receiver<pb::subscribe_item::Item> {
219 self.broadcast.subscribe()
220 }
221
222 #[must_use]
224 pub fn index_basis(&self) -> u64 {
225 self.index_basis.load(Ordering::Acquire)
226 }
227
228 #[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 #[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 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 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
277pub struct TransactorNode {
279 config: NodeConfig,
280 store: Arc<NodeStore>,
281 log_backend: LogBackend,
282 dbs: std::sync::RwLock<HashMap<String, Arc<DbState>>>,
283 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 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 #[must_use]
399 pub fn store(&self) -> &Arc<NodeStore> {
400 &self.store
401 }
402
403 #[must_use]
405 pub fn config(&self) -> &NodeConfig {
406 &self.config
407 }
408
409 #[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 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 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 #[must_use]
441 pub fn shutdown_watch(&self) -> watch::Receiver<Option<String>> {
442 self.shutdown.subscribe()
443 }
444
445 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 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) = ¤t
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 let held = self.acquire_lease(name).await?;
522 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 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 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 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 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 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 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 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 #[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 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 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 #[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 pub async fn gc_deleted(&self) -> Result<u64, NodeError> {
873 self.gc_deleted_with_retention(self.config.gc_retention)
874 .await
875 }
876
877 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 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 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 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 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 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 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 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 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 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}