1use std::collections::{BTreeSet, HashMap, VecDeque};
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::{Datom, IndexOrder, 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::{
20 BlobId, BlobStore, RootStore, StoreError, decode_index_manifest, decode_segment_keys,
21 is_index_manifest, mark_and_sweep_retained, meta_root_name,
22};
23use thiserror::Error;
24use tokio::sync::{broadcast, oneshot, watch};
25use tracing::Instrument;
26
27use crate::backend::{LogBackend, NodeStore, StoreSpec};
28use crate::lease::{self, Lease, LeaseError};
29use crate::metrics::Metrics;
30use crate::{DbRoot, EmbeddedTransactor, Prepared, TransactError, db_root_name};
31
32pub trait TxFnExpander: Send + Sync {
37 fn expand(&self, db: &Db, forms: Vec<Edn>) -> Result<Vec<Edn>, String>;
44}
45
46#[derive(Clone)]
48pub struct NodeConfig {
49 pub store: StoreSpec,
51 pub data_dir: PathBuf,
54 pub owner: String,
56 pub lease_ttl_ms: i64,
58 pub lease_wait_ms: i64,
60 pub ha: bool,
64 pub advertise: Option<String>,
67 pub index_interval: Duration,
69 pub index_backoff: u32,
75 pub index_tail_threshold: u64,
79 pub index_tail_deadline: Duration,
81 pub heartbeat_interval: Duration,
83 pub gc_interval: Option<Duration>,
85 pub gc_retention: Duration,
87 pub max_commit_batch: usize,
93 pub max_commit_batch_bytes: usize,
98 pub tx_fn_expander: Option<Arc<dyn TxFnExpander>>,
100}
101
102impl std::fmt::Debug for NodeConfig {
103 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
104 f.debug_struct("NodeConfig")
105 .field("store", &self.store)
106 .field("data_dir", &self.data_dir)
107 .field("owner", &self.owner)
108 .field("lease_ttl_ms", &self.lease_ttl_ms)
109 .field("lease_wait_ms", &self.lease_wait_ms)
110 .field("ha", &self.ha)
111 .field("advertise", &self.advertise)
112 .field("index_interval", &self.index_interval)
113 .field("index_backoff", &self.index_backoff)
114 .field("index_tail_threshold", &self.index_tail_threshold)
115 .field("index_tail_deadline", &self.index_tail_deadline)
116 .field("heartbeat_interval", &self.heartbeat_interval)
117 .field("gc_interval", &self.gc_interval)
118 .field("gc_retention", &self.gc_retention)
119 .field("max_commit_batch", &self.max_commit_batch)
120 .field("max_commit_batch_bytes", &self.max_commit_batch_bytes)
121 .field("tx_fn_expander", &self.tx_fn_expander.is_some())
122 .finish()
123 }
124}
125
126impl NodeConfig {
127 #[must_use]
129 pub fn new(data_dir: PathBuf) -> Self {
130 Self {
131 store: StoreSpec::Fs,
132 data_dir,
133 owner: format!(
134 "transactor-{}",
135 std::env::var("HOSTNAME").unwrap_or_else(|_| "local".into())
136 ),
137 lease_ttl_ms: 5_000,
138 lease_wait_ms: 15_000,
139 ha: false,
140 advertise: None,
141 index_interval: Duration::from_secs(5),
142 index_backoff: 4,
143 index_tail_threshold: 0,
144 index_tail_deadline: Duration::from_secs(60),
145 heartbeat_interval: Duration::from_secs(10),
146 gc_interval: Some(Duration::from_secs(60 * 60)),
147 gc_retention: Duration::from_secs(72 * 60 * 60),
148 max_commit_batch: 256,
149 max_commit_batch_bytes: 4 * 1024 * 1024,
150 #[cfg(feature = "cljrs")]
151 tx_fn_expander: Some(Arc::new(crate::txfn::DbFnExpander::default())),
152 #[cfg(not(feature = "cljrs"))]
153 tx_fn_expander: None,
154 }
155 }
156}
157
158#[derive(Clone, Copy, Debug, Eq, PartialEq)]
171pub struct IndexPolicy {
172 pub interval: Duration,
174 pub backoff: u32,
177 pub tail_threshold: u64,
180 pub tail_deadline: Duration,
183}
184
185#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
187pub struct IndexPolicyUpdate {
188 pub interval: Option<Duration>,
190 pub backoff: Option<u32>,
192 pub tail_threshold: Option<u64>,
194 pub tail_deadline: Option<Duration>,
196}
197
198impl IndexPolicy {
199 fn from_config(config: &NodeConfig) -> Self {
200 Self {
201 interval: config.index_interval,
202 backoff: config.index_backoff,
203 tail_threshold: config.index_tail_threshold,
204 tail_deadline: config.index_tail_deadline,
205 }
206 }
207
208 fn apply(&mut self, update: IndexPolicyUpdate) {
209 if let Some(interval) = update.interval {
210 self.interval = interval;
211 }
212 if let Some(backoff) = update.backoff {
213 self.backoff = backoff;
214 }
215 if let Some(tail_threshold) = update.tail_threshold {
216 self.tail_threshold = tail_threshold;
217 }
218 if let Some(tail_deadline) = update.tail_deadline {
219 self.tail_deadline = tail_deadline;
220 }
221 }
222
223 fn due(&self, since_publish: Duration, last_duration: Duration, pending: Option<u64>) -> bool {
230 let floor = self
231 .interval
232 .max(last_duration.saturating_mul(self.backoff));
233 if since_publish < floor {
234 return false;
235 }
236 match pending {
237 Some(pending) if pending < self.tail_threshold => since_publish >= self.tail_deadline,
238 _ => true,
239 }
240 }
241}
242
243#[derive(Debug, Error)]
245pub enum NodeError {
246 #[error("unknown database {0:?}")]
248 UnknownDb(String),
249 #[error("invalid database name {0:?}")]
251 InvalidName(String),
252 #[error("storage format {found} is newer than supported format {supported}")]
254 UnsupportedFormat {
255 found: u32,
257 supported: u32,
259 },
260 #[error("deposed: write lease for {0:?} is held elsewhere")]
262 Deposed(String),
263 #[error("standby for {db:?}: lease held by {owner} at {endpoint:?}")]
266 Standby {
267 db: String,
269 owner: String,
271 endpoint: String,
273 },
274 #[error(transparent)]
276 Codec(#[from] CodecError),
277 #[error(transparent)]
279 TxForm(#[from] TxFormError),
280 #[error(transparent)]
282 SchemaForm(#[from] SchemaFormError),
283 #[error(transparent)]
285 Transact(#[from] TransactError),
286 #[error(transparent)]
288 Store(#[from] StoreError),
289 #[error(transparent)]
291 Log(#[from] LogError),
292 #[error(transparent)]
294 Lease(#[from] LeaseError),
295 #[error("bad request: {0}")]
297 BadRequest(String),
298 #[error("group commit aborted: {0}")]
303 GroupCommit(String),
304}
305
306struct Naming {
307 schema: Schema,
308 idents: Idents,
309 interner: KeywordInterner,
310}
311
312struct CommitRequest {
315 forms: Vec<Edn>,
316 resp: oneshot::Sender<Result<pb::TransactResponse, NodeError>>,
317}
318
319pub struct DbState {
321 name: String,
322 transactor: EmbeddedTransactor,
323 log: Arc<dyn TransactionLog>,
324 naming: Mutex<Naming>,
325 commit: tokio::sync::Mutex<()>,
329 pending: Mutex<VecDeque<CommitRequest>>,
331 broadcast: broadcast::Sender<pb::subscribe_item::Item>,
332 basis: watch::Sender<u64>,
333 index_basis: AtomicU64,
334 index_policy: Mutex<IndexPolicy>,
335 held_lease: Mutex<Lease>,
336 deposed: AtomicBool,
337}
338
339impl DbState {
340 #[must_use]
342 pub fn name(&self) -> &str {
343 &self.name
344 }
345
346 #[must_use]
348 pub fn db(&self) -> Db {
349 self.transactor.db()
350 }
351
352 #[must_use]
354 pub fn basis_watch(&self) -> watch::Receiver<u64> {
355 self.basis.subscribe()
356 }
357
358 #[must_use]
361 pub fn stream_items(&self) -> broadcast::Receiver<pb::subscribe_item::Item> {
362 self.broadcast.subscribe()
363 }
364
365 #[must_use]
367 pub fn index_basis(&self) -> u64 {
368 self.index_basis.load(Ordering::Acquire)
369 }
370
371 #[must_use]
373 pub fn index_policy(&self) -> IndexPolicy {
374 *self
375 .index_policy
376 .lock()
377 .unwrap_or_else(std::sync::PoisonError::into_inner)
378 }
379
380 #[must_use]
382 pub fn lease(&self) -> Lease {
383 self.held_lease
384 .lock()
385 .unwrap_or_else(std::sync::PoisonError::into_inner)
386 .clone()
387 }
388
389 #[must_use]
392 pub fn handshake_snapshot(&self) -> (Vec<u8>, KeywordInterner) {
393 let naming = self
394 .naming
395 .lock()
396 .unwrap_or_else(std::sync::PoisonError::into_inner);
397 (
398 codec::encode_schema(&naming.schema, &naming.idents),
399 naming.interner.clone(),
400 )
401 }
402
403 pub async fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, NodeError> {
408 Ok(self.log.tx_range_async(start, end).await?)
409 }
410
411 async fn check_lease(&self, store: &dyn RootStore) -> Result<Lease, NodeError> {
414 if self.deposed.load(Ordering::Acquire) {
415 return Err(NodeError::Deposed(self.name.clone()));
416 }
417 let held = self.lease();
418 match lease::verify(store, &self.name, &held).await {
419 Ok(()) => Ok(held),
420 Err(LeaseError::Lost) => {
421 self.deposed.store(true, Ordering::Release);
422 Err(NodeError::Deposed(self.name.clone()))
423 }
424 Err(error) => Err(error.into()),
425 }
426 }
427}
428
429pub struct TransactorNode {
431 config: NodeConfig,
432 store: Arc<NodeStore>,
433 log_backend: LogBackend,
434 dbs: std::sync::RwLock<HashMap<String, Arc<DbState>>>,
435 standby: std::sync::RwLock<BTreeSet<String>>,
438 gc_lock: tokio::sync::Mutex<()>,
439 fork_lock: tokio::sync::Mutex<()>,
442 metrics: Metrics,
443 shutdown: watch::Sender<Option<String>>,
444}
445
446fn now_unix_ms() -> i64 {
447 i64::try_from(
448 SystemTime::now()
449 .duration_since(UNIX_EPOCH)
450 .unwrap_or_default()
451 .as_millis(),
452 )
453 .unwrap_or(i64::MAX)
454}
455
456fn batch_abort_error(name: &str, error: &NodeError) -> NodeError {
462 match error {
463 NodeError::Deposed(_) => NodeError::Deposed(name.to_owned()),
464 other => NodeError::GroupCommit(other.to_string()),
465 }
466}
467
468fn valid_db_name(name: &str) -> bool {
469 !name.is_empty()
470 && name.len() <= 128
471 && name
472 .bytes()
473 .all(|byte| byte.is_ascii_alphanumeric() || byte == b'-' || byte == b'_')
474}
475
476impl TransactorNode {
477 pub async fn open(config: NodeConfig) -> Result<Arc<Self>, NodeError> {
485 let store = Arc::new(NodeStore::open(&config.store, &config.data_dir).await?);
486 let log_backend = LogBackend::for_spec(&config.store, &config.data_dir, Arc::clone(&store));
487 let node = Arc::new(Self {
488 config,
489 store,
490 log_backend,
491 dbs: std::sync::RwLock::new(HashMap::new()),
492 standby: std::sync::RwLock::new(BTreeSet::new()),
493 gc_lock: tokio::sync::Mutex::new(()),
494 fork_lock: tokio::sync::Mutex::new(()),
495 metrics: Metrics::default(),
496 shutdown: watch::channel(None).0,
497 });
498 let names: Vec<String> = node
499 .store
500 .list_roots("meta:")
501 .await?
502 .into_iter()
503 .filter_map(|root| root.strip_prefix("meta:").map(str::to_owned))
504 .collect();
505 for name in names {
506 match node.open_db(&name).await {
507 Ok(state) => {
508 node.dbs
509 .write()
510 .unwrap_or_else(std::sync::PoisonError::into_inner)
511 .insert(name, state);
512 }
513 Err(NodeError::Lease(LeaseError::Held { owner, .. })) if node.config.ha => {
514 tracing::info!(db = %name, %owner, "standing by; lease held elsewhere");
515 node.standby
516 .write()
517 .unwrap_or_else(std::sync::PoisonError::into_inner)
518 .insert(name);
519 }
520 Err(error) => return Err(error),
521 }
522 }
523 node.spawn_standby_poller();
524 node.spawn_scheduled_gc();
525 Ok(node)
526 }
527
528 #[must_use]
530 pub fn store(&self) -> &Arc<NodeStore> {
531 &self.store
532 }
533
534 #[must_use]
536 pub fn config(&self) -> &NodeConfig {
537 &self.config
538 }
539
540 #[must_use]
542 pub const fn metrics(&self) -> &Metrics {
543 &self.metrics
544 }
545
546 fn spawn_scheduled_gc(self: &Arc<Self>) {
547 let Some(interval) = self.config.gc_interval else {
548 return;
549 };
550 let Ok(runtime) = tokio::runtime::Handle::try_current() else {
551 return;
554 };
555 let node = Arc::clone(self);
556 runtime.spawn(async move {
557 let mut ticker = tokio::time::interval(interval);
558 ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
559 ticker.tick().await;
561 loop {
562 ticker.tick().await;
563 if let Err(error) = node.gc_deleted().await {
564 tracing::warn!(%error, "scheduled garbage collection failed");
565 }
566 }
567 });
568 }
569
570 #[must_use]
572 pub fn shutdown_watch(&self) -> watch::Receiver<Option<String>> {
573 self.shutdown.subscribe()
574 }
575
576 fn depose(&self, state: &DbState, reason: &str) {
580 state.deposed.store(true, Ordering::Release);
581 if self.config.ha {
582 tracing::warn!(db = %state.name, reason, "deposed; returning to standby");
583 self.dbs
584 .write()
585 .unwrap_or_else(std::sync::PoisonError::into_inner)
586 .remove(&state.name);
587 self.standby
588 .write()
589 .unwrap_or_else(std::sync::PoisonError::into_inner)
590 .insert(state.name.clone());
591 } else {
592 let _ = self
593 .shutdown
594 .send(Some(format!("database {:?}: {reason}", state.name)));
595 }
596 }
597
598 fn advertised(&self) -> &str {
599 self.config.advertise.as_deref().unwrap_or("")
600 }
601
602 async fn acquire_lease(&self, name: &str) -> Result<Lease, NodeError> {
606 let deadline = now_unix_ms() + self.config.lease_wait_ms;
607 loop {
608 match lease::acquire(
609 self.store.as_ref(),
610 name,
611 &self.config.owner,
612 self.advertised(),
613 self.config.lease_ttl_ms,
614 now_unix_ms(),
615 )
616 .await
617 {
618 Ok(held) => return Ok(held),
619 Err(LeaseError::Held { .. }) if !self.config.ha && now_unix_ms() < deadline => {
620 tokio::time::sleep(Duration::from_millis(200)).await;
621 }
622 Err(error) => return Err(error.into()),
623 }
624 }
625 }
626
627 async fn open_db(self: &Arc<Self>, name: &str) -> Result<Arc<DbState>, NodeError> {
628 let meta = self
629 .store
630 .get_root(&meta_root_name(name))
631 .await?
632 .ok_or_else(|| NodeError::UnknownDb(name.to_owned()))?;
633 let (schema, idents, interner) = codec::decode_metadata(&meta)?;
634 let root_name = db_root_name(name);
635 let current = self
636 .store
637 .get_root(&root_name)
638 .await?
639 .as_deref()
640 .and_then(DbRoot::decode);
641 if let Some(root) = ¤t
642 && root.format_version > corium_store::FORMAT_VERSION
643 {
644 return Err(NodeError::UnsupportedFormat {
645 found: root.format_version,
646 supported: corium_store::FORMAT_VERSION,
647 });
648 }
649 let held = self.acquire_lease(name).await?;
655 let log = self.log_backend.open(name, held.version).await?;
658 let post_fence = self
659 .store
660 .get_root(&root_name)
661 .await?
662 .as_deref()
663 .and_then(DbRoot::decode);
664 let transactor = self
665 .recover_transactor(name, &schema, &idents, &interner, post_fence.as_ref(), &log)
666 .await?;
667 let basis_t = transactor.db().basis_t();
668 let index_basis = post_fence.map_or(0, |root| root.index_basis_t);
669 let state = Arc::new(DbState {
670 name: name.to_owned(),
671 transactor,
672 log,
673 naming: Mutex::new(Naming {
674 schema,
675 idents,
676 interner,
677 }),
678 commit: tokio::sync::Mutex::new(()),
679 pending: Mutex::new(VecDeque::new()),
680 broadcast: broadcast::channel(1024).0,
681 basis: watch::channel(basis_t).0,
682 index_basis: AtomicU64::new(index_basis),
683 index_policy: Mutex::new(IndexPolicy::from_config(&self.config)),
684 held_lease: Mutex::new(held),
685 deposed: AtomicBool::new(false),
686 });
687 self.spawn_maintenance(&state);
688 Ok(state)
689 }
690
691 async fn recover_transactor(
700 &self,
701 name: &str,
702 schema: &Schema,
703 idents: &Idents,
704 interner: &KeywordInterner,
705 root: Option<&DbRoot>,
706 log: &Arc<dyn TransactionLog>,
707 ) -> Result<EmbeddedTransactor, NodeError> {
708 if let Some(root) = root
711 && let Some(roots) = &root.roots
712 && root.next_entity_id != 0
713 {
714 match self
715 .load_current_snapshot(
716 root,
717 &roots[IndexOrder::Eavt as usize],
718 schema,
719 idents,
720 interner,
721 )
722 .await
723 {
724 Ok(snapshot) => {
725 return Ok(EmbeddedTransactor::recover_from_snapshot_async(
726 snapshot,
727 root.next_entity_id,
728 root.last_tx_instant,
729 Arc::clone(log),
730 )
731 .await?);
732 }
733 Err(error) => {
734 tracing::warn!(
735 db = %name,
736 %error,
737 "index-root recovery failed; falling back to full-log replay"
738 );
739 }
740 }
741 }
742 let base = Db::new(schema.clone()).with_naming(idents.clone(), interner.clone());
743 Ok(EmbeddedTransactor::recover_from_async(base, Arc::clone(log)).await?)
744 }
745
746 async fn load_current_snapshot(
751 &self,
752 root: &DbRoot,
753 eavt: &BlobId,
754 schema: &Schema,
755 idents: &Idents,
756 interner: &KeywordInterner,
757 ) -> Result<Db, StoreError> {
758 let datoms = self
759 .load_index_keys(eavt)
760 .await?
761 .into_iter()
762 .map(|key| Datom::from_key(IndexOrder::Eavt, &key))
763 .collect::<Result<Vec<_>, _>>()
764 .map_err(|error| StoreError::Io(std::io::Error::other(error.to_string())))?;
765 Ok(Db::from_current_snapshot(
766 root.index_basis_t,
767 schema.clone(),
768 idents.clone(),
769 interner.clone(),
770 datoms,
771 ))
772 }
773
774 async fn load_index_keys(&self, id: &BlobId) -> Result<Vec<Vec<u8>>, StoreError> {
777 let blob = self
778 .store
779 .get(id)
780 .await?
781 .ok_or_else(|| StoreError::MissingBlob(id.clone()))?;
782 if !is_index_manifest(&blob) {
783 return decode_segment_keys(&blob);
784 }
785 let mut keys = Vec::new();
786 for child in decode_index_manifest(&blob)? {
787 let chunk = self
788 .store
789 .get(&child)
790 .await?
791 .ok_or_else(|| StoreError::MissingBlob(child.clone()))?;
792 keys.extend(decode_segment_keys(&chunk)?);
793 }
794 Ok(keys)
795 }
796
797 fn spawn_standby_poller(self: &Arc<Self>) {
803 if !self.config.ha {
804 return;
805 }
806 let Ok(runtime) = tokio::runtime::Handle::try_current() else {
807 return;
808 };
809 let ttl = self.config.lease_ttl_ms;
810 let poll_every = Duration::from_millis(u64::try_from(ttl / 3).unwrap_or(1).max(50));
811 let node = Arc::clone(self);
812 runtime.spawn(async move {
813 let mut ticker = tokio::time::interval(poll_every);
814 ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
815 loop {
816 ticker.tick().await;
817 if let Err(error) = node.standby_scan().await {
818 tracing::warn!(%error, "standby scan failed");
819 }
820 }
821 });
822 }
823
824 async fn standby_scan(self: &Arc<Self>) -> Result<(), NodeError> {
827 let names: Vec<String> = self
828 .store
829 .list_roots("meta:")
830 .await?
831 .into_iter()
832 .filter_map(|root| root.strip_prefix("meta:").map(str::to_owned))
833 .collect();
834 {
835 let mut standby = self
836 .standby
837 .write()
838 .unwrap_or_else(std::sync::PoisonError::into_inner);
839 standby.retain(|name| names.contains(name));
840 }
841 for name in names {
842 if self
843 .dbs
844 .read()
845 .unwrap_or_else(std::sync::PoisonError::into_inner)
846 .contains_key(&name)
847 {
848 continue;
849 }
850 match self.open_db(&name).await {
851 Ok(state) => {
852 tracing::info!(db = %name, owner = %self.config.owner, "standby took over write lease");
853 self.standby
854 .write()
855 .unwrap_or_else(std::sync::PoisonError::into_inner)
856 .remove(&name);
857 self.dbs
858 .write()
859 .unwrap_or_else(std::sync::PoisonError::into_inner)
860 .insert(name, state);
861 }
862 Err(NodeError::Lease(LeaseError::Held { .. })) => {
863 self.standby
864 .write()
865 .unwrap_or_else(std::sync::PoisonError::into_inner)
866 .insert(name);
867 }
868 Err(error) => {
869 tracing::warn!(db = %name, %error, "standby takeover attempt failed");
870 }
871 }
872 }
873 Ok(())
874 }
875
876 fn spawn_maintenance(self: &Arc<Self>, state: &Arc<DbState>) {
877 let ttl = self.config.lease_ttl_ms;
878 let renew_every = Duration::from_millis(u64::try_from(ttl / 3).unwrap_or(1).max(50));
879 let node = Arc::clone(self);
881 let db = Arc::clone(state);
882 tokio::spawn(async move {
883 let mut ticker = tokio::time::interval(renew_every);
884 ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
885 loop {
886 ticker.tick().await;
887 if db.deposed.load(Ordering::Acquire) {
888 return;
889 }
890 let _commit = db.commit.lock().await;
894 let held = db.lease();
895 let name = db.name.clone();
896 let renewed =
897 lease::renew(node.store.as_ref(), &name, &held, ttl, now_unix_ms()).await;
898 match renewed {
899 Ok(renewed) => {
900 *db.held_lease
901 .lock()
902 .unwrap_or_else(std::sync::PoisonError::into_inner) = renewed;
903 }
904 Err(LeaseError::Lost) => {
905 node.depose(&db, "write lease lost");
906 return;
907 }
908 Err(_) => {}
909 }
910 }
911 });
912 self.spawn_indexing(state);
913 let node = Arc::clone(self);
915 let db = Arc::clone(state);
916 tokio::spawn(async move {
917 let mut ticker = tokio::time::interval(node.config.heartbeat_interval);
918 ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
919 loop {
920 ticker.tick().await;
921 if db.deposed.load(Ordering::Acquire) {
922 return;
923 }
924 let _ = db
925 .broadcast
926 .send(pb::subscribe_item::Item::Heartbeat(pb::Heartbeat {
927 basis_t: db.db().basis_t(),
928 }));
929 }
930 });
931 }
932
933 fn spawn_indexing(self: &Arc<Self>, state: &Arc<DbState>) {
937 const POLICY_POLL: Duration = Duration::from_secs(1);
941 let node = Arc::clone(self);
942 let db = Arc::clone(state);
943 tokio::spawn(async move {
944 let mut published_at = Instant::now();
945 let mut last_duration = Duration::ZERO;
946 let mut published_len: Option<u64> = None;
947 loop {
948 let policy = db.index_policy();
949 tokio::time::sleep(policy.interval.min(POLICY_POLL)).await;
950 if db.deposed.load(Ordering::Acquire) {
951 return;
952 }
953 let snapshot = db.db();
954 if snapshot.basis_t() <= db.index_basis() {
955 continue;
956 }
957 let recorded_len = u64::try_from(snapshot.recorded_len()).unwrap_or(u64::MAX);
958 let pending = published_len.map(|len| recorded_len.saturating_sub(len));
959 if !policy.due(published_at.elapsed(), last_duration, pending) {
960 continue;
961 }
962 match node.publish_db_indexes(&db).await {
963 Ok((_, duration)) => {
964 last_duration = duration;
965 published_len = Some(recorded_len);
969 }
970 Err(NodeError::Deposed(_)) => return,
971 Err(_) => {}
972 }
973 published_at = Instant::now();
974 }
975 });
976 }
977
978 async fn publish_db_indexes(&self, db: &Arc<DbState>) -> Result<(u64, Duration), NodeError> {
983 let _gc = self.gc_lock.lock().await;
984 let version = db.lease().version;
985 let root_name = db_root_name(&db.name);
986 let started = Instant::now();
987 let published = db
988 .transactor
989 .publish_indexes(self.store.as_ref(), &root_name, version)
990 .await;
991 let duration = started.elapsed();
992 self.metrics.record_index(duration);
993 match published {
994 Ok(root) => {
995 tracing::debug!(db = %db.name, index_basis_t = root.index_basis_t, "published indexes");
996 db.index_basis.store(root.index_basis_t, Ordering::Release);
997 let _ = db
998 .broadcast
999 .send(pb::subscribe_item::Item::IndexBasis(pb::IndexBasis {
1000 index_basis_t: root.index_basis_t,
1001 }));
1002 Ok((root.index_basis_t, duration))
1003 }
1004 Err(TransactError::Deposed { .. }) => {
1005 self.depose(db, "database root fenced by a newer lease");
1006 Err(NodeError::Deposed(db.name.clone()))
1007 }
1008 Err(error) => Err(error.into()),
1009 }
1010 }
1011
1012 pub async fn request_index(&self, name: &str) -> Result<u64, NodeError> {
1021 let state = self.db_state(name).await?;
1022 if state.db().basis_t() <= state.index_basis() {
1023 return Ok(state.index_basis());
1024 }
1025 self.publish_db_indexes(&state)
1026 .await
1027 .map(|(index_basis_t, _)| index_basis_t)
1028 }
1029
1030 pub async fn set_index_policy(
1037 &self,
1038 name: &str,
1039 update: IndexPolicyUpdate,
1040 ) -> Result<IndexPolicy, NodeError> {
1041 let state = self.db_state(name).await?;
1042 let mut policy = state
1043 .index_policy
1044 .lock()
1045 .unwrap_or_else(std::sync::PoisonError::into_inner);
1046 policy.apply(update);
1047 Ok(*policy)
1048 }
1049
1050 pub async fn db_state(&self, name: &str) -> Result<Arc<DbState>, NodeError> {
1056 if let Some(state) = self
1057 .dbs
1058 .read()
1059 .unwrap_or_else(std::sync::PoisonError::into_inner)
1060 .get(name)
1061 .cloned()
1062 {
1063 return Ok(state);
1064 }
1065 if self.config.ha
1066 && self
1067 .standby
1068 .read()
1069 .unwrap_or_else(std::sync::PoisonError::into_inner)
1070 .contains(name)
1071 {
1072 let root = self
1073 .store
1074 .get_root(&db_root_name(name))
1075 .await?
1076 .as_deref()
1077 .and_then(DbRoot::decode);
1078 return Err(NodeError::Standby {
1079 db: name.to_owned(),
1080 owner: root.as_ref().map(|r| r.owner.clone()).unwrap_or_default(),
1081 endpoint: root.map(|r| r.owner_endpoint).unwrap_or_default(),
1082 });
1083 }
1084 Err(NodeError::UnknownDb(name.to_owned()))
1085 }
1086
1087 #[must_use]
1089 pub fn standby_dbs(&self) -> Vec<String> {
1090 self.standby
1091 .read()
1092 .unwrap_or_else(std::sync::PoisonError::into_inner)
1093 .iter()
1094 .cloned()
1095 .collect()
1096 }
1097
1098 pub async fn create_db(
1104 self: &Arc<Self>,
1105 name: &str,
1106 schema_edn: &[u8],
1107 ) -> Result<bool, NodeError> {
1108 if !valid_db_name(name) {
1109 return Err(NodeError::InvalidName(name.to_owned()));
1110 }
1111 if self
1112 .dbs
1113 .read()
1114 .unwrap_or_else(std::sync::PoisonError::into_inner)
1115 .contains_key(name)
1116 {
1117 return Ok(false);
1118 }
1119 let forms = match codec::decode_edn(schema_edn)? {
1120 Edn::Vector(items) | Edn::List(items) => items,
1121 Edn::Nil => Vec::new(),
1122 other => {
1123 return Err(NodeError::BadRequest(format!(
1124 "schema must be a vector of attribute maps, got {other}"
1125 )));
1126 }
1127 };
1128 let (schema, idents) = schema_from_edn(&forms)?;
1129 let meta = codec::encode_metadata(&schema, &idents, &KeywordInterner::default());
1130 match self
1131 .store
1132 .cas_root(&meta_root_name(name), None, &meta)
1133 .await
1134 {
1135 Ok(()) => {}
1136 Err(StoreError::CasFailed { .. }) => return Ok(false),
1137 Err(error) => return Err(error.into()),
1138 }
1139 let state = self.open_db(name).await?;
1140 self.dbs
1141 .write()
1142 .unwrap_or_else(std::sync::PoisonError::into_inner)
1143 .insert(name.to_owned(), state);
1144 Ok(true)
1145 }
1146
1147 pub async fn fork_db(
1158 self: &Arc<Self>,
1159 source: &str,
1160 target: &str,
1161 as_of_t: u64,
1162 ) -> Result<Option<u64>, NodeError> {
1163 if !valid_db_name(target) {
1164 return Err(NodeError::InvalidName(target.to_owned()));
1165 }
1166 if source == target {
1167 return Err(NodeError::BadRequest(
1168 "fork target must differ from the source".into(),
1169 ));
1170 }
1171 let state = self.db_state(source).await?;
1172 let basis = state.db().basis_t();
1173 let t = if as_of_t == 0 { basis } else { as_of_t };
1174 if t > basis {
1175 return Err(NodeError::BadRequest(format!(
1176 "as-of t {t} is ahead of {source:?} basis {basis}"
1177 )));
1178 }
1179 let _guard = self.fork_lock.lock().await;
1180 if self
1181 .dbs
1182 .read()
1183 .unwrap_or_else(std::sync::PoisonError::into_inner)
1184 .contains_key(target)
1185 || self
1186 .store
1187 .get_root(&meta_root_name(target))
1188 .await?
1189 .is_some()
1190 || self.log_backend.exists(target).await
1191 {
1192 return Ok(None);
1193 }
1194 let records = state.log.tx_range_async(0, Some(t + 1)).await?;
1200 let meta = self
1201 .store
1202 .get_root(&meta_root_name(source))
1203 .await?
1204 .ok_or_else(|| NodeError::UnknownDb(source.to_owned()))?;
1205 let log = self.log_backend.open(target, 0).await?;
1210 for record in &records {
1211 log.append_async(record).await?;
1212 }
1213 drop(log);
1214 match self
1215 .store
1216 .cas_root(&meta_root_name(target), None, &meta)
1217 .await
1218 {
1219 Ok(()) => {}
1220 Err(StoreError::CasFailed { .. }) => {
1221 self.log_backend.delete_all(target).await?;
1223 return Ok(None);
1224 }
1225 Err(error) => return Err(error.into()),
1226 }
1227 let state = self.open_db(target).await?;
1228 self.dbs
1229 .write()
1230 .unwrap_or_else(std::sync::PoisonError::into_inner)
1231 .insert(target.to_owned(), state);
1232 Ok(Some(t))
1233 }
1234
1235 pub async fn delete_db(&self, name: &str) -> Result<bool, NodeError> {
1241 let Some(state) = self
1242 .dbs
1243 .write()
1244 .unwrap_or_else(std::sync::PoisonError::into_inner)
1245 .remove(name)
1246 else {
1247 return Ok(false);
1248 };
1249 state.deposed.store(true, Ordering::Release);
1250 self.standby
1251 .write()
1252 .unwrap_or_else(std::sync::PoisonError::into_inner)
1253 .remove(name);
1254 self.store.delete_root(&db_root_name(name)).await?;
1255 self.store.delete_root(&meta_root_name(name)).await?;
1256 self.log_backend.delete_all(name).await?;
1257 Ok(true)
1258 }
1259
1260 #[must_use]
1262 pub fn list_dbs(&self) -> Vec<String> {
1263 let mut names: Vec<String> = self
1264 .dbs
1265 .read()
1266 .unwrap_or_else(std::sync::PoisonError::into_inner)
1267 .keys()
1268 .cloned()
1269 .collect();
1270 names.sort();
1271 names
1272 }
1273
1274 pub async fn gc_deleted(&self) -> Result<u64, NodeError> {
1280 self.gc_deleted_with_retention(self.config.gc_retention)
1281 .await
1282 }
1283
1284 pub async fn gc_deleted_with_retention(&self, retention: Duration) -> Result<u64, NodeError> {
1289 let _gc = self.gc_lock.lock().await;
1290 let mut live = Vec::new();
1291 for root_name in self.store.list_roots("db:").await? {
1292 if let Some(root) = self
1293 .store
1294 .get_root(&root_name)
1295 .await?
1296 .as_deref()
1297 .and_then(DbRoot::decode)
1298 && let Some(roots) = root.roots
1299 {
1300 live.extend(roots);
1301 }
1302 }
1303 let report = mark_and_sweep_retained(
1304 self.store.as_ref(),
1305 live,
1306 |_, bytes| corium_store::index_blob_children(bytes),
1307 retention,
1308 SystemTime::now(),
1309 )
1310 .await?;
1311 self.metrics
1312 .record_gc(report.swept as u64, report.retained as u64);
1313 tracing::info!(
1314 marked = report.marked,
1315 swept = report.swept,
1316 retained = report.retained,
1317 "garbage collection completed"
1318 );
1319 Ok(report.swept as u64)
1320 }
1321
1322 pub async fn transact(
1329 &self,
1330 name: &str,
1331 tx_data: &[u8],
1332 ) -> Result<pb::TransactResponse, NodeError> {
1333 let started = Instant::now();
1334 let result = self
1335 .transact_inner(name, tx_data)
1336 .instrument(tracing::info_span!("transact", db = name))
1337 .await;
1338 self.metrics.record_tx(started.elapsed(), result.is_ok());
1339 if let Err(error) = &result {
1340 tracing::warn!(%error, "transaction failed");
1341 }
1342 result
1343 }
1344
1345 async fn transact_inner(
1346 &self,
1347 name: &str,
1348 tx_data: &[u8],
1349 ) -> Result<pb::TransactResponse, NodeError> {
1350 let state = self.db_state(name).await?;
1351 let decoded = codec::decode_edn(tx_data)?;
1352 let forms = decoded
1353 .as_seq()
1354 .ok_or_else(|| NodeError::BadRequest("tx-data must be a vector".into()))?
1355 .to_vec();
1356 let (resp_tx, mut resp_rx) = oneshot::channel();
1362 state
1363 .pending
1364 .lock()
1365 .unwrap_or_else(std::sync::PoisonError::into_inner)
1366 .push_back(CommitRequest {
1367 forms,
1368 resp: resp_tx,
1369 });
1370 let queued = self.metrics.queue_waiter();
1371 loop {
1372 let commit = state.commit.lock().await;
1373 match resp_rx.try_recv() {
1375 Ok(result) => {
1376 drop(commit);
1377 drop(queued);
1378 return result;
1379 }
1380 Err(oneshot::error::TryRecvError::Empty) => {}
1381 Err(oneshot::error::TryRecvError::Closed) => {
1382 drop(commit);
1383 drop(queued);
1384 return Err(NodeError::GroupCommit("commit response dropped".into()));
1385 }
1386 }
1387 self.flush_commit_batch(&state).await;
1389 drop(commit);
1390 match resp_rx.try_recv() {
1391 Ok(result) => {
1392 drop(queued);
1393 return result;
1394 }
1395 Err(oneshot::error::TryRecvError::Empty) => {}
1398 Err(oneshot::error::TryRecvError::Closed) => {
1399 drop(queued);
1400 return Err(NodeError::GroupCommit("commit response dropped".into()));
1401 }
1402 }
1403 }
1404 }
1405
1406 #[allow(clippy::too_many_lines)]
1415 async fn flush_commit_batch(&self, state: &Arc<DbState>) {
1416 let max_count = self.config.max_commit_batch.max(1);
1417 let max_bytes = self.config.max_commit_batch_bytes;
1418
1419 let mut batch: VecDeque<CommitRequest> = {
1420 let mut pending = state
1421 .pending
1422 .lock()
1423 .unwrap_or_else(std::sync::PoisonError::into_inner);
1424 std::mem::take(&mut *pending)
1425 };
1426 if batch.is_empty() {
1427 return;
1428 }
1429 let now_ms = now_unix_ms();
1438 let mut cursor = state.transactor.batch_cursor();
1439 let mut resps: Vec<oneshot::Sender<Result<pb::TransactResponse, NodeError>>> = Vec::new();
1440 let mut prepared: Vec<Prepared> = Vec::new();
1441 let mut batch_bytes: usize = 0;
1442 let mut measure = Vec::new();
1443 let mut naming_changed = false;
1444 while let Some(request) = batch.pop_front() {
1445 let forms = if let Some(expander) = &self.config.tx_fn_expander {
1449 let expander = Arc::clone(expander);
1450 let db = cursor.db().clone();
1451 let forms = request.forms;
1452 match tokio::task::spawn_blocking(move || expander.expand(&db, forms)).await {
1453 Ok(Ok(forms)) => forms,
1454 Ok(Err(message)) => {
1455 let _ = request.resp.send(Err(NodeError::BadRequest(message)));
1456 continue;
1457 }
1458 Err(error) => {
1459 let _ = request.resp.send(Err(NodeError::BadRequest(format!(
1460 "expander task failed: {error}"
1461 ))));
1462 continue;
1463 }
1464 }
1465 } else {
1466 request.forms
1467 };
1468 let (items, this_changed) = {
1471 let mut naming = state
1472 .naming
1473 .lock()
1474 .unwrap_or_else(std::sync::PoisonError::into_inner);
1475 let before = naming.interner.len();
1476 match tx_items_from_edn(cursor.db(), &mut naming.interner, &forms) {
1477 Ok(items) => (items, naming.interner.len() > before),
1478 Err(error) => {
1479 drop(naming);
1480 let _ = request.resp.send(Err(error.into()));
1481 continue;
1482 }
1483 }
1484 };
1485 match cursor.prepare(items, now_ms) {
1486 Ok(prep) => {
1487 measure.clear();
1488 let _ = corium_log::append_framed_record(&mut measure, &prep.record);
1489 batch_bytes += measure.len();
1490 resps.push(request.resp);
1491 prepared.push(prep);
1492 }
1493 Err(error) => {
1494 let _ = request.resp.send(Err(NodeError::Transact(error.into())));
1495 continue;
1496 }
1497 }
1498 if this_changed {
1499 naming_changed = true;
1500 break;
1501 }
1502 if prepared.len() >= max_count || batch_bytes >= max_bytes {
1506 break;
1507 }
1508 }
1509 if !batch.is_empty() {
1511 let mut pending = state
1512 .pending
1513 .lock()
1514 .unwrap_or_else(std::sync::PoisonError::into_inner);
1515 while let Some(request) = batch.pop_back() {
1516 pending.push_front(request);
1517 }
1518 }
1519 if prepared.is_empty() {
1520 return;
1521 }
1522 let interner = {
1525 let naming = state
1526 .naming
1527 .lock()
1528 .unwrap_or_else(std::sync::PoisonError::into_inner);
1529 naming.interner.clone()
1530 };
1531 let changed_idents = if naming_changed {
1532 if let Err(error) = state.check_lease(self.store.as_ref()).await {
1536 if matches!(error, NodeError::Deposed(_)) {
1537 self.depose(state, "write lease lost before metadata publish");
1538 }
1539 for resp in resps {
1540 let _ = resp.send(Err(batch_abort_error(&state.name, &error)));
1541 }
1542 return;
1543 }
1544 let (idents, schema) = {
1545 let naming = state
1546 .naming
1547 .lock()
1548 .unwrap_or_else(std::sync::PoisonError::into_inner);
1549 (naming.idents.clone(), naming.schema.clone())
1550 };
1551 let meta = codec::encode_metadata(&schema, &idents, &interner);
1554 loop {
1555 let cas = match self.store.get_root(&meta_root_name(&state.name)).await {
1556 Ok(current) => {
1557 self.store
1558 .cas_root(&meta_root_name(&state.name), current.as_deref(), &meta)
1559 .await
1560 }
1561 Err(error) => Err(error),
1562 };
1563 match cas {
1564 Ok(()) => break,
1565 Err(StoreError::CasFailed { .. }) => {}
1566 Err(error) => {
1567 let error = NodeError::Store(error);
1568 for resp in resps {
1569 let _ = resp.send(Err(batch_abort_error(&state.name, &error)));
1570 }
1571 return;
1572 }
1573 }
1574 }
1575 Some(idents)
1576 } else {
1577 None
1578 };
1579 let records: Vec<TxRecord> = prepared.iter().map(|prep| prep.record.clone()).collect();
1581 if let Err(error) = state.log.append_batch_async(&records).await {
1582 let error = NodeError::Log(error);
1583 for resp in resps {
1584 let _ = resp.send(Err(batch_abort_error(&state.name, &error)));
1585 }
1586 return;
1587 }
1588 let reports = state.transactor.install_batch(cursor, prepared);
1594 if let Some(idents) = changed_idents {
1595 state.transactor.update_naming(idents, interner.clone());
1596 }
1597 if let Err(error) = state.check_lease(self.store.as_ref()).await {
1606 if matches!(error, NodeError::Deposed(_)) {
1607 self.depose(state, "write lease lost after durable append");
1608 }
1609 for resp in resps {
1610 let _ = resp.send(Err(batch_abort_error(&state.name, &error)));
1611 }
1612 return;
1613 }
1614 let mut last_t = 0;
1615 for (resp, report) in resps.into_iter().zip(reports) {
1616 let t = report.db_after.basis_t();
1617 last_t = last_t.max(t);
1618 let datoms = match codec::encode_datoms(&report.tx.datoms, &interner) {
1619 Ok(datoms) => datoms,
1620 Err(error) => {
1621 let _ = resp.send(Err(NodeError::Codec(error)));
1622 continue;
1623 }
1624 };
1625 let tempids = codec::encode_edn(&Edn::Map(
1626 report
1627 .tx
1628 .tempids
1629 .iter()
1630 .map(|(tempid, eid)| {
1631 (
1632 Edn::Str(tempid.clone()),
1633 Edn::Long(i64::try_from(eid.raw()).unwrap_or(i64::MAX)),
1634 )
1635 })
1636 .collect(),
1637 ));
1638 let _ = state
1639 .broadcast
1640 .send(pb::subscribe_item::Item::Report(pb::TxReport {
1641 t,
1642 tx_instant: report.tx_instant,
1643 datoms: datoms.clone(),
1644 }));
1645 let _ = resp.send(Ok(pb::TransactResponse {
1646 basis_before: report.db_before.basis_t(),
1647 basis_t: t,
1648 tx_instant: report.tx_instant,
1649 tempids,
1650 tx_data: datoms,
1651 }));
1652 }
1653 if last_t > 0 {
1654 let _ = state.basis.send(last_t);
1655 }
1656 }
1657
1658 pub async fn status(&self, name: &str) -> Result<pb::StatusResponse, NodeError> {
1663 let state = self.db_state(name).await?;
1664 let db = state.db();
1665 let counts = db.stats();
1666 let held = state.lease();
1667 let metrics = self.metrics.snapshot();
1668 Ok(pb::StatusResponse {
1669 basis_t: db.basis_t(),
1670 index_basis_t: state.index_basis(),
1671 lease_owner: held.owner,
1672 lease_version: held.version,
1673 lease_expires_unix_ms: held.expires_unix_ms,
1674 datom_count: counts.datoms as u64,
1675 entity_count: counts.entities as u64,
1676 attribute_count: counts.attributes as u64,
1677 transaction_count: metrics.tx_total,
1678 transaction_failure_count: metrics.tx_failed,
1679 transaction_queue_depth: metrics.queue_depth,
1680 index_lag: db.basis_t().saturating_sub(state.index_basis()),
1681 indexing_runs: metrics.index_runs,
1682 gc_runs: metrics.gc_runs,
1683 gc_swept_blobs: metrics.gc_swept,
1684 lease_owner_endpoint: held.endpoint,
1685 })
1686 }
1687
1688 pub async fn backup_info(&self, name: &str) -> Result<pb::GetStorageInfoResponse, NodeError> {
1696 let state = self.db_state(name).await?;
1697 let _commit = state.commit.lock().await;
1701 state.check_lease(self.store.as_ref()).await?;
1702 let basis_t = state.db().basis_t();
1703 let storage = self
1704 .config
1705 .store
1706 .connection_info(&self.config.data_dir)
1707 .map_err(NodeError::BadRequest)?;
1708 Ok(pb::GetStorageInfoResponse {
1709 basis_t,
1710 storage: Some(storage),
1711 })
1712 }
1713
1714 pub async fn release_leases(&self) {
1719 let states: Vec<Arc<DbState>> = self
1720 .dbs
1721 .write()
1722 .unwrap_or_else(std::sync::PoisonError::into_inner)
1723 .drain()
1724 .map(|(_, state)| state)
1725 .collect();
1726 for state in states {
1727 state.deposed.store(true, Ordering::Release);
1728 if let Err(error) =
1729 lease::release(self.store.as_ref(), &state.name, &state.lease()).await
1730 {
1731 tracing::warn!(db = %state.name, %error, "lease release failed at shutdown");
1732 }
1733 }
1734 }
1735
1736 pub async fn sync(&self, name: &str, t: u64) -> Result<u64, NodeError> {
1741 let state = self.db_state(name).await?;
1742 let mut basis = state.basis_watch();
1743 let target = if t == 0 { *basis.borrow() } else { t };
1744 loop {
1745 let current = *basis.borrow();
1746 if current >= target {
1747 return Ok(current);
1748 }
1749 if basis.changed().await.is_err() {
1750 return Ok(*basis.borrow());
1751 }
1752 }
1753 }
1754}
1755
1756#[cfg(test)]
1757mod tests {
1758 use super::IndexPolicy;
1759 use std::time::Duration;
1760
1761 fn pacing(interval_ms: u64, backoff: u32, threshold: u64, deadline_ms: u64) -> IndexPolicy {
1762 IndexPolicy {
1763 interval: Duration::from_millis(interval_ms),
1764 backoff,
1765 tail_threshold: threshold,
1766 tail_deadline: Duration::from_millis(deadline_ms),
1767 }
1768 }
1769
1770 #[test]
1771 fn base_interval_gates_publication() {
1772 let pacing = pacing(100, 4, 0, 60_000);
1773 assert!(!pacing.due(Duration::from_millis(99), Duration::ZERO, None));
1774 assert!(pacing.due(Duration::from_millis(100), Duration::ZERO, None));
1775 }
1776
1777 #[test]
1778 fn backoff_stretches_the_floor_past_the_interval() {
1779 let pacing = pacing(100, 4, 0, 60_000);
1780 let last = Duration::from_millis(300);
1781 assert!(!pacing.due(Duration::from_millis(1_199), last, Some(10)));
1782 assert!(pacing.due(Duration::from_millis(1_200), last, Some(10)));
1783 }
1784
1785 #[test]
1786 fn zero_backoff_keeps_the_base_interval() {
1787 let pacing = pacing(100, 0, 0, 60_000);
1788 assert!(pacing.due(
1789 Duration::from_millis(100),
1790 Duration::from_secs(30),
1791 Some(10)
1792 ));
1793 }
1794
1795 #[test]
1796 fn fast_publications_leave_the_interval_untouched() {
1797 let pacing = pacing(5_000, 4, 0, 60_000);
1798 assert!(pacing.due(Duration::from_secs(5), Duration::from_millis(3), Some(1)));
1799 }
1800
1801 #[test]
1802 fn small_tail_defers_until_the_deadline() {
1803 let pacing = pacing(100, 4, 1_000, 60_000);
1804 assert!(!pacing.due(Duration::from_secs(30), Duration::ZERO, Some(999)));
1805 assert!(pacing.due(Duration::from_secs(60), Duration::ZERO, Some(999)));
1806 }
1807
1808 #[test]
1809 fn tail_at_threshold_publishes_at_base_pacing() {
1810 let pacing = pacing(100, 4, 1_000, 60_000);
1811 assert!(pacing.due(Duration::from_millis(100), Duration::ZERO, Some(1_000)));
1812 }
1813
1814 #[test]
1815 fn unknown_tail_publishes_at_base_pacing() {
1816 let pacing = pacing(100, 4, 1_000, 60_000);
1817 assert!(pacing.due(Duration::from_millis(100), Duration::ZERO, None));
1818 }
1819
1820 #[test]
1821 fn deadline_never_overrides_the_backoff_floor() {
1822 let pacing = pacing(100, 4, 1_000, 200);
1823 let last = Duration::from_millis(300);
1824 assert!(!pacing.due(Duration::from_millis(400), last, Some(1)));
1825 assert!(pacing.due(Duration::from_millis(1_200), last, Some(1)));
1826 }
1827}