khive_db/pool/
admission.rs1#[cfg(test)]
4use super::Arc;
5use super::{Condvar, Connection, Mutex, SqliteError};
6
7#[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
35pub(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 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 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}