diesel/sqlite/connection/hooks.rs
1use super::SqliteConnection;
2use core::num::NonZeroU32;
3
4pub(super) use super::{BusyDecision, CommitDecision, ProgressDecision};
5
6impl SqliteConnection {
7 /// Registers a callback invoked when a transaction is about to be
8 /// committed.
9 ///
10 /// The callback returns a [`CommitDecision`]: `Proceed` lets the commit
11 /// complete, `Rollback` converts it into a rollback.
12 ///
13 /// Only one commit hook can be active at a time per connection.
14 /// Registering a new one replaces the previous.
15 ///
16 /// The callback runs synchronously as part of the committing
17 /// `sqlite3_step()` call, on the thread performing the commit, so it is
18 /// never invoked concurrently. Per SQLite, the callback must not use the
19 /// connection that triggered it (running any SQL, including a `SELECT`,
20 /// counts as use) and is not reentrant. A panic in the callback aborts the
21 /// process.
22 ///
23 /// See: [`sqlite3_commit_hook`](https://www.sqlite.org/c3ref/commit_hook.html)
24 ///
25 /// # Example
26 ///
27 /// ```rust
28 /// use diesel::prelude::*;
29 /// use diesel::sqlite::{SqliteConnection, CommitDecision};
30 /// use std::sync::{Arc, Mutex};
31 ///
32 /// diesel::table! {
33 /// users (id) {
34 /// id -> Integer,
35 /// name -> Text,
36 /// }
37 /// }
38 ///
39 /// # use diesel::connection::SimpleConnection;
40 /// # let conn = &mut SqliteConnection::establish(":memory:").unwrap();
41 /// # conn.batch_execute("CREATE TABLE users (id INTEGER PRIMARY KEY, name TEXT NOT NULL)").unwrap();
42 /// let commits = Arc::new(Mutex::new(0u32));
43 /// let commits2 = commits.clone();
44 ///
45 /// conn.on_commit(move || {
46 /// *commits2.lock().unwrap() += 1;
47 /// CommitDecision::Proceed
48 /// });
49 ///
50 /// conn.immediate_transaction(|conn| {
51 /// diesel::insert_into(users::table)
52 /// .values(users::name.eq("Alice"))
53 /// .execute(conn)?;
54 /// Ok::<_, diesel::result::Error>(())
55 /// }).unwrap();
56 ///
57 /// assert_eq!(*commits.lock().unwrap(), 1);
58 /// ```
59 pub fn on_commit<F>(&mut self, hook: F)
60 where
61 F: FnMut() -> CommitDecision + Send + 'static,
62 {
63 self.raw_connection.set_commit_hook(hook);
64 }
65
66 /// Removes the commit hook. Subsequent commits will not invoke any
67 /// callback.
68 ///
69 /// See [`on_commit`](Self::on_commit) for usage example.
70 pub fn remove_commit_hook(&mut self) {
71 self.raw_connection.remove_commit_hook();
72 }
73
74 /// Registers a callback invoked after a transaction is rolled back.
75 ///
76 /// This is **not** invoked for the implicit rollback that occurs when
77 /// the connection is closed. It **is** invoked when a commit hook forces
78 /// a rollback by returning [`CommitDecision::Rollback`].
79 ///
80 /// Only one rollback hook can be active at a time per connection.
81 /// Registering a new one replaces the previous.
82 ///
83 /// The callback must not use the database connection. It is invoked
84 /// synchronously on the thread driving the connection, so it is never
85 /// called concurrently, and like the commit hook it is not reentrant.
86 /// Panics in the callback abort the process.
87 ///
88 /// See: [`sqlite3_rollback_hook`](https://www.sqlite.org/c3ref/commit_hook.html)
89 ///
90 /// # Example
91 ///
92 /// ```rust
93 /// use diesel::prelude::*;
94 /// use diesel::sqlite::SqliteConnection;
95 /// use std::sync::Arc;
96 /// use std::sync::atomic::{AtomicU32, Ordering};
97 ///
98 /// # use diesel::connection::SimpleConnection;
99 /// # let conn = &mut SqliteConnection::establish(":memory:").unwrap();
100 /// # conn.batch_execute("CREATE TABLE users (id INTEGER PRIMARY KEY, name TEXT NOT NULL)").unwrap();
101 /// let rollbacks = Arc::new(AtomicU32::new(0));
102 /// let rb2 = rollbacks.clone();
103 ///
104 /// conn.on_rollback(move || {
105 /// rb2.fetch_add(1, Ordering::Relaxed);
106 /// });
107 ///
108 /// // Force a rollback by returning an error.
109 /// let _ = conn.immediate_transaction(|_conn| {
110 /// Err::<(), _>(diesel::result::Error::RollbackTransaction)
111 /// });
112 ///
113 /// assert_eq!(rollbacks.load(Ordering::Relaxed), 1);
114 /// ```
115 pub fn on_rollback<F>(&mut self, hook: F)
116 where
117 F: FnMut() + Send + 'static,
118 {
119 self.raw_connection.set_rollback_hook(hook);
120 }
121
122 /// Removes the rollback hook. Subsequent rollbacks will not invoke any
123 /// callback.
124 ///
125 /// See [`on_rollback`](Self::on_rollback) for usage example.
126 pub fn remove_rollback_hook(&mut self) {
127 self.raw_connection.remove_rollback_hook();
128 }
129
130 /// Registers a progress handler that can interrupt long-running queries.
131 ///
132 /// The callback is invoked periodically while a query runs. `n` is the
133 /// approximate number of virtual-machine instructions between callbacks. It
134 /// is a [`NonZeroU32`] so the handler cannot be disabled implicitly by
135 /// passing zero. Use [`remove_progress_handler`](Self::remove_progress_handler)
136 /// to disable it. Since SQLite 3.41.0 the callback may also fire during
137 /// statement preparation.
138 ///
139 /// The callback returns a [`ProgressDecision`]: `Continue` lets the query
140 /// keep executing, `Interrupt` aborts it (causes `SQLITE_INTERRUPT`).
141 ///
142 /// Only one progress handler can be active at a time per connection.
143 /// Registering a new one replaces the previous.
144 ///
145 /// The callback must not use the database connection. It is invoked
146 /// synchronously on the thread driving the connection, so it is never
147 /// called concurrently. Panics in the callback abort the process.
148 ///
149 /// See: [`sqlite3_progress_handler`](https://www.sqlite.org/c3ref/progress_handler.html)
150 ///
151 /// # Example
152 ///
153 /// ```rust
154 /// use diesel::prelude::*;
155 /// use diesel::sqlite::{SqliteConnection, ProgressDecision};
156 /// use std::num::NonZeroU32;
157 /// use std::sync::Arc;
158 /// use std::sync::atomic::{AtomicBool, Ordering};
159 ///
160 /// # let conn = &mut SqliteConnection::establish(":memory:").unwrap();
161 /// let cancelled = Arc::new(AtomicBool::new(false));
162 /// let cancelled2 = cancelled.clone();
163 ///
164 /// conn.on_progress(NonZeroU32::new(1000).unwrap(), move || {
165 /// if cancelled2.load(Ordering::Relaxed) {
166 /// ProgressDecision::Interrupt
167 /// } else {
168 /// ProgressDecision::Continue
169 /// }
170 /// });
171 ///
172 /// // Later: remove the handler.
173 /// conn.remove_progress_handler();
174 /// ```
175 pub fn on_progress<F>(&mut self, n: NonZeroU32, hook: F)
176 where
177 F: FnMut() -> ProgressDecision + Send + 'static,
178 {
179 self.raw_connection.set_progress_handler(n, hook);
180 }
181
182 /// Removes the progress handler. Subsequent queries will not invoke any
183 /// callback.
184 ///
185 /// See [`on_progress`](Self::on_progress) for usage example.
186 pub fn remove_progress_handler(&mut self) {
187 self.raw_connection.remove_progress_handler();
188 }
189
190 /// Registers a callback invoked after each commit of a transaction in
191 /// [WAL mode](https://www.sqlite.org/wal.html). It receives a borrowed
192 /// connection, the database name (`"main"`, `"temp"`, or an `ATTACH`
193 /// alias), and the current WAL page count.
194 ///
195 /// The hook fires after the commit completes and the write-lock is
196 /// released, so the callback may read, write, or
197 /// [checkpoint](https://www.sqlite.org/wal.html#ckpt) through `conn`,
198 /// provided it leaves no open transaction. A write that commits inside the
199 /// callback re-fires the hook re-entrantly, so a callback that writes
200 /// unconditionally will recurse until the stack overflows. Guard against
201 /// that yourself if the callback writes.
202 ///
203 /// Only one WAL hook is active at a time, and re-registering replaces it.
204 /// `PRAGMA wal_autocheckpoint` installs its own WAL hook and overwrites
205 /// this one. A panic in the callback aborts the process.
206 ///
207 /// See: [`sqlite3_wal_hook`](https://www.sqlite.org/c3ref/wal_hook.html)
208 ///
209 /// # Example
210 ///
211 /// ```rust,no_run
212 /// # use diesel::prelude::*;
213 /// # use diesel::connection::SimpleConnection;
214 /// # use diesel::sqlite::SqliteConnection;
215 /// # use diesel::sqlite::WalCheckpointMode;
216 /// # let conn = &mut SqliteConnection::establish("test.db").unwrap();
217 /// # conn.batch_execute("PRAGMA journal_mode = WAL;").unwrap();
218 /// conn.on_wal(|conn, db_name, n_pages| {
219 /// println!("WAL for {db_name}: {n_pages} pages");
220 /// if n_pages > 1000 {
221 /// // The connection may be used here, e.g. to force a checkpoint.
222 /// let _ = conn.wal_checkpoint(None, WalCheckpointMode::Truncate);
223 /// }
224 /// });
225 /// ```
226 pub fn on_wal<F>(&mut self, hook: F)
227 where
228 F: Fn(&mut SqliteConnection, &str, u32) + Send + 'static,
229 {
230 self.raw_connection.set_wal_hook(hook);
231 }
232
233 /// Removes the WAL hook. Subsequent commits will not invoke any callback.
234 ///
235 /// See [`on_wal`](Self::on_wal) for usage example.
236 pub fn remove_wal_hook(&mut self) {
237 self.raw_connection.remove_wal_hook();
238 }
239
240 /// Registers a custom busy handler for lock contention.
241 ///
242 /// The callback receives the retry count (starting from 0) and returns a
243 /// [`BusyDecision`]: `Retry` retries the locked operation, `GiveUp` aborts
244 /// and returns `SQLITE_BUSY` to the caller.
245 ///
246 /// Setting this clears any timeout previously set with
247 /// [`set_busy_timeout`](Self::set_busy_timeout). Conversely, calling
248 /// `set_busy_timeout` clears this handler. Only one busy handler can be
249 /// active at a time per connection.
250 ///
251 /// The callback must not use the database connection. If the callback
252 /// modifies the database, behavior is undefined. SQLite may return
253 /// `SQLITE_BUSY` instead of calling the handler to prevent deadlocks.
254 ///
255 /// The callback is invoked synchronously on the thread driving the
256 /// connection, so it is never called concurrently, and per SQLite it is
257 /// not reentrant.
258 ///
259 /// Panics in the callback abort the process.
260 ///
261 /// See: [`sqlite3_busy_handler`](https://www.sqlite.org/c3ref/busy_handler.html)
262 ///
263 /// # Example
264 ///
265 /// ```rust
266 /// use diesel::prelude::*;
267 /// use diesel::sqlite::{SqliteConnection, BusyDecision};
268 /// use std::thread;
269 /// use std::time::Duration;
270 ///
271 /// # let conn = &mut SqliteConnection::establish(":memory:").unwrap();
272 /// conn.on_busy(|retry_count| {
273 /// if retry_count < 5 {
274 /// thread::sleep(Duration::from_millis(100));
275 /// BusyDecision::Retry
276 /// } else {
277 /// BusyDecision::GiveUp
278 /// }
279 /// });
280 ///
281 /// // Later: remove the handler
282 /// conn.remove_busy_handler();
283 /// ```
284 pub fn on_busy<F>(&mut self, hook: F)
285 where
286 F: FnMut(i32) -> BusyDecision + Send + 'static,
287 {
288 self.raw_connection.set_busy_handler(hook);
289 }
290
291 /// Removes the custom busy handler.
292 ///
293 /// See [`on_busy`](Self::on_busy) for usage example.
294 pub fn remove_busy_handler(&mut self) {
295 self.raw_connection.remove_busy_handler();
296 }
297
298 /// Sets a simple timeout-based busy handler.
299 ///
300 /// When a table is locked, SQLite will sleep and retry until `ms`
301 /// milliseconds have elapsed. Pass 0 to disable (return `SQLITE_BUSY`
302 /// immediately).
303 ///
304 /// Setting this clears any custom [`on_busy`](Self::on_busy) handler.
305 /// Conversely, calling `on_busy` clears this timeout. For most use cases,
306 /// this is simpler than a custom busy handler.
307 ///
308 /// See: [`sqlite3_busy_timeout`](https://www.sqlite.org/c3ref/busy_timeout.html)
309 ///
310 /// # Example
311 ///
312 /// ```rust
313 /// use diesel::prelude::*;
314 /// use diesel::sqlite::SqliteConnection;
315 ///
316 /// # let conn = &mut SqliteConnection::establish(":memory:").unwrap();
317 /// // Wait up to 5 seconds for locked tables
318 /// conn.set_busy_timeout(5000);
319 /// ```
320 pub fn set_busy_timeout(&mut self, ms: i32) {
321 self.raw_connection.set_busy_timeout(ms);
322 }
323}
324
325#[cfg(test)]
326mod tests {
327 use super::*;
328 use crate::connection::Connection;
329 use crate::query_dsl::RunQueryDsl;
330 use std::sync::Arc;
331 use std::sync::atomic::{AtomicU32, Ordering};
332
333 fn connection() -> SqliteConnection {
334 SqliteConnection::establish(":memory:").unwrap()
335 }
336
337 // Only used by tests that are compiled out under miri (see the
338 // individual tests below).
339 #[cfg(not(miri))]
340 #[derive(crate::QueryableByName)]
341 struct CountResult {
342 #[diesel(sql_type = crate::sql_types::BigInt)]
343 c: i64,
344 }
345
346 // Registers a callback that is invoked by the native library, which is not
347 // supported when running under miri with a native libsqlite3 (`-Zmiri-native-lib`).
348 #[cfg(not(miri))] // ffi call
349 #[diesel_test_helper::test]
350 fn on_commit_fires_on_commit() {
351 let conn = &mut connection();
352
353 let count = Arc::new(AtomicU32::new(0));
354 let c2 = count.clone();
355
356 conn.on_commit(move || {
357 c2.fetch_add(1, Ordering::Relaxed);
358 CommitDecision::Proceed
359 });
360
361 conn.immediate_transaction(|conn| {
362 crate::sql_query("CREATE TABLE t1 (id INTEGER PRIMARY KEY)")
363 .execute(conn)
364 .unwrap();
365 Ok::<_, crate::result::Error>(())
366 })
367 .unwrap();
368
369 assert_eq!(count.load(Ordering::Relaxed), 1);
370 }
371
372 // Registers a callback that is invoked by the native library, which is not
373 // supported when running under miri with a native libsqlite3 (`-Zmiri-native-lib`).
374 #[cfg(not(miri))] // ffi call
375 #[diesel_test_helper::test]
376 fn on_commit_returning_true_forces_rollback() {
377 let conn = &mut connection();
378
379 crate::sql_query("CREATE TABLE t_commit (id INTEGER PRIMARY KEY)")
380 .execute(conn)
381 .unwrap();
382
383 conn.on_commit(|| CommitDecision::Rollback);
384
385 // The transaction will attempt to commit, but the hook will convert
386 // it to a rollback. diesel's AnsiTransactionManager will see the
387 // failure from the COMMIT statement (sqlite returns an error when
388 // the commit hook returns non-zero and the commit is aborted).
389 let result = conn.immediate_transaction(|conn| {
390 crate::sql_query("INSERT INTO t_commit (id) VALUES (1)")
391 .execute(conn)
392 .unwrap();
393 Ok::<_, crate::result::Error>(())
394 });
395
396 // The transaction should have been rolled back.
397 assert!(
398 matches!(result, Err(crate::result::Error::DatabaseError(..))),
399 "{result:?}"
400 );
401 assert_eq!(
402 Ok(None),
403 <crate::connection::AnsiTransactionManager as crate::connection::TransactionManager<
404 SqliteConnection,
405 >>::transaction_manager_status_mut(conn)
406 .transaction_depth()
407 );
408
409 // Remove the hook so subsequent queries don't fail.
410 conn.remove_commit_hook();
411 conn.immediate_transaction(|_| Ok::<_, crate::result::Error>(()))
412 .unwrap();
413
414 // Verify the row was not persisted.
415 let cnt: i64 = crate::sql_query("SELECT COUNT(*) as c FROM t_commit")
416 .get_result::<CountResult>(conn)
417 .unwrap()
418 .c;
419 assert_eq!(cnt, 0);
420 }
421
422 // Registers a callback that is invoked by the native library, which is not
423 // supported when running under miri with a native libsqlite3 (`-Zmiri-native-lib`).
424 #[cfg(not(miri))] // ffi call
425 #[diesel_test_helper::test]
426 fn replacing_commit_hook_drops_old() {
427 let conn = &mut connection();
428
429 let old_count = Arc::new(AtomicU32::new(0));
430 let new_count = Arc::new(AtomicU32::new(0));
431 let oc = old_count.clone();
432 let nc = new_count.clone();
433
434 conn.on_commit(move || {
435 oc.fetch_add(1, Ordering::Relaxed);
436 CommitDecision::Proceed
437 });
438
439 // Replace with a new hook.
440 conn.on_commit(move || {
441 nc.fetch_add(1, Ordering::Relaxed);
442 CommitDecision::Proceed
443 });
444
445 conn.immediate_transaction(|conn| {
446 crate::sql_query("CREATE TABLE t_replace (id INTEGER PRIMARY KEY)")
447 .execute(conn)
448 .unwrap();
449 Ok::<_, crate::result::Error>(())
450 })
451 .unwrap();
452
453 assert_eq!(old_count.load(Ordering::Relaxed), 0);
454 assert_eq!(new_count.load(Ordering::Relaxed), 1);
455 }
456
457 #[diesel_test_helper::test]
458 fn remove_commit_hook_disables_callback() {
459 let conn = &mut connection();
460
461 let count = Arc::new(AtomicU32::new(0));
462 let c2 = count.clone();
463
464 conn.on_commit(move || {
465 c2.fetch_add(1, Ordering::Relaxed);
466 CommitDecision::Proceed
467 });
468
469 conn.remove_commit_hook();
470
471 conn.immediate_transaction(|conn| {
472 crate::sql_query("CREATE TABLE t_rem (id INTEGER PRIMARY KEY)")
473 .execute(conn)
474 .unwrap();
475 Ok::<_, crate::result::Error>(())
476 })
477 .unwrap();
478
479 assert_eq!(count.load(Ordering::Relaxed), 0);
480 }
481
482 // Registers a callback that is invoked by the native library, which is not
483 // supported when running under miri with a native libsqlite3 (`-Zmiri-native-lib`).
484 #[cfg(not(miri))] // ffi call
485 #[diesel_test_helper::test]
486 fn on_rollback_fires_on_explicit_rollback() {
487 let conn = &mut connection();
488
489 crate::sql_query("CREATE TABLE t_rb (id INTEGER PRIMARY KEY)")
490 .execute(conn)
491 .unwrap();
492
493 let count = Arc::new(AtomicU32::new(0));
494 let c2 = count.clone();
495
496 conn.on_rollback(move || {
497 c2.fetch_add(1, Ordering::Relaxed);
498 });
499
500 // Force a rollback by returning Err from the transaction closure.
501 let _ = conn.immediate_transaction(|conn| {
502 crate::sql_query("INSERT INTO t_rb (id) VALUES (1)")
503 .execute(conn)
504 .unwrap();
505 Err::<(), _>(crate::result::Error::RollbackTransaction)
506 });
507
508 assert_eq!(count.load(Ordering::Relaxed), 1);
509 }
510
511 // Registers a callback that is invoked by the native library, which is not
512 // supported when running under miri with a native libsqlite3 (`-Zmiri-native-lib`).
513 #[cfg(not(miri))] // ffi call
514 #[diesel_test_helper::test]
515 fn on_rollback_fires_when_commit_hook_forces_rollback() {
516 let conn = &mut connection();
517
518 crate::sql_query("CREATE TABLE t_rb2 (id INTEGER PRIMARY KEY)")
519 .execute(conn)
520 .unwrap();
521
522 let rb_count = Arc::new(AtomicU32::new(0));
523 let rb2 = rb_count.clone();
524
525 conn.on_commit(|| CommitDecision::Rollback);
526 conn.on_rollback(move || {
527 rb2.fetch_add(1, Ordering::Relaxed);
528 });
529
530 let _ = conn.immediate_transaction(|conn| {
531 crate::sql_query("INSERT INTO t_rb2 (id) VALUES (1)")
532 .execute(conn)
533 .unwrap();
534 Ok::<_, crate::result::Error>(())
535 });
536
537 // Rollback hook should have fired.
538 assert_eq!(rb_count.load(Ordering::Relaxed), 1);
539
540 conn.remove_commit_hook();
541 conn.remove_rollback_hook();
542
543 // Verify the row was not persisted.
544 let cnt: i64 = crate::sql_query("SELECT COUNT(*) as c FROM t_rb2")
545 .get_result::<CountResult>(conn)
546 .unwrap()
547 .c;
548 assert_eq!(cnt, 0);
549 }
550
551 #[diesel_test_helper::test]
552 fn on_rollback_does_not_fire_on_connection_close() {
553 let count = Arc::new(AtomicU32::new(0));
554 let c2 = count.clone();
555
556 {
557 let conn = &mut connection();
558 conn.on_rollback(move || {
559 c2.fetch_add(1, Ordering::Relaxed);
560 });
561 // conn is dropped here: implicit close, not a rollback.
562 }
563
564 assert_eq!(count.load(Ordering::Relaxed), 0);
565 }
566
567 #[diesel_test_helper::test]
568 fn remove_rollback_hook_disables_callback() {
569 let conn = &mut connection();
570
571 crate::sql_query("CREATE TABLE t_rem_rb (id INTEGER PRIMARY KEY)")
572 .execute(conn)
573 .unwrap();
574
575 let count = Arc::new(AtomicU32::new(0));
576 let c2 = count.clone();
577
578 conn.on_rollback(move || {
579 c2.fetch_add(1, Ordering::Relaxed);
580 });
581
582 conn.remove_rollback_hook();
583
584 let _ = conn.immediate_transaction(|conn| {
585 crate::sql_query("INSERT INTO t_rem_rb (id) VALUES (1)")
586 .execute(conn)
587 .unwrap();
588 Err::<(), _>(crate::result::Error::RollbackTransaction)
589 });
590
591 assert_eq!(count.load(Ordering::Relaxed), 0);
592 }
593
594 // A recursive CTE heavy enough that the progress handler fires while it runs.
595 const HEAVY_QUERY: &str = "WITH RECURSIVE c(x) AS \
596 (SELECT 1 UNION ALL SELECT x + 1 FROM c WHERE x < 100000) SELECT count(*) FROM c";
597
598 // Registers a callback that is invoked by the native library, which is not
599 // supported when running under miri with a native libsqlite3 (`-Zmiri-native-lib`).
600 #[cfg(not(miri))] // ffi call
601 #[diesel_test_helper::test]
602 fn on_progress_interrupts_query() {
603 let conn = &mut connection();
604
605 conn.on_progress(NonZeroU32::new(1).unwrap(), || ProgressDecision::Interrupt);
606
607 let result = crate::sql_query(HEAVY_QUERY).execute(conn);
608 assert!(
609 result.is_err(),
610 "the query should be interrupted by the progress handler"
611 );
612 }
613
614 #[diesel_test_helper::test]
615 fn remove_progress_handler_stops_interruption() {
616 let conn = &mut connection();
617
618 conn.on_progress(NonZeroU32::new(1).unwrap(), || ProgressDecision::Interrupt);
619 conn.remove_progress_handler();
620
621 // With the handler removed the same query runs to completion.
622 let result = crate::sql_query(HEAVY_QUERY).execute(conn);
623 assert!(
624 result.is_ok(),
625 "the query should complete after the handler is removed"
626 );
627 }
628
629 // WAL hook tests
630 //
631 // Gated out on WASM because these tests need a file-backed database
632 // (WAL mode does not work with `:memory:`), and `tempfile::tempdir()`
633 // panics on WASM due to the lack of a filesystem.
634 //
635 // The WAL API itself (`on_wal`, `remove_wal_hook`) is available on all
636 // platforms, including WASM. The file-backed databases used here cannot
637 // be created on WASM, so only the tests are gated out.
638
639 #[cfg(not(any(all(target_family = "wasm", target_os = "unknown"), miri)))]
640 /// Helper: create a file-backed connection in WAL mode (WAL requires a real
641 /// file, and every WAL test below wants the connection already in WAL mode).
642 /// Its only consumers are also compiled out under miri (see the
643 /// individual tests below).
644 fn wal_connection() -> (SqliteConnection, tempfile::TempDir) {
645 let dir = tempfile::tempdir().unwrap();
646 let path = dir.path().join("test.db");
647 let mut conn = SqliteConnection::establish(path.to_str().unwrap()).unwrap();
648 crate::sql_query("PRAGMA journal_mode=WAL")
649 .execute(&mut conn)
650 .unwrap();
651 (conn, dir)
652 }
653
654 #[cfg(not(all(target_family = "wasm", target_os = "unknown")))]
655 // Registers a callback that is invoked by the native library, which is not
656 // supported when running under miri with a native libsqlite3 (`-Zmiri-native-lib`).
657 #[cfg(not(miri))] // ffi call
658 #[diesel_test_helper::test]
659 fn on_wal_fires_in_wal_mode() {
660 let (conn, _dir) = &mut wal_connection();
661
662 crate::sql_query("CREATE TABLE t_wal (id INTEGER PRIMARY KEY)")
663 .execute(conn)
664 .unwrap();
665
666 let events: Arc<std::sync::Mutex<Vec<(String, u32)>>> =
667 Arc::new(std::sync::Mutex::new(Vec::new()));
668 let events2 = events.clone();
669
670 conn.on_wal(move |_, db_name, n_pages| {
671 events2.lock().unwrap().push((db_name.to_owned(), n_pages));
672 });
673
674 crate::sql_query("INSERT INTO t_wal (id) VALUES (1)")
675 .execute(conn)
676 .unwrap();
677
678 let events = events.lock().unwrap();
679 assert!(
680 !events.is_empty(),
681 "WAL hook should have fired at least once"
682 );
683 assert!(
684 events.iter().all(|(db_name, _)| db_name == "main"),
685 "db_name should always be \"main\""
686 );
687 assert!(
688 events.iter().any(|(_, n_pages)| *n_pages > 0),
689 "n_pages should be positive"
690 );
691 }
692
693 #[cfg(not(all(target_family = "wasm", target_os = "unknown")))]
694 // Registers a callback that is invoked by the native library, which is not
695 // supported when running under miri with a native libsqlite3 (`-Zmiri-native-lib`).
696 #[cfg(not(miri))] // ffi call
697 #[diesel_test_helper::test]
698 fn replacing_wal_hook_drops_old() {
699 let (conn, _dir) = &mut wal_connection();
700
701 crate::sql_query("CREATE TABLE t_wal2 (id INTEGER PRIMARY KEY)")
702 .execute(conn)
703 .unwrap();
704
705 let old_count = Arc::new(AtomicU32::new(0));
706 let new_count = Arc::new(AtomicU32::new(0));
707
708 let c_old = old_count.clone();
709 conn.on_wal(move |_, _, _| {
710 c_old.fetch_add(1, Ordering::Relaxed);
711 });
712
713 // Replace with a new hook.
714 let c_new = new_count.clone();
715 conn.on_wal(move |_, _, _| {
716 c_new.fetch_add(1, Ordering::Relaxed);
717 });
718
719 crate::sql_query("INSERT INTO t_wal2 (id) VALUES (1)")
720 .execute(conn)
721 .unwrap();
722
723 // Old hook should NOT have fired after replacement.
724 let old_before = old_count.load(Ordering::Relaxed);
725 crate::sql_query("INSERT INTO t_wal2 (id) VALUES (2)")
726 .execute(conn)
727 .unwrap();
728 assert_eq!(
729 old_count.load(Ordering::Relaxed),
730 old_before,
731 "old WAL hook should not fire after replacement"
732 );
733 assert!(
734 new_count.load(Ordering::Relaxed) > 0,
735 "new WAL hook should fire"
736 );
737 }
738
739 #[cfg(not(all(target_family = "wasm", target_os = "unknown")))]
740 // Registers a callback that is invoked by the native library, which is not
741 // supported when running under miri with a native libsqlite3 (`-Zmiri-native-lib`).
742 #[cfg(not(miri))] // ffi call
743 #[diesel_test_helper::test]
744 fn remove_wal_hook_disables_callback() {
745 let (conn, _dir) = &mut wal_connection();
746
747 crate::sql_query("CREATE TABLE t_wal3 (id INTEGER PRIMARY KEY)")
748 .execute(conn)
749 .unwrap();
750
751 let count = Arc::new(AtomicU32::new(0));
752 let c2 = count.clone();
753
754 conn.on_wal(move |_, _, _| {
755 c2.fetch_add(1, Ordering::Relaxed);
756 });
757
758 conn.remove_wal_hook();
759
760 crate::sql_query("INSERT INTO t_wal3 (id) VALUES (1)")
761 .execute(conn)
762 .unwrap();
763
764 assert_eq!(count.load(Ordering::Relaxed), 0);
765 }
766
767 #[cfg(not(all(target_family = "wasm", target_os = "unknown")))]
768 // Registers a callback that is invoked by the native library, which is not
769 // supported when running under miri with a native libsqlite3 (`-Zmiri-native-lib`).
770 #[cfg(not(miri))] // ffi call
771 #[diesel_test_helper::test]
772 fn wal_hook_does_not_fire_in_default_journal_mode() {
773 // A plain file connection left in the default journal mode ("delete"),
774 // not WAL, so `wal_connection()` (which enables WAL) is not used here.
775 let dir = tempfile::tempdir().unwrap();
776 let path = dir.path().join("test.db");
777 let conn = &mut SqliteConnection::establish(path.to_str().unwrap()).unwrap();
778
779 crate::sql_query("CREATE TABLE t_wal4 (id INTEGER PRIMARY KEY)")
780 .execute(conn)
781 .unwrap();
782
783 let count = Arc::new(AtomicU32::new(0));
784 let c2 = count.clone();
785
786 conn.on_wal(move |_, _, _| {
787 c2.fetch_add(1, Ordering::Relaxed);
788 });
789
790 crate::sql_query("INSERT INTO t_wal4 (id) VALUES (1)")
791 .execute(conn)
792 .unwrap();
793
794 assert_eq!(
795 count.load(Ordering::Relaxed),
796 0,
797 "WAL hook should not fire when not in WAL mode"
798 );
799 }
800
801 #[cfg(not(all(target_family = "wasm", target_os = "unknown")))]
802 // Registers a callback that is invoked by the native library, which is not
803 // supported when running under miri with a native libsqlite3 (`-Zmiri-native-lib`).
804 #[cfg(not(miri))] // ffi call
805 #[diesel_test_helper::test]
806 fn on_wal_can_use_borrowed_connection() {
807 let (conn, _dir) = &mut wal_connection();
808
809 crate::sql_query("CREATE TABLE t_wal_use (id INTEGER PRIMARY KEY)")
810 .execute(conn)
811 .unwrap();
812
813 let counts: Arc<std::sync::Mutex<Vec<i64>>> = Arc::new(std::sync::Mutex::new(Vec::new()));
814 let counts2 = counts.clone();
815
816 conn.on_wal(move |conn, _db_name, _n_pages| {
817 // Read through the borrowed connection, a fresh implicit read
818 // transaction that finalizes on return.
819 let c = crate::sql_query("SELECT COUNT(*) AS c FROM t_wal_use")
820 .get_result::<CountResult>(conn)
821 .unwrap()
822 .c;
823 counts2.lock().unwrap().push(c);
824 });
825
826 crate::sql_query("INSERT INTO t_wal_use (id) VALUES (1)")
827 .execute(conn)
828 .unwrap();
829
830 let observed = counts.lock().unwrap();
831 assert!(!observed.is_empty(), "WAL hook should have fired");
832 assert!(
833 observed.contains(&1),
834 "callback should observe the committed row through the connection"
835 );
836 }
837
838 #[cfg(not(all(target_family = "wasm", target_os = "unknown")))]
839 // Registers a callback that is invoked by the native library, which is not
840 // supported when running under miri with a native libsqlite3 (`-Zmiri-native-lib`).
841 #[cfg(not(miri))] // ffi call
842 #[diesel_test_helper::test]
843 fn on_wal_callback_write_re_enters_hook() {
844 let (conn, _dir) = &mut wal_connection();
845
846 crate::sql_query("CREATE TABLE t_wal_re (id INTEGER PRIMARY KEY)")
847 .execute(conn)
848 .unwrap();
849 crate::sql_query("CREATE TABLE t_wal_log (id INTEGER PRIMARY KEY AUTOINCREMENT)")
850 .execute(conn)
851 .unwrap();
852
853 let calls = Arc::new(AtomicU32::new(0));
854 let calls2 = calls.clone();
855
856 conn.on_wal(move |conn, _db_name, _n_pages| {
857 let n = calls2.fetch_add(1, Ordering::Relaxed);
858 // Write on the first invocation only. The write commits in WAL mode
859 // and re-fires the hook re-entrantly. Gating on `n == 0` bounds the
860 // recursion to a single nested call instead of overflowing the stack.
861 if n == 0 {
862 crate::sql_query("INSERT INTO t_wal_log DEFAULT VALUES")
863 .execute(conn)
864 .unwrap();
865 }
866 });
867
868 crate::sql_query("INSERT INTO t_wal_re (id) VALUES (1)")
869 .execute(conn)
870 .unwrap();
871
872 // The outer commit fired the hook, and the write inside it re-fired the
873 // hook once more, so the `Fn` callback ran twice.
874 assert_eq!(
875 calls.load(Ordering::Relaxed),
876 2,
877 "a committing write inside the callback re-enters the hook"
878 );
879
880 // The write performed inside the callback took effect.
881 let logged: i64 = crate::sql_query("SELECT COUNT(*) AS c FROM t_wal_log")
882 .get_result::<CountResult>(conn)
883 .unwrap()
884 .c;
885 assert_eq!(
886 logged, 1,
887 "the write performed inside the callback should persist"
888 );
889 }
890
891 #[cfg(not(all(target_family = "wasm", target_os = "unknown")))]
892 // Registers a callback that is invoked by the native library, which is not
893 // supported when running under miri with a native libsqlite3 (`-Zmiri-native-lib`).
894 #[cfg(not(miri))] // ffi call
895 #[diesel_test_helper::test]
896 fn on_wal_fires_once_per_transaction_commit() {
897 let (conn, _dir) = &mut wal_connection();
898
899 crate::sql_query("CREATE TABLE t_wal_txn (id INTEGER PRIMARY KEY)")
900 .execute(conn)
901 .unwrap();
902
903 let count = Arc::new(AtomicU32::new(0));
904 let c2 = count.clone();
905
906 conn.on_wal(move |_, _, _| {
907 c2.fetch_add(1, Ordering::Relaxed);
908 });
909
910 // Several writes inside one explicit transaction commit together, so the
911 // WAL hook fires once for the single commit, not once per statement.
912 conn.immediate_transaction(|conn| {
913 crate::sql_query("INSERT INTO t_wal_txn (id) VALUES (1)").execute(conn)?;
914 crate::sql_query("INSERT INTO t_wal_txn (id) VALUES (2)").execute(conn)?;
915 crate::sql_query("INSERT INTO t_wal_txn (id) VALUES (3)").execute(conn)?;
916 Ok::<_, crate::result::Error>(())
917 })
918 .unwrap();
919
920 assert_eq!(
921 count.load(Ordering::Relaxed),
922 1,
923 "the WAL hook should fire once per transaction commit, not per statement"
924 );
925 }
926
927 // Busy handler test.
928 //
929 // Gated out on WASM: it needs two connections to a shared file-backed
930 // database (`:memory:` connections do not share a lock), and
931 // `tempfile::tempdir()` panics on WASM due to the lack of a filesystem.
932 #[cfg(not(all(target_family = "wasm", target_os = "unknown")))]
933 // Registers a callback that is invoked by the native library, which is not
934 // supported when running under miri with a native libsqlite3 (`-Zmiri-native-lib`).
935 #[cfg(not(miri))] // ffi call
936 #[diesel_test_helper::test]
937 fn on_busy_handler_is_invoked_on_lock_contention() {
938 let dir = tempfile::tempdir().unwrap();
939 let path = dir.path().join("busy.db");
940 let url = path.to_str().unwrap();
941
942 // One connection acquires and holds a write lock.
943 let mut holder = SqliteConnection::establish(url).unwrap();
944 crate::sql_query("CREATE TABLE t_busy (id INTEGER PRIMARY KEY)")
945 .execute(&mut holder)
946 .unwrap();
947 crate::sql_query("BEGIN IMMEDIATE")
948 .execute(&mut holder)
949 .unwrap();
950
951 // A second connection registers a busy handler that records the call
952 // and gives up.
953 let mut contender = SqliteConnection::establish(url).unwrap();
954 let calls = Arc::new(AtomicU32::new(0));
955 let calls2 = calls.clone();
956 contender.on_busy(move |_retry_count| {
957 calls2.fetch_add(1, Ordering::Relaxed);
958 BusyDecision::GiveUp
959 });
960
961 // The write contends with the held lock, so the busy handler fires.
962 // Because it gives up, the write fails instead of blocking.
963 let result = crate::sql_query("INSERT INTO t_busy (id) VALUES (1)").execute(&mut contender);
964
965 assert!(
966 result.is_err(),
967 "the contended write should fail once the busy handler gives up"
968 );
969 assert!(
970 calls.load(Ordering::Relaxed) >= 1,
971 "the busy handler should have been invoked at least once"
972 );
973 }
974}