1use super::{
4 BTreeMap, ConnectionPool, Duration, HashMap, Instant, Mutex, OnceLock, Ordering, Path, PathBuf,
5 RawCheckpointObservation, CHECKPOINT_CONSECUTIVE_SKIPS, CHECKPOINT_LAST_SKIP_WAL_PAGES,
6 CHECKPOINT_LIFECYCLE_APPEND_ATTEMPTS, CHECKPOINT_LIFECYCLE_APPEND_FAILURES,
7 CHECKPOINT_LIFECYCLE_ENQUEUE_DROPS, CHECKPOINT_PRESSURE_ELEVATED_TICKS,
8 CHECKPOINT_PRESSURE_EPISODES_RECOVERED, CHECKPOINT_PRESSURE_EPISODES_STARTED,
9 CHECKPOINT_SKIPPED_TICKS, LAST_WAL_PAGES, READ_TX_MAX_AGE_EVICTIONS, TRUNCATE_ATTEMPTS,
10 TRUNCATE_CONSECUTIVE_FAILURES,
11};
12
13#[derive(Debug, Clone, PartialEq, Eq)]
18pub struct RoutineWalObservation {
19 pub busy: i64,
20 pub log_frames: u64,
21 pub checkpointed_frames: u64,
22 pub pending_frames: u64,
23 pub physical_wal_bytes: Option<u64>,
24 pub observed_at_unix_ms: u64,
25}
26
27#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize)]
32pub struct CheckpointRun {
33 pub frame: i64,
35 pub first_observed_at_unix_ms: u64,
37}
38
39#[derive(Debug, Clone, Copy, PartialEq, Eq)]
40pub(crate) enum CheckpointRunStatus {
41 NoTask,
42 NoObservation,
43 Observed(CheckpointRun),
44}
45
46#[derive(Debug, Clone, Copy, PartialEq, Eq)]
47pub(super) struct CheckpointRunEntry {
48 pub(super) run: CheckpointRun,
49 first_observed_at: Instant,
50 pub(super) last_log_frames: i64,
51 last_informative_at: Instant,
52 pub(super) busy_since_last_informative: bool,
53}
54
55#[derive(Debug, Default)]
56pub(super) struct CheckpointRunState {
57 active_tasks: usize,
58 pub(super) checkpoint_interval_ms: u64,
59 owner_intervals_ms: BTreeMap<u64, usize>,
60 pub(super) entry: Option<CheckpointRunEntry>,
61}
62
63static CHECKPOINT_RUNS: OnceLock<Mutex<HashMap<Option<PathBuf>, CheckpointRunState>>> =
64 OnceLock::new();
65
66pub(super) fn checkpoint_runs() -> &'static Mutex<HashMap<Option<PathBuf>, CheckpointRunState>> {
67 CHECKPOINT_RUNS.get_or_init(|| Mutex::new(HashMap::new()))
68}
69
70pub(crate) struct CheckpointRunTaskGuard {
71 key: Option<PathBuf>,
72 interval_ms: u64,
73}
74
75impl CheckpointRunTaskGuard {
76 pub(crate) fn start(pool: &ConnectionPool, interval: Duration) -> Self {
77 let key = checkpoint_db_key(pool);
78 let interval_ms = interval.as_millis().min(u128::from(u64::MAX)) as u64;
79 let interval_ms = interval_ms.max(1);
80 let mut runs = checkpoint_runs()
81 .lock()
82 .unwrap_or_else(std::sync::PoisonError::into_inner);
83 let state = runs.entry(key.clone()).or_default();
84 if state.active_tasks == 0 {
85 state.entry = None;
86 }
87 *state.owner_intervals_ms.entry(interval_ms).or_default() += 1;
88 state.active_tasks = state.active_tasks.saturating_add(1);
89 state.checkpoint_interval_ms = *state
90 .owner_intervals_ms
91 .first_key_value()
92 .expect("active owner has an interval")
93 .0;
94 Self { key, interval_ms }
95 }
96}
97
98impl Drop for CheckpointRunTaskGuard {
99 fn drop(&mut self) {
100 let mut runs = checkpoint_runs()
101 .lock()
102 .unwrap_or_else(std::sync::PoisonError::into_inner);
103 let Some(state) = runs.get_mut(&self.key) else {
104 return;
105 };
106 let Some(count) = state.owner_intervals_ms.get_mut(&self.interval_ms) else {
107 return;
108 };
109 *count -= 1;
110 if *count == 0 {
111 state.owner_intervals_ms.remove(&self.interval_ms);
112 }
113 state.active_tasks = state.active_tasks.saturating_sub(1);
114 if state.active_tasks == 0 {
115 runs.remove(&self.key);
116 } else {
117 state.checkpoint_interval_ms = *state
118 .owner_intervals_ms
119 .first_key_value()
120 .expect("surviving owner has an interval")
121 .0;
122 }
123 }
124}
125
126pub(super) fn advance_checkpoint_run_at(
127 entry: &mut Option<CheckpointRunEntry>,
128 observation: Option<(i64, i64, i64)>,
129 observed_at_unix_ms: u64,
130 observed_at: Instant,
131 checkpoint_interval_ms: u64,
132) {
133 let Some((busy, log_frames, checkpointed_frames)) = observation else {
134 *entry = None;
135 return;
136 };
137 if busy != 0 {
140 if let Some(current) = entry {
141 current.busy_since_last_informative = true;
142 }
143 return;
144 }
145 if log_frames < 0 || checkpointed_frames < 0 || checkpointed_frames >= log_frames {
146 *entry = None;
147 return;
148 }
149
150 match entry {
151 Some(current)
152 if current.run.frame == checkpointed_frames
153 && log_frames >= current.last_log_frames
154 && (!current.busy_since_last_informative
155 || observed_at
156 .checked_duration_since(current.last_informative_at)
157 .is_some_and(|elapsed| {
158 elapsed
159 <= Duration::from_millis(checkpoint_interval_ms.saturating_mul(2))
160 })) =>
161 {
162 current.last_log_frames = log_frames;
163 current.last_informative_at = observed_at;
164 current.busy_since_last_informative = false;
165 }
166 _ => {
167 *entry = Some(CheckpointRunEntry {
168 run: CheckpointRun {
169 frame: checkpointed_frames,
170 first_observed_at_unix_ms: observed_at_unix_ms,
171 },
172 first_observed_at: observed_at,
173 last_log_frames: log_frames,
174 last_informative_at: observed_at,
175 busy_since_last_informative: false,
176 });
177 }
178 }
179}
180
181#[cfg(test)]
182pub(super) fn advance_checkpoint_run(
183 entry: &mut Option<CheckpointRunEntry>,
184 observation: Option<(i64, i64, i64)>,
185 observed_at_unix_ms: u64,
186 checkpoint_interval_ms: u64,
187) {
188 static TEST_ORIGIN: OnceLock<Instant> = OnceLock::new();
189 let observed_at = TEST_ORIGIN
190 .get_or_init(Instant::now)
191 .checked_add(Duration::from_millis(observed_at_unix_ms))
192 .expect("test monotonic timestamp");
193 advance_checkpoint_run_at(
194 entry,
195 observation,
196 observed_at_unix_ms,
197 observed_at,
198 checkpoint_interval_ms,
199 );
200}
201
202pub(crate) fn record_checkpoint_run_result(
203 pool: &ConnectionPool,
204 observation: Option<(i64, i64, i64)>,
205) -> CheckpointRunStatus {
206 let mut runs = checkpoint_runs()
207 .lock()
208 .unwrap_or_else(std::sync::PoisonError::into_inner);
209 let Some(state) = runs.get_mut(&checkpoint_db_key(pool)) else {
210 return CheckpointRunStatus::NoTask;
211 };
212 if state.active_tasks == 0 {
213 return CheckpointRunStatus::NoTask;
214 }
215 let checkpoint_interval_ms = state.checkpoint_interval_ms;
216 advance_checkpoint_run_at(
217 &mut state.entry,
218 observation,
219 observed_at_unix_ms(),
220 Instant::now(),
221 checkpoint_interval_ms,
222 );
223 state
224 .entry
225 .map_or(CheckpointRunStatus::NoObservation, |entry| {
226 CheckpointRunStatus::Observed(entry.run)
227 })
228}
229
230#[cfg(test)]
231pub(crate) fn checkpoint_run_status(pool: &ConnectionPool) -> CheckpointRunStatus {
232 checkpoint_run_snapshot(pool).0
233}
234
235pub(crate) fn checkpoint_run_snapshot(
238 pool: &ConnectionPool,
239) -> (CheckpointRunStatus, Option<Duration>) {
240 let runs = checkpoint_runs()
241 .lock()
242 .unwrap_or_else(std::sync::PoisonError::into_inner);
243 let Some(state) = runs.get(&checkpoint_db_key(pool)) else {
244 return (CheckpointRunStatus::NoTask, None);
245 };
246 if state.active_tasks == 0 {
247 return (CheckpointRunStatus::NoTask, None);
248 }
249 state
250 .entry
251 .map_or((CheckpointRunStatus::NoObservation, None), |entry| {
252 (
253 CheckpointRunStatus::Observed(entry.run),
254 Some(entry.first_observed_at.elapsed()),
255 )
256 })
257}
258
259#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
263#[serde(default)]
264pub struct CheckpointTiming {
265 pub ticks: u64,
266 pub elapsed_us_sum: u64,
267 pub elapsed_us_max: u64,
268 pub busy_ticks: u64,
269 pub error_ticks: u64,
270}
271
272static CHECKPOINT_TIMINGS: OnceLock<Mutex<HashMap<Option<PathBuf>, CheckpointTiming>>> =
273 OnceLock::new();
274
275pub(super) fn checkpoint_timings() -> &'static Mutex<HashMap<Option<PathBuf>, CheckpointTiming>> {
276 CHECKPOINT_TIMINGS.get_or_init(|| Mutex::new(HashMap::new()))
277}
278
279pub(super) fn record_checkpoint_timing(pool: &ConnectionPool, elapsed_us: u64, busy: Option<i64>) {
280 let mut timings = checkpoint_timings()
281 .lock()
282 .unwrap_or_else(std::sync::PoisonError::into_inner);
283 let timing = timings.entry(checkpoint_db_key(pool)).or_default();
284 timing.ticks = timing.ticks.saturating_add(1);
285 timing.elapsed_us_sum = timing.elapsed_us_sum.saturating_add(elapsed_us);
286 timing.elapsed_us_max = timing.elapsed_us_max.max(elapsed_us);
287 timing.busy_ticks = timing
288 .busy_ticks
289 .saturating_add(u64::from(busy.is_some_and(|value| value != 0)));
290 timing.error_ticks = timing.error_ticks.saturating_add(u64::from(busy.is_none()));
291}
292
293pub fn checkpoint_timing(pool: &ConnectionPool) -> CheckpointTiming {
295 checkpoint_timings()
296 .lock()
297 .unwrap_or_else(std::sync::PoisonError::into_inner)
298 .get(&checkpoint_db_key(pool))
299 .copied()
300 .unwrap_or_default()
301}
302
303static ROUTINE_WAL_OBSERVATIONS: OnceLock<Mutex<HashMap<Option<PathBuf>, RoutineWalObservation>>> =
307 OnceLock::new();
308
309fn routine_wal_observations() -> &'static Mutex<HashMap<Option<PathBuf>, RoutineWalObservation>> {
310 ROUTINE_WAL_OBSERVATIONS.get_or_init(|| Mutex::new(HashMap::new()))
311}
312
313pub(super) fn checkpoint_db_key_from_path(path: Option<&Path>) -> Option<PathBuf> {
314 path.map(Path::to_path_buf)
315}
316
317pub(super) fn checkpoint_db_key(pool: &ConnectionPool) -> Option<PathBuf> {
318 checkpoint_db_key_from_path(pool.canonical_path())
319}
320
321fn observed_at_unix_ms() -> u64 {
322 std::time::SystemTime::now()
323 .duration_since(std::time::UNIX_EPOCH)
324 .map(|duration| duration.as_millis() as u64)
325 .unwrap_or(0)
326}
327
328fn physical_wal_bytes(pool: &ConnectionPool) -> Option<u64> {
329 let path = pool.canonical_path()?;
330 let mut sidecar = path.as_os_str().to_os_string();
331 sidecar.push("-wal");
332 std::fs::metadata(PathBuf::from(sidecar))
333 .ok()
334 .map(|metadata| metadata.len())
335}
336
337pub(super) fn record_routine_wal_observation(
338 pool: &ConnectionPool,
339 raw: RawCheckpointObservation,
340) -> RoutineWalObservation {
341 let log_frames = raw.log_frames.max(0) as u64;
342 let checkpointed_frames = raw.checkpointed_frames.max(0) as u64;
343 let observation = RoutineWalObservation {
344 busy: raw.busy,
345 log_frames,
346 checkpointed_frames,
347 pending_frames: log_frames.saturating_sub(checkpointed_frames),
348 physical_wal_bytes: physical_wal_bytes(pool),
349 observed_at_unix_ms: observed_at_unix_ms(),
350 };
351 routine_wal_observations()
352 .lock()
353 .unwrap_or_else(std::sync::PoisonError::into_inner)
354 .insert(checkpoint_db_key(pool), observation.clone());
355 observation
356}
357
358pub fn routine_wal_observation(pool: &ConnectionPool) -> Option<RoutineWalObservation> {
361 routine_wal_observations()
362 .lock()
363 .unwrap_or_else(std::sync::PoisonError::into_inner)
364 .get(&checkpoint_db_key(pool))
365 .cloned()
366}
367
368pub fn last_observed_wal_pages() -> Option<u64> {
371 match LAST_WAL_PAGES.load(Ordering::Relaxed) {
372 u64::MAX => None,
373 pages => Some(pages),
374 }
375}
376
377pub fn truncate_attempts() -> u64 {
379 TRUNCATE_ATTEMPTS.load(Ordering::Relaxed)
380}
381
382pub fn truncate_consecutive_failures() -> u64 {
384 TRUNCATE_CONSECUTIVE_FAILURES.load(Ordering::Relaxed)
385}
386
387pub fn checkpoint_skipped_ticks() -> u64 {
390 CHECKPOINT_SKIPPED_TICKS.load(Ordering::Relaxed)
391}
392
393pub fn checkpoint_consecutive_skips() -> u64 {
395 CHECKPOINT_CONSECUTIVE_SKIPS.load(Ordering::Relaxed)
396}
397
398pub fn checkpoint_last_skip_wal_pages() -> Option<u64> {
401 match CHECKPOINT_LAST_SKIP_WAL_PAGES.load(Ordering::Relaxed) {
402 u64::MAX => None,
403 pages => Some(pages),
404 }
405}
406
407pub fn checkpoint_pressure_elevated_ticks() -> u64 {
409 CHECKPOINT_PRESSURE_ELEVATED_TICKS.load(Ordering::Relaxed)
410}
411
412pub fn checkpoint_pressure_episodes_started() -> u64 {
414 CHECKPOINT_PRESSURE_EPISODES_STARTED.load(Ordering::Relaxed)
415}
416
417pub fn checkpoint_pressure_episodes_recovered() -> u64 {
419 CHECKPOINT_PRESSURE_EPISODES_RECOVERED.load(Ordering::Relaxed)
420}
421
422pub fn checkpoint_lifecycle_append_attempts() -> u64 {
424 CHECKPOINT_LIFECYCLE_APPEND_ATTEMPTS.load(Ordering::Relaxed)
425}
426
427pub fn checkpoint_lifecycle_append_failures() -> u64 {
429 CHECKPOINT_LIFECYCLE_APPEND_FAILURES.load(Ordering::Relaxed)
430}
431
432pub fn read_tx_max_age_evictions() -> u64 {
435 READ_TX_MAX_AGE_EVICTIONS.load(Ordering::Relaxed)
436}
437
438pub(crate) fn note_read_tx_max_age_eviction() {
444 READ_TX_MAX_AGE_EVICTIONS.fetch_add(1, Ordering::Relaxed);
445}
446
447pub fn checkpoint_lifecycle_enqueue_drops() -> u64 {
449 CHECKPOINT_LIFECYCLE_ENQUEUE_DROPS.load(Ordering::Relaxed)
450}
451
452pub(super) fn note_checkpoint_skipped() {
457 CHECKPOINT_SKIPPED_TICKS.fetch_add(1, Ordering::Relaxed);
458 CHECKPOINT_CONSECUTIVE_SKIPS.fetch_add(1, Ordering::Relaxed);
459 if let Some(pages) = last_observed_wal_pages() {
460 CHECKPOINT_LAST_SKIP_WAL_PAGES.store(pages, Ordering::Relaxed);
461 }
462}
463
464pub(super) fn note_checkpoint_observed(_wal_pages: u64) {
469 CHECKPOINT_CONSECUTIVE_SKIPS.store(0, Ordering::Relaxed);
470}
471
472pub(super) fn note_checkpoint_pressure_observation(above_warn: bool, was_above_warn: bool) {
473 if above_warn {
474 CHECKPOINT_PRESSURE_ELEVATED_TICKS.fetch_add(1, Ordering::Relaxed);
475 if !was_above_warn {
476 CHECKPOINT_PRESSURE_EPISODES_STARTED.fetch_add(1, Ordering::Relaxed);
477 }
478 } else if was_above_warn {
479 CHECKPOINT_PRESSURE_EPISODES_RECOVERED.fetch_add(1, Ordering::Relaxed);
480 }
481}
482
483#[cfg(test)]
487pub(crate) fn reset_checkpoint_metrics_for_tests() {
488 CHECKPOINT_SKIPPED_TICKS.store(0, Ordering::Relaxed);
489 CHECKPOINT_CONSECUTIVE_SKIPS.store(0, Ordering::Relaxed);
490 CHECKPOINT_LAST_SKIP_WAL_PAGES.store(u64::MAX, Ordering::Relaxed);
491 CHECKPOINT_PRESSURE_ELEVATED_TICKS.store(0, Ordering::Relaxed);
492 CHECKPOINT_PRESSURE_EPISODES_STARTED.store(0, Ordering::Relaxed);
493 CHECKPOINT_PRESSURE_EPISODES_RECOVERED.store(0, Ordering::Relaxed);
494 CHECKPOINT_LIFECYCLE_APPEND_ATTEMPTS.store(0, Ordering::Relaxed);
495 CHECKPOINT_LIFECYCLE_APPEND_FAILURES.store(0, Ordering::Relaxed);
496 CHECKPOINT_LIFECYCLE_ENQUEUE_DROPS.store(0, Ordering::Relaxed);
497 READ_TX_MAX_AGE_EVICTIONS.store(0, Ordering::Relaxed);
498}
499
500#[derive(Debug, Clone, Copy, PartialEq, Eq)]
508pub enum CheckpointTick {
509 Skipped,
511 Observed(u64),
513}