use super::{
BTreeMap, ConnectionPool, Duration, HashMap, Instant, Mutex, OnceLock, Ordering, Path, PathBuf,
RawCheckpointObservation, CHECKPOINT_CONSECUTIVE_SKIPS, CHECKPOINT_LAST_SKIP_WAL_PAGES,
CHECKPOINT_LIFECYCLE_APPEND_ATTEMPTS, CHECKPOINT_LIFECYCLE_APPEND_FAILURES,
CHECKPOINT_LIFECYCLE_ENQUEUE_DROPS, CHECKPOINT_PRESSURE_ELEVATED_TICKS,
CHECKPOINT_PRESSURE_EPISODES_RECOVERED, CHECKPOINT_PRESSURE_EPISODES_STARTED,
CHECKPOINT_SKIPPED_TICKS, LAST_WAL_PAGES, READ_TX_MAX_AGE_EVICTIONS, TRUNCATE_ATTEMPTS,
TRUNCATE_CONSECUTIVE_FAILURES,
};
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct RoutineWalObservation {
pub busy: i64,
pub log_frames: u64,
pub checkpointed_frames: u64,
pub pending_frames: u64,
pub physical_wal_bytes: Option<u64>,
pub observed_at_unix_ms: u64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize)]
pub struct CheckpointRun {
pub frame: i64,
pub first_observed_at_unix_ms: u64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum CheckpointRunStatus {
NoTask,
NoObservation,
Observed(CheckpointRun),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(super) struct CheckpointRunEntry {
pub(super) run: CheckpointRun,
first_observed_at: Instant,
pub(super) last_log_frames: i64,
last_informative_at: Instant,
pub(super) busy_since_last_informative: bool,
}
#[derive(Debug, Default)]
pub(super) struct CheckpointRunState {
active_tasks: usize,
pub(super) checkpoint_interval_ms: u64,
owner_intervals_ms: BTreeMap<u64, usize>,
pub(super) entry: Option<CheckpointRunEntry>,
}
static CHECKPOINT_RUNS: OnceLock<Mutex<HashMap<Option<PathBuf>, CheckpointRunState>>> =
OnceLock::new();
pub(super) fn checkpoint_runs() -> &'static Mutex<HashMap<Option<PathBuf>, CheckpointRunState>> {
CHECKPOINT_RUNS.get_or_init(|| Mutex::new(HashMap::new()))
}
pub(crate) struct CheckpointRunTaskGuard {
key: Option<PathBuf>,
interval_ms: u64,
}
impl CheckpointRunTaskGuard {
pub(crate) fn start(pool: &ConnectionPool, interval: Duration) -> Self {
let key = checkpoint_db_key(pool);
let interval_ms = interval.as_millis().min(u128::from(u64::MAX)) as u64;
let interval_ms = interval_ms.max(1);
let mut runs = checkpoint_runs()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let state = runs.entry(key.clone()).or_default();
if state.active_tasks == 0 {
state.entry = None;
}
*state.owner_intervals_ms.entry(interval_ms).or_default() += 1;
state.active_tasks = state.active_tasks.saturating_add(1);
state.checkpoint_interval_ms = *state
.owner_intervals_ms
.first_key_value()
.expect("active owner has an interval")
.0;
Self { key, interval_ms }
}
}
impl Drop for CheckpointRunTaskGuard {
fn drop(&mut self) {
let mut runs = checkpoint_runs()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let Some(state) = runs.get_mut(&self.key) else {
return;
};
let Some(count) = state.owner_intervals_ms.get_mut(&self.interval_ms) else {
return;
};
*count -= 1;
if *count == 0 {
state.owner_intervals_ms.remove(&self.interval_ms);
}
state.active_tasks = state.active_tasks.saturating_sub(1);
if state.active_tasks == 0 {
runs.remove(&self.key);
} else {
state.checkpoint_interval_ms = *state
.owner_intervals_ms
.first_key_value()
.expect("surviving owner has an interval")
.0;
}
}
}
pub(super) fn advance_checkpoint_run_at(
entry: &mut Option<CheckpointRunEntry>,
observation: Option<(i64, i64, i64)>,
observed_at_unix_ms: u64,
observed_at: Instant,
checkpoint_interval_ms: u64,
) {
let Some((busy, log_frames, checkpointed_frames)) = observation else {
*entry = None;
return;
};
if busy != 0 {
if let Some(current) = entry {
current.busy_since_last_informative = true;
}
return;
}
if log_frames < 0 || checkpointed_frames < 0 || checkpointed_frames >= log_frames {
*entry = None;
return;
}
match entry {
Some(current)
if current.run.frame == checkpointed_frames
&& log_frames >= current.last_log_frames
&& (!current.busy_since_last_informative
|| observed_at
.checked_duration_since(current.last_informative_at)
.is_some_and(|elapsed| {
elapsed
<= Duration::from_millis(checkpoint_interval_ms.saturating_mul(2))
})) =>
{
current.last_log_frames = log_frames;
current.last_informative_at = observed_at;
current.busy_since_last_informative = false;
}
_ => {
*entry = Some(CheckpointRunEntry {
run: CheckpointRun {
frame: checkpointed_frames,
first_observed_at_unix_ms: observed_at_unix_ms,
},
first_observed_at: observed_at,
last_log_frames: log_frames,
last_informative_at: observed_at,
busy_since_last_informative: false,
});
}
}
}
#[cfg(test)]
pub(super) fn advance_checkpoint_run(
entry: &mut Option<CheckpointRunEntry>,
observation: Option<(i64, i64, i64)>,
observed_at_unix_ms: u64,
checkpoint_interval_ms: u64,
) {
static TEST_ORIGIN: OnceLock<Instant> = OnceLock::new();
let observed_at = TEST_ORIGIN
.get_or_init(Instant::now)
.checked_add(Duration::from_millis(observed_at_unix_ms))
.expect("test monotonic timestamp");
advance_checkpoint_run_at(
entry,
observation,
observed_at_unix_ms,
observed_at,
checkpoint_interval_ms,
);
}
pub(crate) fn record_checkpoint_run_result(
pool: &ConnectionPool,
observation: Option<(i64, i64, i64)>,
) -> CheckpointRunStatus {
let mut runs = checkpoint_runs()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let Some(state) = runs.get_mut(&checkpoint_db_key(pool)) else {
return CheckpointRunStatus::NoTask;
};
if state.active_tasks == 0 {
return CheckpointRunStatus::NoTask;
}
let checkpoint_interval_ms = state.checkpoint_interval_ms;
advance_checkpoint_run_at(
&mut state.entry,
observation,
observed_at_unix_ms(),
Instant::now(),
checkpoint_interval_ms,
);
state
.entry
.map_or(CheckpointRunStatus::NoObservation, |entry| {
CheckpointRunStatus::Observed(entry.run)
})
}
#[cfg(test)]
pub(crate) fn checkpoint_run_status(pool: &ConnectionPool) -> CheckpointRunStatus {
checkpoint_run_snapshot(pool).0
}
pub(crate) fn checkpoint_run_snapshot(
pool: &ConnectionPool,
) -> (CheckpointRunStatus, Option<Duration>) {
let runs = checkpoint_runs()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let Some(state) = runs.get(&checkpoint_db_key(pool)) else {
return (CheckpointRunStatus::NoTask, None);
};
if state.active_tasks == 0 {
return (CheckpointRunStatus::NoTask, None);
}
state
.entry
.map_or((CheckpointRunStatus::NoObservation, None), |entry| {
(
CheckpointRunStatus::Observed(entry.run),
Some(entry.first_observed_at.elapsed()),
)
})
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
#[serde(default)]
pub struct CheckpointTiming {
pub ticks: u64,
pub elapsed_us_sum: u64,
pub elapsed_us_max: u64,
pub busy_ticks: u64,
pub error_ticks: u64,
}
static CHECKPOINT_TIMINGS: OnceLock<Mutex<HashMap<Option<PathBuf>, CheckpointTiming>>> =
OnceLock::new();
pub(super) fn checkpoint_timings() -> &'static Mutex<HashMap<Option<PathBuf>, CheckpointTiming>> {
CHECKPOINT_TIMINGS.get_or_init(|| Mutex::new(HashMap::new()))
}
pub(super) fn record_checkpoint_timing(pool: &ConnectionPool, elapsed_us: u64, busy: Option<i64>) {
let mut timings = checkpoint_timings()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let timing = timings.entry(checkpoint_db_key(pool)).or_default();
timing.ticks = timing.ticks.saturating_add(1);
timing.elapsed_us_sum = timing.elapsed_us_sum.saturating_add(elapsed_us);
timing.elapsed_us_max = timing.elapsed_us_max.max(elapsed_us);
timing.busy_ticks = timing
.busy_ticks
.saturating_add(u64::from(busy.is_some_and(|value| value != 0)));
timing.error_ticks = timing.error_ticks.saturating_add(u64::from(busy.is_none()));
}
pub fn checkpoint_timing(pool: &ConnectionPool) -> CheckpointTiming {
checkpoint_timings()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.get(&checkpoint_db_key(pool))
.copied()
.unwrap_or_default()
}
static ROUTINE_WAL_OBSERVATIONS: OnceLock<Mutex<HashMap<Option<PathBuf>, RoutineWalObservation>>> =
OnceLock::new();
fn routine_wal_observations() -> &'static Mutex<HashMap<Option<PathBuf>, RoutineWalObservation>> {
ROUTINE_WAL_OBSERVATIONS.get_or_init(|| Mutex::new(HashMap::new()))
}
pub(super) fn checkpoint_db_key_from_path(path: Option<&Path>) -> Option<PathBuf> {
path.map(Path::to_path_buf)
}
pub(super) fn checkpoint_db_key(pool: &ConnectionPool) -> Option<PathBuf> {
checkpoint_db_key_from_path(pool.canonical_path())
}
fn observed_at_unix_ms() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|duration| duration.as_millis() as u64)
.unwrap_or(0)
}
fn physical_wal_bytes(pool: &ConnectionPool) -> Option<u64> {
let path = pool.canonical_path()?;
let mut sidecar = path.as_os_str().to_os_string();
sidecar.push("-wal");
std::fs::metadata(PathBuf::from(sidecar))
.ok()
.map(|metadata| metadata.len())
}
pub(super) fn record_routine_wal_observation(
pool: &ConnectionPool,
raw: RawCheckpointObservation,
) -> RoutineWalObservation {
let log_frames = raw.log_frames.max(0) as u64;
let checkpointed_frames = raw.checkpointed_frames.max(0) as u64;
let observation = RoutineWalObservation {
busy: raw.busy,
log_frames,
checkpointed_frames,
pending_frames: log_frames.saturating_sub(checkpointed_frames),
physical_wal_bytes: physical_wal_bytes(pool),
observed_at_unix_ms: observed_at_unix_ms(),
};
routine_wal_observations()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.insert(checkpoint_db_key(pool), observation.clone());
observation
}
pub fn routine_wal_observation(pool: &ConnectionPool) -> Option<RoutineWalObservation> {
routine_wal_observations()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.get(&checkpoint_db_key(pool))
.cloned()
}
pub fn last_observed_wal_pages() -> Option<u64> {
match LAST_WAL_PAGES.load(Ordering::Relaxed) {
u64::MAX => None,
pages => Some(pages),
}
}
pub fn truncate_attempts() -> u64 {
TRUNCATE_ATTEMPTS.load(Ordering::Relaxed)
}
pub fn truncate_consecutive_failures() -> u64 {
TRUNCATE_CONSECUTIVE_FAILURES.load(Ordering::Relaxed)
}
pub fn checkpoint_skipped_ticks() -> u64 {
CHECKPOINT_SKIPPED_TICKS.load(Ordering::Relaxed)
}
pub fn checkpoint_consecutive_skips() -> u64 {
CHECKPOINT_CONSECUTIVE_SKIPS.load(Ordering::Relaxed)
}
pub fn checkpoint_last_skip_wal_pages() -> Option<u64> {
match CHECKPOINT_LAST_SKIP_WAL_PAGES.load(Ordering::Relaxed) {
u64::MAX => None,
pages => Some(pages),
}
}
pub fn checkpoint_pressure_elevated_ticks() -> u64 {
CHECKPOINT_PRESSURE_ELEVATED_TICKS.load(Ordering::Relaxed)
}
pub fn checkpoint_pressure_episodes_started() -> u64 {
CHECKPOINT_PRESSURE_EPISODES_STARTED.load(Ordering::Relaxed)
}
pub fn checkpoint_pressure_episodes_recovered() -> u64 {
CHECKPOINT_PRESSURE_EPISODES_RECOVERED.load(Ordering::Relaxed)
}
pub fn checkpoint_lifecycle_append_attempts() -> u64 {
CHECKPOINT_LIFECYCLE_APPEND_ATTEMPTS.load(Ordering::Relaxed)
}
pub fn checkpoint_lifecycle_append_failures() -> u64 {
CHECKPOINT_LIFECYCLE_APPEND_FAILURES.load(Ordering::Relaxed)
}
pub fn read_tx_max_age_evictions() -> u64 {
READ_TX_MAX_AGE_EVICTIONS.load(Ordering::Relaxed)
}
pub(crate) fn note_read_tx_max_age_eviction() {
READ_TX_MAX_AGE_EVICTIONS.fetch_add(1, Ordering::Relaxed);
}
pub fn checkpoint_lifecycle_enqueue_drops() -> u64 {
CHECKPOINT_LIFECYCLE_ENQUEUE_DROPS.load(Ordering::Relaxed)
}
pub(super) fn note_checkpoint_skipped() {
CHECKPOINT_SKIPPED_TICKS.fetch_add(1, Ordering::Relaxed);
CHECKPOINT_CONSECUTIVE_SKIPS.fetch_add(1, Ordering::Relaxed);
if let Some(pages) = last_observed_wal_pages() {
CHECKPOINT_LAST_SKIP_WAL_PAGES.store(pages, Ordering::Relaxed);
}
}
pub(super) fn note_checkpoint_observed(_wal_pages: u64) {
CHECKPOINT_CONSECUTIVE_SKIPS.store(0, Ordering::Relaxed);
}
pub(super) fn note_checkpoint_pressure_observation(above_warn: bool, was_above_warn: bool) {
if above_warn {
CHECKPOINT_PRESSURE_ELEVATED_TICKS.fetch_add(1, Ordering::Relaxed);
if !was_above_warn {
CHECKPOINT_PRESSURE_EPISODES_STARTED.fetch_add(1, Ordering::Relaxed);
}
} else if was_above_warn {
CHECKPOINT_PRESSURE_EPISODES_RECOVERED.fetch_add(1, Ordering::Relaxed);
}
}
#[cfg(test)]
pub(crate) fn reset_checkpoint_metrics_for_tests() {
CHECKPOINT_SKIPPED_TICKS.store(0, Ordering::Relaxed);
CHECKPOINT_CONSECUTIVE_SKIPS.store(0, Ordering::Relaxed);
CHECKPOINT_LAST_SKIP_WAL_PAGES.store(u64::MAX, Ordering::Relaxed);
CHECKPOINT_PRESSURE_ELEVATED_TICKS.store(0, Ordering::Relaxed);
CHECKPOINT_PRESSURE_EPISODES_STARTED.store(0, Ordering::Relaxed);
CHECKPOINT_PRESSURE_EPISODES_RECOVERED.store(0, Ordering::Relaxed);
CHECKPOINT_LIFECYCLE_APPEND_ATTEMPTS.store(0, Ordering::Relaxed);
CHECKPOINT_LIFECYCLE_APPEND_FAILURES.store(0, Ordering::Relaxed);
CHECKPOINT_LIFECYCLE_ENQUEUE_DROPS.store(0, Ordering::Relaxed);
READ_TX_MAX_AGE_EVICTIONS.store(0, Ordering::Relaxed);
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum CheckpointTick {
Skipped,
Observed(u64),
}