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, 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
28pub trait TxFnExpander: Send + Sync {
33 fn expand(&self, db: &Db, forms: Vec<Edn>) -> Result<Vec<Edn>, String>;
40}
41
42#[derive(Clone)]
44pub struct NodeConfig {
45 pub data_dir: PathBuf,
47 pub owner: String,
49 pub lease_ttl_ms: i64,
51 pub lease_wait_ms: i64,
53 pub ha: bool,
57 pub advertise: Option<String>,
60 pub index_interval: Duration,
62 pub heartbeat_interval: Duration,
64 pub gc_interval: Option<Duration>,
66 pub gc_retention: Duration,
68 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 #[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#[derive(Debug, Error)]
115pub enum NodeError {
116 #[error("unknown database {0:?}")]
118 UnknownDb(String),
119 #[error("invalid database name {0:?}")]
121 InvalidName(String),
122 #[error("storage format {found} is newer than supported format {supported}")]
124 UnsupportedFormat {
125 found: u32,
127 supported: u32,
129 },
130 #[error("deposed: write lease for {0:?} is held elsewhere")]
132 Deposed(String),
133 #[error("standby for {db:?}: lease held by {owner} at {endpoint:?}")]
136 Standby {
137 db: String,
139 owner: String,
141 endpoint: String,
143 },
144 #[error(transparent)]
146 Codec(#[from] CodecError),
147 #[error(transparent)]
149 TxForm(#[from] TxFormError),
150 #[error(transparent)]
152 SchemaForm(#[from] SchemaFormError),
153 #[error(transparent)]
155 Transact(#[from] TransactError),
156 #[error(transparent)]
158 Store(#[from] StoreError),
159 #[error(transparent)]
161 Log(#[from] LogError),
162 #[error(transparent)]
164 Lease(#[from] LeaseError),
165 #[error("bad request: {0}")]
167 BadRequest(String),
168}
169
170struct Naming {
171 schema: Schema,
172 idents: Idents,
173 interner: KeywordInterner,
174}
175
176pub 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 #[must_use]
193 pub fn name(&self) -> &str {
194 &self.name
195 }
196
197 #[must_use]
199 pub fn db(&self) -> Db {
200 self.transactor.db()
201 }
202
203 #[must_use]
205 pub fn basis_watch(&self) -> watch::Receiver<u64> {
206 self.basis.subscribe()
207 }
208
209 #[must_use]
212 pub fn stream_items(&self) -> broadcast::Receiver<pb::subscribe_item::Item> {
213 self.broadcast.subscribe()
214 }
215
216 #[must_use]
218 pub fn index_basis(&self) -> u64 {
219 self.index_basis.load(Ordering::Acquire)
220 }
221
222 #[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 #[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 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 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
271pub struct TransactorNode {
273 config: NodeConfig,
274 store: Arc<FsStore>,
275 dbs: std::sync::RwLock<HashMap<String, Arc<DbState>>>,
276 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 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 #[must_use]
390 pub fn store(&self) -> &Arc<FsStore> {
391 &self.store
392 }
393
394 #[must_use]
396 pub fn config(&self) -> &NodeConfig {
397 &self.config
398 }
399
400 #[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 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 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 #[must_use]
432 pub fn shutdown_watch(&self) -> watch::Receiver<Option<String>> {
433 self.shutdown.subscribe()
434 }
435
436 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 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) = ¤t {
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 let held = self.acquire_lease(name).await?;
513 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 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 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 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 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 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 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 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 #[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 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 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 #[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 pub async fn gc_deleted(&self) -> Result<u64, NodeError> {
868 self.gc_deleted_with_retention(self.config.gc_retention)
869 .await
870 }
871
872 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 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 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 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 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 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 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 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 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 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}