use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::{Duration, Instant};
use tokio::sync::{Mutex, MutexGuard};
const DEFAULT_EMBEDDER_INIT_SECS: u64 = 180;
const DEFAULT_EMBED_BATCH_SECS: u64 = 30;
const DEFAULT_DREAM_PERMIT_WAIT_SECS: u64 = 30;
const DEFAULT_WRITE_LOCK_SECS: u64 = 60;
const DEFAULT_OPEN_QUEUE_SECS: u64 = 60;
const DEFAULT_WRITE_OP_BUDGET_SECS: u64 = 60;
const _: () = assert!(
DEFAULT_WRITE_OP_BUDGET_SECS < DEFAULT_WRITE_LOCK_SECS + DEFAULT_OPEN_QUEUE_SECS,
"the joint write budget must be smaller than the additive worst case it replaces (#4002)"
);
const DEFAULT_WRITE_PIPELINE_SECS: u64 = 240;
const _: () = assert!(
DEFAULT_WRITE_PIPELINE_SECS > DEFAULT_EMBEDDER_INIT_SECS + DEFAULT_EMBED_BATCH_SECS,
"the write-pipeline ceiling must exceed the embedder legs it contains (#6366)"
);
const _: () = assert!(
DEFAULT_WRITE_PIPELINE_SECS > DEFAULT_WRITE_OP_BUDGET_SECS,
"the write-pipeline ceiling must exceed the acquisition budget it follows (#6366)"
);
const DEFAULT_WRITE_TXN_SECS: u64 = 30;
const _: () = assert!(
DEFAULT_WRITE_TXN_SECS < WRITE_PIPELINE_FLOOR_SECS,
"the write-transaction deadline must sit inside the pipeline ceiling (#8749)"
);
const _: () = assert!(
DEFAULT_WRITE_TXN_SECS < DEFAULT_WRITE_LOCK_SECS,
"the write-transaction deadline must be shorter than the write-lock wait (#8749)"
);
const SLOW_WRITE_WARN_SECS: u64 = 5;
const _: () = assert!(
SLOW_WRITE_WARN_SECS < DEFAULT_WRITE_PIPELINE_SECS,
"the slow-write warning must be reachable below the pipeline ceiling (#6366)"
);
pub fn embedder_init_timeout() -> Duration {
parse_secs_env(
"TRUSTY_EMBEDDER_INIT_TIMEOUT_SECS",
DEFAULT_EMBEDDER_INIT_SECS,
)
}
pub fn embed_batch_timeout() -> Duration {
parse_secs_env("TRUSTY_EMBED_BATCH_TIMEOUT_SECS", DEFAULT_EMBED_BATCH_SECS)
}
pub fn dream_permit_wait_timeout() -> Duration {
parse_secs_env(
"TRUSTY_DREAM_PERMIT_WAIT_SECS",
DEFAULT_DREAM_PERMIT_WAIT_SECS,
)
}
pub fn write_lock_timeout() -> Duration {
parse_secs_env("TRUSTY_WRITE_LOCK_TIMEOUT_SECS", DEFAULT_WRITE_LOCK_SECS)
}
pub fn open_queue_timeout() -> Duration {
parse_secs_env("TRUSTY_OPEN_QUEUE_TIMEOUT_SECS", DEFAULT_OPEN_QUEUE_SECS)
}
pub fn write_op_budget() -> Duration {
parse_secs_env("TRUSTY_WRITE_OP_BUDGET_SECS", DEFAULT_WRITE_OP_BUDGET_SECS)
}
pub fn write_pipeline_timeout() -> Duration {
let configured = parse_secs_env(
"TRUSTY_WRITE_PIPELINE_TIMEOUT_SECS",
DEFAULT_WRITE_PIPELINE_SECS,
);
let floored = floor_write_pipeline(configured);
if floored != configured && first_warning(&PIPELINE_FLOOR_WARNED) {
tracing::warn!(
configured_secs = configured.as_secs(),
floor_secs = WRITE_PIPELINE_FLOOR_SECS,
applied_secs = floored.as_secs(),
"#6366: TRUSTY_WRITE_PIPELINE_TIMEOUT_SECS is below the embedder \
legs the write pipeline contains, so every cold write would fail \
the moment the embedder initialises. Clamped to the compiled-in \
default; raise the value above the floor to set it deliberately"
);
}
floored
}
const WRITE_PIPELINE_FLOOR_SECS: u64 = DEFAULT_EMBEDDER_INIT_SECS + DEFAULT_EMBED_BATCH_SECS;
const _: () = assert!(
DEFAULT_WRITE_PIPELINE_SECS > WRITE_PIPELINE_FLOOR_SECS,
"the compiled-in default must clear the floor it is clamped to (#6366)"
);
pub fn floor_write_pipeline(configured: Duration) -> Duration {
if configured > Duration::from_secs(WRITE_PIPELINE_FLOOR_SECS) {
configured
} else {
Duration::from_secs(DEFAULT_WRITE_PIPELINE_SECS)
}
}
pub fn write_txn_deadline() -> Duration {
let configured = parse_secs_env("TRUSTY_WRITE_TXN_DEADLINE_SECS", DEFAULT_WRITE_TXN_SECS);
let pipeline = write_pipeline_timeout();
let lock_wait = write_lock_timeout();
let applied = cap_write_txn_deadline(configured, pipeline, lock_wait);
if applied != configured && first_warning(&TXN_DEADLINE_WARNED) {
tracing::warn!(
configured_secs = configured.as_secs(),
pipeline_secs = pipeline.as_secs(),
lock_wait_secs = lock_wait.as_secs(),
applied_secs = applied.as_secs(),
"#8749: TRUSTY_WRITE_TXN_DEADLINE_SECS must be above zero and below \
both the write-pipeline ceiling and the write-lock wait; using the \
compiled-in default"
);
}
applied
}
pub fn cap_write_txn_deadline(
configured: Duration,
pipeline: Duration,
lock_wait: Duration,
) -> Duration {
if configured.is_zero() || configured >= pipeline.min(lock_wait) {
Duration::from_secs(DEFAULT_WRITE_TXN_SECS)
} else {
configured
}
}
static PIPELINE_FLOOR_WARNED: AtomicBool = AtomicBool::new(false);
static TXN_DEADLINE_WARNED: AtomicBool = AtomicBool::new(false);
fn first_warning(flag: &AtomicBool) -> bool {
!flag.swap(true, Ordering::Relaxed)
}
pub fn slow_write_warn_threshold() -> Duration {
parse_secs_env("TRUSTY_SLOW_WRITE_WARN_SECS", SLOW_WRITE_WARN_SECS)
}
#[derive(Debug, Clone, Copy)]
pub struct OpBudget {
started: Instant,
total: Duration,
}
impl OpBudget {
#[must_use]
pub fn start(total: Duration) -> Self {
Self {
started: Instant::now(),
total,
}
}
#[must_use]
pub fn start_default() -> Self {
Self::start(write_op_budget())
}
#[must_use]
pub fn remaining(&self) -> Duration {
self.total.saturating_sub(self.started.elapsed())
}
#[must_use]
pub fn leg(&self, configured: Duration) -> Duration {
configured.min(self.remaining())
}
}
pub async fn lock_with_timeout<'a>(
mutex: &'a Arc<Mutex<()>>,
duration: Duration,
label: &str,
) -> anyhow::Result<MutexGuard<'a, ()>> {
tokio::time::timeout(duration, mutex.lock())
.await
.map_err(|_| {
anyhow::Error::new(WriteTimeout::LockWait {
palace: label.to_string(),
waited: duration,
})
})
}
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
pub enum WriteTimeout {
#[error(
"palace '{palace}' write-lock acquisition timed out after {waited:?} \
(issue #906); a previous writer may be stuck — retry or increase \
TRUSTY_WRITE_LOCK_TIMEOUT_SECS"
)]
LockWait { palace: String, waited: Duration },
#[error(
"palace '{}' write pipeline exceeded its {:?} budget after {:?} \
(issue #6366); the palace write mutex has been released so other \
writers proceed. kg.redb is {} bytes — a large store makes commits \
slower; raise TRUSTY_WRITE_PIPELINE_TIMEOUT_SECS if writes on this \
palace are legitimately this slow",
.palace,
.budget,
.elapsed,
.kg_redb_bytes.map_or_else(|| "unknown".to_string(), |b| b.to_string())
)]
PipelineBudget {
palace: String,
budget: Duration,
elapsed: Duration,
kg_redb_bytes: Option<u64>,
},
}
pub fn parse_secs_with(
lookup: impl Fn(&str) -> Option<String>,
key: &str,
default_secs: u64,
) -> Duration {
lookup(key)
.and_then(|v| v.parse::<u64>().ok())
.map(Duration::from_secs)
.unwrap_or_else(|| Duration::from_secs(default_secs))
}
fn parse_secs_env(name: &str, default_secs: u64) -> Duration {
parse_secs_with(|k| std::env::var(k).ok(), name, default_secs)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn parse_secs_with_uses_default_when_absent() {
let t = parse_secs_with(|_| None, "ANY_KEY", DEFAULT_EMBEDDER_INIT_SECS);
assert_eq!(t, Duration::from_secs(DEFAULT_EMBEDDER_INIT_SECS));
}
#[test]
fn parse_secs_with_falls_back_on_bad_value() {
let t = parse_secs_with(
|_| Some("notanumber".to_string()),
"ANY_KEY",
DEFAULT_EMBEDDER_INIT_SECS,
);
assert_eq!(t, Duration::from_secs(DEFAULT_EMBEDDER_INIT_SECS));
}
#[test]
fn parse_secs_with_reads_custom_value() {
let t = parse_secs_with(
|_| Some("5".to_string()),
"ANY_KEY",
DEFAULT_EMBED_BATCH_SECS,
);
assert_eq!(t, Duration::from_secs(5));
}
fn env_lock() -> std::sync::MutexGuard<'static, ()> {
crate::data_dir::ENV_LOCK
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
}
#[test]
fn embedder_init_timeout_default() {
let _guard = env_lock();
unsafe { std::env::remove_var("TRUSTY_EMBEDDER_INIT_TIMEOUT_SECS") };
let t = embedder_init_timeout();
assert_eq!(t, Duration::from_secs(DEFAULT_EMBEDDER_INIT_SECS));
}
#[test]
fn embed_batch_timeout_default() {
let _guard = env_lock();
unsafe { std::env::remove_var("TRUSTY_EMBED_BATCH_TIMEOUT_SECS") };
let t = embed_batch_timeout();
assert_eq!(t, Duration::from_secs(DEFAULT_EMBED_BATCH_SECS));
}
#[test]
fn dream_permit_wait_timeout_default() {
let t = parse_secs_with(
|_| None,
"TRUSTY_DREAM_PERMIT_WAIT_SECS",
DEFAULT_DREAM_PERMIT_WAIT_SECS,
);
assert_eq!(t, Duration::from_secs(DEFAULT_DREAM_PERMIT_WAIT_SECS));
assert_eq!(DEFAULT_DREAM_PERMIT_WAIT_SECS, 30);
}
#[test]
fn write_lock_timeout_default() {
let _guard = env_lock();
unsafe { std::env::remove_var("TRUSTY_WRITE_LOCK_TIMEOUT_SECS") };
let t = write_lock_timeout();
assert_eq!(t, Duration::from_secs(DEFAULT_WRITE_LOCK_SECS));
}
#[test]
fn open_queue_timeout_default() {
let _guard = env_lock();
unsafe { std::env::remove_var("TRUSTY_OPEN_QUEUE_TIMEOUT_SECS") };
let t = open_queue_timeout();
assert_eq!(t, Duration::from_secs(DEFAULT_OPEN_QUEUE_SECS));
}
#[test]
fn write_op_budget_default() {
let _guard = env_lock();
unsafe { std::env::remove_var("TRUSTY_WRITE_OP_BUDGET_SECS") };
let t = write_op_budget();
assert_eq!(t, Duration::from_secs(DEFAULT_WRITE_OP_BUDGET_SECS));
}
#[test]
fn write_pipeline_timeout_default() {
let _guard = env_lock();
unsafe { std::env::remove_var("TRUSTY_WRITE_PIPELINE_TIMEOUT_SECS") };
let t = write_pipeline_timeout();
assert_eq!(t, Duration::from_secs(DEFAULT_WRITE_PIPELINE_SECS));
}
#[test]
fn write_txn_deadline_default() {
let _guard = env_lock();
unsafe { std::env::remove_var("TRUSTY_WRITE_TXN_DEADLINE_SECS") };
assert_eq!(
write_txn_deadline(),
Duration::from_secs(DEFAULT_WRITE_TXN_SECS)
);
unsafe { std::env::set_var("TRUSTY_WRITE_TXN_DEADLINE_SECS", "0") };
let zero = write_txn_deadline();
unsafe { std::env::remove_var("TRUSTY_WRITE_TXN_DEADLINE_SECS") };
assert_eq!(zero, Duration::from_secs(DEFAULT_WRITE_TXN_SECS));
}
#[test]
fn a_txn_deadline_outside_the_pipeline_ceiling_falls_back_to_the_default() {
let default = Duration::from_secs(DEFAULT_WRITE_TXN_SECS);
let pipeline = Duration::from_secs(DEFAULT_WRITE_PIPELINE_SECS);
let lock_wait = pipeline * 4;
for configured in [Duration::ZERO, pipeline, pipeline * 2] {
assert_eq!(
cap_write_txn_deadline(configured, pipeline, lock_wait),
default
);
}
let inside = Duration::from_secs(90);
assert_eq!(cap_write_txn_deadline(inside, pipeline, lock_wait), inside);
}
#[test]
fn a_txn_deadline_at_or_above_the_write_lock_wait_falls_back_to_the_default() {
const TXN: &str = "TRUSTY_WRITE_TXN_DEADLINE_SECS";
const LOCK: &str = "TRUSTY_WRITE_LOCK_TIMEOUT_SECS";
let _guard = env_lock();
unsafe { std::env::remove_var("TRUSTY_WRITE_PIPELINE_TIMEOUT_SECS") };
let default = Duration::from_secs(DEFAULT_WRITE_TXN_SECS);
let mut seen = Vec::new();
for (txn, lock) in [("90", None), ("25", Some("20")), ("45", None)] {
unsafe { std::env::set_var(TXN, txn) };
match lock {
Some(l) => unsafe { std::env::set_var(LOCK, l) },
None => unsafe { std::env::remove_var(LOCK) },
}
seen.push(write_txn_deadline());
}
unsafe { std::env::remove_var(TXN) };
unsafe { std::env::remove_var(LOCK) };
assert_eq!(
seen,
[default, default, Duration::from_secs(45)],
"#8749: the deadline must stay below the write-lock wait"
);
}
#[test]
fn a_zero_pipeline_override_is_clamped_to_the_default() {
assert_eq!(
floor_write_pipeline(Duration::ZERO),
Duration::from_secs(DEFAULT_WRITE_PIPELINE_SECS)
);
}
#[test]
fn an_override_below_the_embedder_legs_is_clamped() {
let default = Duration::from_secs(DEFAULT_WRITE_PIPELINE_SECS);
assert_eq!(floor_write_pipeline(Duration::from_secs(1)), default);
assert_eq!(
floor_write_pipeline(Duration::from_secs(WRITE_PIPELINE_FLOOR_SECS - 1)),
default
);
assert_eq!(
floor_write_pipeline(Duration::from_secs(WRITE_PIPELINE_FLOOR_SECS)),
default
);
}
#[test]
fn an_override_above_the_floor_is_honoured() {
for secs in [WRITE_PIPELINE_FLOOR_SECS + 1, 600, 3600] {
assert_eq!(
floor_write_pipeline(Duration::from_secs(secs)),
Duration::from_secs(secs),
"a {secs}s ceiling clears the floor and must be honoured"
);
}
}
#[test]
fn slow_write_warn_threshold_default() {
let _guard = env_lock();
unsafe { std::env::remove_var("TRUSTY_SLOW_WRITE_WARN_SECS") };
let t = slow_write_warn_threshold();
assert_eq!(t, Duration::from_secs(SLOW_WRITE_WARN_SECS));
}
#[test]
fn budget_clamps_a_leg_to_what_is_left() {
let budget = OpBudget::start(Duration::from_secs(5));
let leg = budget.leg(Duration::from_secs(60));
assert!(
leg <= Duration::from_secs(5),
"a 60s leg under a 5s budget must be clamped to the budget; got {leg:?}"
);
}
#[test]
fn an_exhausted_budget_leaves_a_later_leg_nothing() {
let budget = OpBudget::start(Duration::ZERO);
assert_eq!(budget.remaining(), Duration::ZERO);
assert_eq!(budget.leg(Duration::from_secs(60)), Duration::ZERO);
}
#[test]
fn budget_never_extends_a_shorter_leg() {
let budget = OpBudget::start(Duration::from_secs(60));
assert_eq!(
budget.leg(Duration::from_millis(1)),
Duration::from_millis(1)
);
}
}