1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
use super::{ConnectionPool, StorageError, WriterAcquisitionCounters, WriterAcquisitionSnapshot};
use std::sync::atomic::Ordering;
impl WriterAcquisitionCounters {
pub(crate) fn record_writer_task_acquisition(&self) {
self.writer_task_acquisitions
.fetch_add(1, Ordering::Relaxed);
}
/// Records one writer-task `BEGIN IMMEDIATE` refused busy or locked.
/// Called for every such refusal, whether or not a bounded retry goes
/// on to absorb it — this is the caller-facing contention count, and it
/// alone must equal the number of busy/locked refusals SQLite actually
/// returned, independent of retry policy.
pub(crate) fn record_writer_task_begin_busy(&self) {
self.writer_task_begin_busy.fetch_add(1, Ordering::Relaxed);
}
/// Records one busy or locked `BEGIN IMMEDIATE` refusal hidden from the
/// caller by a subsequent bounded retry. This counter moves before the
/// next BEGIN attempt, in addition to (never instead of) the
/// `writer_task_begin_busy` call for the same refusal; it never implies
/// that the request closure ran.
pub(crate) fn record_writer_task_begin_busy_absorbed(&self) {
self.writer_task_begin_busy_absorbed
.fetch_add(1, Ordering::Relaxed);
}
/// Records one writer-task `BEGIN IMMEDIATE` that failed for any other
/// reason. Without this the non-busy arm reproduces, one level down, the
/// same silent-failure gap the busy counter closes.
pub(crate) fn record_writer_task_begin_error(&self) {
self.writer_task_begin_errors
.fetch_add(1, Ordering::Relaxed);
}
/// Records one dequeued writer-task request that reached the writer seam
/// and terminated in error. Called exactly once per such request,
/// regardless of which terminal state it produced.
pub(crate) fn record_writer_task_request_failure(&self) {
self.writer_task_request_failures
.fetch_add(1, Ordering::Relaxed);
}
/// Records the subset of [`Self::record_writer_task_request_failure`]
/// whose terminal state was `SideEffectsUnknown`. Callers pair this call
/// with a `record_writer_task_request_failure()` call for the same
/// request rather than in place of it.
pub(crate) fn record_writer_task_side_effects_unknown(&self) {
self.writer_task_side_effects_unknown
.fetch_add(1, Ordering::Relaxed);
}
/// `lease_timeouts` is owned by the pool's write admission, which takes
/// the volume lease for every writer class, so the pool supplies it.
pub(super) fn snapshot(&self, lease_timeouts: u64) -> WriterAcquisitionSnapshot {
let pooled_acquisitions = self.pooled_acquisitions.load(Ordering::Relaxed);
let standalone_acquisitions = self.standalone_acquisitions.load(Ordering::Relaxed);
let writer_task_acquisitions = self.writer_task_acquisitions.load(Ordering::Relaxed);
WriterAcquisitionSnapshot {
acquisitions: pooled_acquisitions
.saturating_add(standalone_acquisitions)
.saturating_add(writer_task_acquisitions),
pooled_acquisitions,
standalone_acquisitions,
writer_task_acquisitions,
timeouts: self.pooled_timeouts.load(Ordering::Relaxed),
lease_timeouts,
direct_busy_refusals: self.direct_busy_refusals.load(Ordering::Relaxed),
writer_task_begin_busy: self.writer_task_begin_busy.load(Ordering::Relaxed),
writer_task_begin_busy_absorbed: self
.writer_task_begin_busy_absorbed
.load(Ordering::Relaxed),
writer_task_begin_errors: self.writer_task_begin_errors.load(Ordering::Relaxed),
writer_task_request_failures: self.writer_task_request_failures.load(Ordering::Relaxed),
writer_task_side_effects_unknown: self
.writer_task_side_effects_unknown
.load(Ordering::Relaxed),
writer_guard_drop_rollbacks: self.writer_guard_drop_rollbacks.load(Ordering::Relaxed),
}
}
}
impl ConnectionPool {
/// Observe one final direct execution error without changing its classification.
/// Checkout, connection-open, reader and writer-task errors do not call this boundary.
pub(crate) fn record_direct_writer_error(&self, error: &StorageError) {
if direct_writer_sqlite_code(error) == Some(rusqlite::ErrorCode::DatabaseBusy) {
self.writer_acquisition_counters
.direct_busy_refusals
.fetch_add(1, Ordering::Relaxed);
}
}
/// Record a pooled writer guard dropped inside a transaction: count it in
/// [`WriterAcquisitionSnapshot::writer_guard_drop_rollbacks`] and write a
/// `writer_guard_drop` sink row naming how the drop settled it.
pub(crate) fn record_writer_guard_drop(
&self,
settlement: &Result<(), crate::error::SqliteError>,
) {
self.writer_acquisition_counters
.writer_guard_drop_rollbacks
.fetch_add(1, Ordering::Relaxed);
let outcome = match settlement {
Ok(()) => "guard dropped with open transaction; rolled back".to_string(),
Err(error) => {
format!("guard dropped with open transaction; writer retired: {error}")
}
};
let db = crate::timeout_sink::db_label(self);
tracing::warn!(db = %db, %outcome, "pooled writer guard settled on drop");
crate::timeout_sink::emit_writer_guard_drop(&db, &outcome);
}
}
// Public atomic callbacks can supply cyclic Error source chains.
// A finite walk bounds observation when each source() call returns.
const DIRECT_WRITER_SOURCE_LIMIT: usize = 32;
/// Inspect at most 32 nodes, counting the returned StorageError as node one,
/// and make at most 32 source() calls. Preserved storage wrappers are included;
/// message text, deeper causes and cause-free outcomes cannot supply a code.
fn direct_writer_sqlite_code(error: &StorageError) -> Option<rusqlite::ErrorCode> {
let mut cause: &(dyn std::error::Error + 'static) = error;
for _ in 0..DIRECT_WRITER_SOURCE_LIMIT {
if let Some(sqlite) = cause.downcast_ref::<rusqlite::Error>() {
return sqlite.sqlite_error_code();
}
if let Some(crate::error::SqliteError::Rusqlite(sqlite)) =
cause.downcast_ref::<crate::error::SqliteError>()
{
return sqlite.sqlite_error_code();
}
cause = cause.source()?;
}
None
}