1use std::collections::HashMap;
5use std::path::PathBuf;
6use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
7use std::sync::{Arc, Mutex};
8use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
9
10use corium_core::{KeywordInterner, Schema};
11use corium_db::{Db, Idents};
12use corium_log::{FileLog, LogError, TransactionLog, TxRecord};
13use corium_protocol::codec::{self, CodecError};
14use corium_protocol::pb;
15use corium_protocol::schemaform::{SchemaFormError, schema_from_edn};
16use corium_protocol::txforms::{TxFormError, tx_items_from_edn};
17use corium_query::edn::Edn;
18use corium_store::{FsStore, RootStore, StoreError, mark_and_sweep_retained};
19use thiserror::Error;
20use tokio::sync::{broadcast, watch};
21use tracing::Instrument;
22
23use crate::lease::{self, Lease, LeaseError};
24use crate::metrics::Metrics;
25use crate::{DbRoot, EmbeddedTransactor, TransactError, db_root_name, publish_root};
26
27pub trait TxFnExpander: Send + Sync {
32 fn expand(&self, db: &Db, forms: Vec<Edn>) -> Result<Vec<Edn>, String>;
39}
40
41#[derive(Clone)]
43pub struct NodeConfig {
44 pub data_dir: PathBuf,
46 pub owner: String,
48 pub lease_ttl_ms: i64,
50 pub lease_wait_ms: i64,
52 pub index_interval: Duration,
54 pub heartbeat_interval: Duration,
56 pub gc_interval: Option<Duration>,
58 pub gc_retention: Duration,
60 pub tx_fn_expander: Option<Arc<dyn TxFnExpander>>,
62}
63
64impl std::fmt::Debug for NodeConfig {
65 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
66 f.debug_struct("NodeConfig")
67 .field("data_dir", &self.data_dir)
68 .field("owner", &self.owner)
69 .field("lease_ttl_ms", &self.lease_ttl_ms)
70 .field("lease_wait_ms", &self.lease_wait_ms)
71 .field("index_interval", &self.index_interval)
72 .field("heartbeat_interval", &self.heartbeat_interval)
73 .field("gc_interval", &self.gc_interval)
74 .field("gc_retention", &self.gc_retention)
75 .field("tx_fn_expander", &self.tx_fn_expander.is_some())
76 .finish()
77 }
78}
79
80impl NodeConfig {
81 #[must_use]
83 pub fn new(data_dir: PathBuf) -> Self {
84 Self {
85 data_dir,
86 owner: format!(
87 "transactor-{}",
88 std::env::var("HOSTNAME").unwrap_or_else(|_| "local".into())
89 ),
90 lease_ttl_ms: 5_000,
91 lease_wait_ms: 15_000,
92 index_interval: Duration::from_secs(5),
93 heartbeat_interval: Duration::from_secs(10),
94 gc_interval: Some(Duration::from_secs(60 * 60)),
95 gc_retention: Duration::from_secs(72 * 60 * 60),
96 tx_fn_expander: None,
97 }
98 }
99}
100
101#[derive(Debug, Error)]
103pub enum NodeError {
104 #[error("unknown database {0:?}")]
106 UnknownDb(String),
107 #[error("invalid database name {0:?}")]
109 InvalidName(String),
110 #[error("storage format {found} is newer than supported format {supported}")]
112 UnsupportedFormat {
113 found: u32,
115 supported: u32,
117 },
118 #[error("deposed: write lease for {0:?} is held elsewhere")]
120 Deposed(String),
121 #[error(transparent)]
123 Codec(#[from] CodecError),
124 #[error(transparent)]
126 TxForm(#[from] TxFormError),
127 #[error(transparent)]
129 SchemaForm(#[from] SchemaFormError),
130 #[error(transparent)]
132 Transact(#[from] TransactError),
133 #[error(transparent)]
135 Store(#[from] StoreError),
136 #[error(transparent)]
138 Log(#[from] LogError),
139 #[error(transparent)]
141 Lease(#[from] LeaseError),
142 #[error("bad request: {0}")]
144 BadRequest(String),
145}
146
147struct Naming {
148 schema: Schema,
149 idents: Idents,
150 interner: KeywordInterner,
151}
152
153pub struct DbState {
155 name: String,
156 transactor: EmbeddedTransactor,
157 log: Arc<FileLog>,
158 naming: Mutex<Naming>,
159 commit: tokio::sync::Mutex<()>,
160 broadcast: broadcast::Sender<pb::subscribe_item::Item>,
161 basis: watch::Sender<u64>,
162 index_basis: AtomicU64,
163 held_lease: Mutex<Lease>,
164 deposed: AtomicBool,
165}
166
167impl DbState {
168 #[must_use]
170 pub fn name(&self) -> &str {
171 &self.name
172 }
173
174 #[must_use]
176 pub fn db(&self) -> Db {
177 self.transactor.db()
178 }
179
180 #[must_use]
182 pub fn basis_watch(&self) -> watch::Receiver<u64> {
183 self.basis.subscribe()
184 }
185
186 #[must_use]
189 pub fn stream_items(&self) -> broadcast::Receiver<pb::subscribe_item::Item> {
190 self.broadcast.subscribe()
191 }
192
193 #[must_use]
195 pub fn index_basis(&self) -> u64 {
196 self.index_basis.load(Ordering::Acquire)
197 }
198
199 #[must_use]
201 pub fn lease(&self) -> Lease {
202 self.held_lease
203 .lock()
204 .unwrap_or_else(std::sync::PoisonError::into_inner)
205 .clone()
206 }
207
208 #[must_use]
211 pub fn handshake_snapshot(&self) -> (Vec<u8>, KeywordInterner) {
212 let naming = self
213 .naming
214 .lock()
215 .unwrap_or_else(std::sync::PoisonError::into_inner);
216 (
217 codec::encode_schema(&naming.schema, &naming.idents),
218 naming.interner.clone(),
219 )
220 }
221
222 pub fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, NodeError> {
227 Ok(self.log.tx_range(start, end)?)
228 }
229
230 fn check_lease(&self, store: &dyn RootStore) -> Result<Lease, NodeError> {
231 if self.deposed.load(Ordering::Acquire) {
232 return Err(NodeError::Deposed(self.name.clone()));
233 }
234 let held = self.lease();
235 let stored = store.get_root(&lease::lease_root(&self.name))?;
236 if stored.as_deref() == Some(held.encode().as_slice()) {
237 return Ok(held);
238 }
239 if let Some(stored) = stored.as_deref().and_then(Lease::decode) {
243 if stored.owner == held.owner && stored.version == held.version {
244 *self
245 .held_lease
246 .lock()
247 .unwrap_or_else(std::sync::PoisonError::into_inner) = stored.clone();
248 return Ok(stored);
249 }
250 }
251 self.deposed.store(true, Ordering::Release);
252 Err(NodeError::Deposed(self.name.clone()))
253 }
254}
255
256pub struct TransactorNode {
258 config: NodeConfig,
259 store: Arc<FsStore>,
260 dbs: std::sync::RwLock<HashMap<String, Arc<DbState>>>,
261 gc_lock: tokio::sync::Mutex<()>,
262 metrics: Metrics,
263 shutdown: watch::Sender<Option<String>>,
264}
265
266fn now_unix_ms() -> i64 {
267 i64::try_from(
268 SystemTime::now()
269 .duration_since(UNIX_EPOCH)
270 .unwrap_or_default()
271 .as_millis(),
272 )
273 .unwrap_or(i64::MAX)
274}
275
276fn valid_db_name(name: &str) -> bool {
277 !name.is_empty()
278 && name.len() <= 128
279 && name
280 .bytes()
281 .all(|byte| byte.is_ascii_alphanumeric() || byte == b'-' || byte == b'_')
282}
283
284fn meta_root_name(db: &str) -> String {
285 format!("meta:{db}")
286}
287
288fn encode_meta(schema: &Schema, idents: &Idents, interner: &KeywordInterner) -> Vec<u8> {
289 let schema_bytes = codec::encode_schema(schema, idents);
290 let naming_bytes = codec::encode_naming(interner);
291 let mut out = Vec::with_capacity(8 + schema_bytes.len() + naming_bytes.len());
292 out.extend_from_slice(&u32::try_from(schema_bytes.len()).unwrap_or(0).to_be_bytes());
293 out.extend_from_slice(&schema_bytes);
294 out.extend_from_slice(&u32::try_from(naming_bytes.len()).unwrap_or(0).to_be_bytes());
295 out.extend_from_slice(&naming_bytes);
296 out
297}
298
299fn decode_meta(bytes: &[u8]) -> Result<(Schema, Idents, KeywordInterner), NodeError> {
300 let take = |input: &mut &[u8]| -> Result<Vec<u8>, NodeError> {
301 let len_bytes = input
302 .get(..4)
303 .ok_or(NodeError::Codec(CodecError::Truncated))?;
304 let len = usize::try_from(u32::from_be_bytes(len_bytes.try_into().unwrap_or_default()))
305 .map_err(|_| NodeError::Codec(CodecError::Length))?;
306 let payload = input
307 .get(4..4 + len)
308 .ok_or(NodeError::Codec(CodecError::Truncated))?
309 .to_vec();
310 *input = &input[4 + len..];
311 Ok(payload)
312 };
313 let mut input = bytes;
314 let schema_bytes = take(&mut input)?;
315 let naming_bytes = take(&mut input)?;
316 let (schema, idents) = codec::decode_schema(&schema_bytes)?;
317 let interner = codec::decode_naming(&naming_bytes)?;
318 Ok((schema, idents, interner))
319}
320
321impl TransactorNode {
322 pub fn open(config: NodeConfig) -> Result<Arc<Self>, NodeError> {
330 let store = Arc::new(FsStore::open(config.data_dir.join("store"))?);
331 let node = Arc::new(Self {
332 config,
333 store,
334 dbs: std::sync::RwLock::new(HashMap::new()),
335 gc_lock: tokio::sync::Mutex::new(()),
336 metrics: Metrics::default(),
337 shutdown: watch::channel(None).0,
338 });
339 let names: Vec<String> = node
340 .store
341 .list_roots("meta:")?
342 .into_iter()
343 .filter_map(|root| root.strip_prefix("meta:").map(str::to_owned))
344 .collect();
345 for name in names {
346 let state = node.open_db(&name)?;
347 node.dbs
348 .write()
349 .unwrap_or_else(std::sync::PoisonError::into_inner)
350 .insert(name, state);
351 }
352 node.spawn_scheduled_gc();
353 Ok(node)
354 }
355
356 #[must_use]
358 pub fn store(&self) -> &Arc<FsStore> {
359 &self.store
360 }
361
362 #[must_use]
364 pub fn config(&self) -> &NodeConfig {
365 &self.config
366 }
367
368 #[must_use]
370 pub const fn metrics(&self) -> &Metrics {
371 &self.metrics
372 }
373
374 fn spawn_scheduled_gc(self: &Arc<Self>) {
375 let Some(interval) = self.config.gc_interval else {
376 return;
377 };
378 let Ok(runtime) = tokio::runtime::Handle::try_current() else {
379 return;
382 };
383 let node = Arc::clone(self);
384 runtime.spawn(async move {
385 let mut ticker = tokio::time::interval(interval);
386 ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
387 ticker.tick().await;
389 loop {
390 ticker.tick().await;
391 if let Err(error) = node.gc_deleted().await {
392 tracing::warn!(%error, "scheduled garbage collection failed");
393 }
394 }
395 });
396 }
397
398 #[must_use]
400 pub fn shutdown_watch(&self) -> watch::Receiver<Option<String>> {
401 self.shutdown.subscribe()
402 }
403
404 fn depose(&self, state: &DbState, reason: &str) {
405 state.deposed.store(true, Ordering::Release);
406 let _ = self
407 .shutdown
408 .send(Some(format!("database {:?}: {reason}", state.name)));
409 }
410
411 fn log_path(&self, name: &str) -> PathBuf {
412 self.config
413 .data_dir
414 .join("logs")
415 .join(format!("{name}.log"))
416 }
417
418 fn acquire_with_wait(&self, name: &str) -> Result<Lease, NodeError> {
419 let deadline = now_unix_ms() + self.config.lease_wait_ms;
420 loop {
421 match lease::acquire(
422 self.store.as_ref(),
423 name,
424 &self.config.owner,
425 self.config.lease_ttl_ms,
426 now_unix_ms(),
427 ) {
428 Ok(held) => return Ok(held),
429 Err(LeaseError::Held { .. }) if now_unix_ms() < deadline => {
430 std::thread::sleep(Duration::from_millis(200));
431 }
432 Err(error) => return Err(error.into()),
433 }
434 }
435 }
436
437 fn open_db(self: &Arc<Self>, name: &str) -> Result<Arc<DbState>, NodeError> {
438 let meta = self
439 .store
440 .get_root(&meta_root_name(name))?
441 .ok_or_else(|| NodeError::UnknownDb(name.to_owned()))?;
442 let (schema, idents, interner) = decode_meta(&meta)?;
443 let root_name = db_root_name(name);
444 let current = self
445 .store
446 .get_root(&root_name)?
447 .as_deref()
448 .and_then(DbRoot::decode);
449 if let Some(root) = ¤t {
450 if root.format_version > corium_store::FORMAT_VERSION {
451 return Err(NodeError::UnsupportedFormat {
452 found: root.format_version,
453 supported: corium_store::FORMAT_VERSION,
454 });
455 }
456 }
457 let held = self.acquire_with_wait(name)?;
458 if current
461 .as_ref()
462 .is_none_or(|root| root.lease_version < held.version)
463 {
464 publish_root(
465 self.store.as_ref(),
466 &root_name,
467 &DbRoot {
468 format_version: corium_store::FORMAT_VERSION,
469 lease_version: held.version,
470 index_basis_t: current.as_ref().map_or(0, |root| root.index_basis_t),
471 roots: current.and_then(|root| root.roots),
472 },
473 )?;
474 }
475 let log = Arc::new(FileLog::open(self.log_path(name))?);
476 let base = Db::new(schema.clone()).with_naming(idents.clone(), interner.clone());
477 let transactor = EmbeddedTransactor::recover_from(base, Arc::clone(&log) as _)?;
478 let basis_t = transactor.db().basis_t();
479 let index_basis = self
480 .store
481 .get_root(&root_name)?
482 .as_deref()
483 .and_then(DbRoot::decode)
484 .map_or(0, |root| root.index_basis_t);
485 let state = Arc::new(DbState {
486 name: name.to_owned(),
487 transactor,
488 log,
489 naming: Mutex::new(Naming {
490 schema,
491 idents,
492 interner,
493 }),
494 commit: tokio::sync::Mutex::new(()),
495 broadcast: broadcast::channel(1024).0,
496 basis: watch::channel(basis_t).0,
497 index_basis: AtomicU64::new(index_basis),
498 held_lease: Mutex::new(held),
499 deposed: AtomicBool::new(false),
500 });
501 self.spawn_maintenance(&state);
502 Ok(state)
503 }
504
505 fn spawn_maintenance(self: &Arc<Self>, state: &Arc<DbState>) {
506 let ttl = self.config.lease_ttl_ms;
507 let renew_every = Duration::from_millis(u64::try_from(ttl / 3).unwrap_or(1).max(50));
508 let node = Arc::clone(self);
510 let db = Arc::clone(state);
511 tokio::spawn(async move {
512 let mut ticker = tokio::time::interval(renew_every);
513 ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
514 loop {
515 ticker.tick().await;
516 if db.deposed.load(Ordering::Acquire) {
517 return;
518 }
519 let held = db.lease();
520 let store = Arc::clone(&node.store);
521 let name = db.name.clone();
522 let renewed = tokio::task::spawn_blocking(move || {
523 lease::renew(store.as_ref(), &name, &held, ttl, now_unix_ms())
524 })
525 .await;
526 match renewed {
527 Ok(Ok(renewed)) => {
528 *db.held_lease
529 .lock()
530 .unwrap_or_else(std::sync::PoisonError::into_inner) = renewed;
531 }
532 Ok(Err(LeaseError::Lost)) => {
533 node.depose(&db, "write lease lost");
534 return;
535 }
536 Ok(Err(_)) | Err(_) => {}
537 }
538 }
539 });
540 let node = Arc::clone(self);
542 let db = Arc::clone(state);
543 tokio::spawn(async move {
544 let mut ticker = tokio::time::interval(node.config.index_interval);
545 ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
546 loop {
547 ticker.tick().await;
548 if db.deposed.load(Ordering::Acquire) {
549 return;
550 }
551 if db.db().basis_t() <= db.index_basis() {
552 continue;
553 }
554 let _gc = node.gc_lock.lock().await;
555 let store = Arc::clone(&node.store);
556 let version = db.lease().version;
557 let root_name = db_root_name(&db.name);
558 let worker = Arc::clone(&db);
559 let started = Instant::now();
560 let published = tokio::task::spawn_blocking(move || {
561 worker
562 .transactor
563 .publish_indexes(store.as_ref(), &root_name, version)
564 })
565 .await;
566 node.metrics.record_index(started.elapsed());
567 match published {
568 Ok(Ok(root)) => {
569 tracing::debug!(db = %db.name, index_basis_t = root.index_basis_t, "published indexes");
570 db.index_basis.store(root.index_basis_t, Ordering::Release);
571 let _ = db.broadcast.send(pb::subscribe_item::Item::IndexBasis(
572 pb::IndexBasis {
573 index_basis_t: root.index_basis_t,
574 },
575 ));
576 }
577 Ok(Err(TransactError::Deposed { .. })) => {
578 node.depose(&db, "database root fenced by a newer lease");
579 return;
580 }
581 Ok(Err(_)) | Err(_) => {}
582 }
583 }
584 });
585 let node = Arc::clone(self);
587 let db = Arc::clone(state);
588 tokio::spawn(async move {
589 let mut ticker = tokio::time::interval(node.config.heartbeat_interval);
590 ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
591 loop {
592 ticker.tick().await;
593 if db.deposed.load(Ordering::Acquire) {
594 return;
595 }
596 let _ = db
597 .broadcast
598 .send(pb::subscribe_item::Item::Heartbeat(pb::Heartbeat {
599 basis_t: db.db().basis_t(),
600 }));
601 }
602 });
603 }
604
605 pub fn db_state(&self, name: &str) -> Result<Arc<DbState>, NodeError> {
610 self.dbs
611 .read()
612 .unwrap_or_else(std::sync::PoisonError::into_inner)
613 .get(name)
614 .cloned()
615 .ok_or_else(|| NodeError::UnknownDb(name.to_owned()))
616 }
617
618 pub fn create_db(self: &Arc<Self>, name: &str, schema_edn: &[u8]) -> Result<bool, NodeError> {
624 if !valid_db_name(name) {
625 return Err(NodeError::InvalidName(name.to_owned()));
626 }
627 if self
628 .dbs
629 .read()
630 .unwrap_or_else(std::sync::PoisonError::into_inner)
631 .contains_key(name)
632 {
633 return Ok(false);
634 }
635 let forms = match codec::decode_edn(schema_edn)? {
636 Edn::Vector(items) | Edn::List(items) => items,
637 Edn::Nil => Vec::new(),
638 other => {
639 return Err(NodeError::BadRequest(format!(
640 "schema must be a vector of attribute maps, got {other}"
641 )));
642 }
643 };
644 let (schema, idents) = schema_from_edn(&forms)?;
645 let meta = encode_meta(&schema, &idents, &KeywordInterner::default());
646 match self.store.cas_root(&meta_root_name(name), None, &meta) {
647 Ok(()) => {}
648 Err(StoreError::CasFailed { .. }) => return Ok(false),
649 Err(error) => return Err(error.into()),
650 }
651 let state = self.open_db(name)?;
652 self.dbs
653 .write()
654 .unwrap_or_else(std::sync::PoisonError::into_inner)
655 .insert(name.to_owned(), state);
656 Ok(true)
657 }
658
659 pub fn delete_db(&self, name: &str) -> Result<bool, NodeError> {
665 let Some(state) = self
666 .dbs
667 .write()
668 .unwrap_or_else(std::sync::PoisonError::into_inner)
669 .remove(name)
670 else {
671 return Ok(false);
672 };
673 state.deposed.store(true, Ordering::Release);
674 let _ = lease::release(self.store.as_ref(), name, &state.lease());
675 self.store.delete_root(&db_root_name(name))?;
676 self.store.delete_root(&meta_root_name(name))?;
677 self.store.delete_root(&lease::lease_root(name))?;
678 match std::fs::remove_file(self.log_path(name)) {
679 Ok(()) => {}
680 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
681 Err(error) => return Err(NodeError::Store(StoreError::Io(error))),
682 }
683 Ok(true)
684 }
685
686 #[must_use]
688 pub fn list_dbs(&self) -> Vec<String> {
689 let mut names: Vec<String> = self
690 .dbs
691 .read()
692 .unwrap_or_else(std::sync::PoisonError::into_inner)
693 .keys()
694 .cloned()
695 .collect();
696 names.sort();
697 names
698 }
699
700 pub async fn gc_deleted(&self) -> Result<u64, NodeError> {
706 self.gc_deleted_with_retention(self.config.gc_retention)
707 .await
708 }
709
710 pub async fn gc_deleted_with_retention(&self, retention: Duration) -> Result<u64, NodeError> {
715 let _gc = self.gc_lock.lock().await;
716 let store = Arc::clone(&self.store);
717 let report =
718 tokio::task::spawn_blocking(move || -> Result<corium_store::GcReport, NodeError> {
719 let mut live = Vec::new();
720 for root_name in store.list_roots("db:")? {
721 if let Some(root) = store
722 .get_root(&root_name)?
723 .as_deref()
724 .and_then(DbRoot::decode)
725 {
726 if let Some(roots) = root.roots {
727 live.extend(roots);
728 }
729 }
730 }
731 Ok(mark_and_sweep_retained(
732 store.as_ref(),
733 live,
734 |_, _| Ok(Vec::new()),
735 retention,
736 SystemTime::now(),
737 )?)
738 })
739 .await
740 .map_err(|error| NodeError::BadRequest(format!("gc task failed: {error}")))??;
741 self.metrics
742 .record_gc(report.swept as u64, report.retained as u64);
743 tracing::info!(
744 marked = report.marked,
745 swept = report.swept,
746 retained = report.retained,
747 "garbage collection completed"
748 );
749 Ok(report.swept as u64)
750 }
751
752 pub async fn transact(
759 &self,
760 name: &str,
761 tx_data: &[u8],
762 ) -> Result<pb::TransactResponse, NodeError> {
763 let started = Instant::now();
764 let result = self
765 .transact_inner(name, tx_data)
766 .instrument(tracing::info_span!("transact", db = name))
767 .await;
768 self.metrics.record_tx(started.elapsed(), result.is_ok());
769 if let Err(error) = &result {
770 tracing::warn!(%error, "transaction failed");
771 }
772 result
773 }
774
775 async fn transact_inner(
776 &self,
777 name: &str,
778 tx_data: &[u8],
779 ) -> Result<pb::TransactResponse, NodeError> {
780 let state = self.db_state(name)?;
781 let decoded = codec::decode_edn(tx_data)?;
782 let forms = decoded
783 .as_seq()
784 .ok_or_else(|| NodeError::BadRequest("tx-data must be a vector".into()))?
785 .to_vec();
786 let queued = self.metrics.queue_waiter();
787 let commit = state.commit.lock().await;
788 drop(queued);
789 let _commit = commit;
790 state.check_lease(self.store.as_ref())?;
791 let forms = if let Some(expander) = &self.config.tx_fn_expander {
796 let expander = Arc::clone(expander);
797 let db = state.transactor.db();
798 tokio::task::spawn_blocking(move || expander.expand(&db, forms))
799 .await
800 .map_err(|error| NodeError::BadRequest(format!("expander task failed: {error}")))?
801 .map_err(NodeError::BadRequest)?
802 } else {
803 forms
804 };
805 let (items, naming_changed, idents, interner, schema) = {
807 let mut naming = state
808 .naming
809 .lock()
810 .unwrap_or_else(std::sync::PoisonError::into_inner);
811 let db = state.transactor.db();
812 let before = naming.interner.len();
813 let items = tx_items_from_edn(&db, &mut naming.interner, &forms)?;
814 (
815 items,
816 naming.interner.len() > before,
817 naming.idents.clone(),
818 naming.interner.clone(),
819 naming.schema.clone(),
820 )
821 };
822 if naming_changed {
823 let meta = encode_meta(&schema, &idents, &interner);
826 loop {
827 let current = self.store.get_root(&meta_root_name(name))?;
828 match self
829 .store
830 .cas_root(&meta_root_name(name), current.as_deref(), &meta)
831 {
832 Ok(()) => break,
833 Err(StoreError::CasFailed { .. }) => {}
834 Err(error) => return Err(error.into()),
835 }
836 }
837 state.transactor.update_naming(idents, interner.clone());
838 }
839 let worker = Arc::clone(&state);
840 let report = tokio::task::spawn_blocking(move || worker.transactor.transact(items))
841 .await
842 .map_err(|error| NodeError::BadRequest(format!("transact task failed: {error}")))??;
843 let t = report.db_after.basis_t();
844 let datoms = codec::encode_datoms(&report.tx.datoms, &interner)?;
845 let tempids = codec::encode_edn(&Edn::Map(
846 report
847 .tx
848 .tempids
849 .iter()
850 .map(|(tempid, eid)| {
851 (
852 Edn::Str(tempid.clone()),
853 Edn::Long(i64::try_from(eid.raw()).unwrap_or(i64::MAX)),
854 )
855 })
856 .collect(),
857 ));
858 let _ = state
859 .broadcast
860 .send(pb::subscribe_item::Item::Report(pb::TxReport {
861 t,
862 tx_instant: report.tx_instant,
863 datoms: datoms.clone(),
864 }));
865 let _ = state.basis.send(t);
866 Ok(pb::TransactResponse {
867 basis_before: report.db_before.basis_t(),
868 basis_t: t,
869 tx_instant: report.tx_instant,
870 tempids,
871 tx_data: datoms,
872 })
873 }
874
875 pub fn status(&self, name: &str) -> Result<pb::StatusResponse, NodeError> {
880 let state = self.db_state(name)?;
881 let db = state.db();
882 let counts = db.stats();
883 let held = state.lease();
884 let metrics = self.metrics.snapshot();
885 Ok(pb::StatusResponse {
886 basis_t: db.basis_t(),
887 index_basis_t: state.index_basis(),
888 lease_owner: held.owner,
889 lease_version: held.version,
890 lease_expires_unix_ms: held.expires_unix_ms,
891 datom_count: counts.datoms as u64,
892 entity_count: counts.entities as u64,
893 attribute_count: counts.attributes as u64,
894 transaction_count: metrics.tx_total,
895 transaction_failure_count: metrics.tx_failed,
896 transaction_queue_depth: metrics.queue_depth,
897 index_lag: db.basis_t().saturating_sub(state.index_basis()),
898 indexing_runs: metrics.index_runs,
899 gc_runs: metrics.gc_runs,
900 gc_swept_blobs: metrics.gc_swept,
901 })
902 }
903
904 pub async fn sync(&self, name: &str, t: u64) -> Result<u64, NodeError> {
909 let state = self.db_state(name)?;
910 let mut basis = state.basis_watch();
911 let target = if t == 0 { *basis.borrow() } else { t };
912 loop {
913 let current = *basis.borrow();
914 if current >= target {
915 return Ok(current);
916 }
917 if basis.changed().await.is_err() {
918 return Ok(*basis.borrow());
919 }
920 }
921 }
922}