Skip to main content

khive_db/pool/
admission.rs

1//! Runtime write-route classification and checkpoint ownership admission.
2
3#[cfg(test)]
4use super::Arc;
5use super::{Condvar, Connection, Mutex, SqliteError};
6
7/// Runtime-owned SQL transactions that share the store write-routing policy.
8#[derive(Clone, Copy, Debug)]
9pub enum RuntimeWriteOperation {
10    MergeEntity,
11    MergeNote,
12    UpdateSymmetricEdge,
13}
14
15impl RuntimeWriteOperation {
16    pub(super) fn operation(self) -> &'static str {
17        match self {
18            Self::MergeEntity => "merge_entity",
19            Self::MergeNote => "merge_note",
20            Self::UpdateSymmetricEdge => "update_edge",
21        }
22    }
23
24    pub(super) fn fallback_site(self) -> crate::timeout_sink::Site {
25        match self {
26            Self::MergeEntity => crate::timeout_sink::Site::DirectRouteRuntimeMergeEntity,
27            Self::MergeNote => crate::timeout_sink::Site::DirectRouteRuntimeMergeNote,
28            Self::UpdateSymmetricEdge => {
29                crate::timeout_sink::Site::DirectRouteRuntimeUpdateSymmetricEdge
30            }
31        }
32    }
33}
34
35/// Bounded WAL autocheckpoint applied to writer-capable connections while no
36/// dedicated checkpoint owner has claimed the pool (4,000 pages ≈ 16 MiB at
37/// SQLite's default 4 KiB page size — SQLite's historic behaviour for this
38/// pool). Not a tuning parameter: there is no config field or environment
39/// override, and the only way to change the effective value is an actual
40/// ownership claim ([`ConnectionPool::claim_checkpoint_ownership`](super::ConnectionPool::claim_checkpoint_ownership)), which a
41/// runtime may make only when it really runs the scheduled checkpoint task.
42pub(crate) const FALLBACK_WAL_AUTOCHECKPOINT_PAGES: u32 = 4_000;
43
44#[derive(Clone, Copy, Debug, PartialEq, Eq)]
45pub(super) enum CheckpointOwnership {
46    Unclaimed,
47    Claiming,
48    Claimed,
49}
50
51pub(super) struct CheckpointOwnershipState {
52    pub(super) phase: CheckpointOwnership,
53    #[cfg(test)]
54    pub(super) connection_waiters: usize,
55}
56
57#[cfg(test)]
58pub(super) struct CheckpointConnectionConfigPause {
59    pub(super) selected: std::sync::Barrier,
60    pub(super) resume: std::sync::Barrier,
61}
62
63#[cfg(test)]
64impl CheckpointConnectionConfigPause {
65    pub(super) fn new() -> Self {
66        Self {
67            selected: std::sync::Barrier::new(2),
68            resume: std::sync::Barrier::new(2),
69        }
70    }
71}
72
73pub(super) struct CheckpointOwnershipGate {
74    pub(super) state: Mutex<CheckpointOwnershipState>,
75    pub(super) changed: Condvar,
76    #[cfg(test)]
77    pub(super) connection_config_pause: Mutex<Option<Arc<CheckpointConnectionConfigPause>>>,
78    #[cfg(test)]
79    pub(super) claim_lock_observed: Mutex<Option<std::sync::mpsc::SyncSender<bool>>>,
80}
81
82impl CheckpointOwnershipGate {
83    pub(super) fn new() -> Self {
84        Self {
85            state: Mutex::new(CheckpointOwnershipState {
86                phase: CheckpointOwnership::Unclaimed,
87                #[cfg(test)]
88                connection_waiters: 0,
89            }),
90            changed: Condvar::new(),
91            #[cfg(test)]
92            connection_config_pause: Mutex::new(None),
93            #[cfg(test)]
94            claim_lock_observed: Mutex::new(None),
95        }
96    }
97
98    /// Join an in-flight claim, or become the one caller that configures it.
99    /// Returns `false` when another caller has already completed the claim.
100    pub(super) fn begin_claim(&self) -> bool {
101        #[cfg(test)]
102        let claim_lock_observed = self.claim_lock_observed.lock().take();
103        #[cfg(test)]
104        let mut state = if let Some(observed) = claim_lock_observed {
105            match self.state.try_lock() {
106                Some(state) => {
107                    let _ = observed.send(false);
108                    state
109                }
110                None => {
111                    let _ = observed.send(true);
112                    self.state.lock()
113                }
114            }
115        } else {
116            self.state.lock()
117        };
118        #[cfg(not(test))]
119        let mut state = self.state.lock();
120        loop {
121            match state.phase {
122                CheckpointOwnership::Unclaimed => {
123                    state.phase = CheckpointOwnership::Claiming;
124                    self.changed.notify_all();
125                    return true;
126                }
127                CheckpointOwnership::Claiming => self.changed.wait(&mut state),
128                CheckpointOwnership::Claimed => return false,
129            }
130        }
131    }
132
133    pub(super) fn finish_claim(&self, succeeded: bool) {
134        let mut state = self.state.lock();
135        debug_assert_eq!(state.phase, CheckpointOwnership::Claiming);
136        state.phase = if succeeded {
137            CheckpointOwnership::Claimed
138        } else {
139            CheckpointOwnership::Unclaimed
140        };
141        self.changed.notify_all();
142    }
143
144    fn settled_state(&self) -> parking_lot::MutexGuard<'_, CheckpointOwnershipState> {
145        let mut state = self.state.lock();
146        while state.phase == CheckpointOwnership::Claiming {
147            #[cfg(test)]
148            {
149                state.connection_waiters += 1;
150                self.changed.notify_all();
151            }
152            self.changed.wait(&mut state);
153            #[cfg(test)]
154            {
155                state.connection_waiters -= 1;
156                self.changed.notify_all();
157            }
158        }
159        state
160    }
161
162    #[cfg(test)]
163    pub(super) fn wal_autocheckpoint_pages(&self) -> u32 {
164        let state = self.settled_state();
165        match state.phase {
166            CheckpointOwnership::Unclaimed => FALLBACK_WAL_AUTOCHECKPOINT_PAGES,
167            CheckpointOwnership::Claimed => 0,
168            CheckpointOwnership::Claiming => unreachable!("claim wait must settle the state"),
169        }
170    }
171
172    /// Wait for any in-flight claim, select the resulting posture, and retain
173    /// the gate until SQLite has applied that connection-local PRAGMA. A claim
174    /// therefore linearizes entirely before or after this configuration,
175    /// never between its state sample and side effect.
176    pub(super) fn configure_wal_autocheckpoint(
177        &self,
178        conn: &Connection,
179    ) -> Result<(), SqliteError> {
180        let state = self.settled_state();
181        let pages = match state.phase {
182            CheckpointOwnership::Unclaimed => FALLBACK_WAL_AUTOCHECKPOINT_PAGES,
183            CheckpointOwnership::Claimed => 0,
184            CheckpointOwnership::Claiming => unreachable!("claim wait must settle the state"),
185        };
186        #[cfg(test)]
187        if let Some(pause) = self.connection_config_pause.lock().take() {
188            pause.selected.wait();
189            pause.resume.wait();
190        }
191        conn.pragma_update(None, "wal_autocheckpoint", pages)?;
192        drop(state);
193        Ok(())
194    }
195}