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
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
//! Typed errors at the SQLite adapter boundary.
use dovecote::{DeliveryState, RowId};
use thiserror::Error;
/// Errors that can be retried by the adapter's bounded busy policy.
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
#[non_exhaustive]
pub enum TransientKind {
/// SQLite could not acquire its single-writer lock before the configured
/// busy timeout. The complete operation has been rolled back.
BusyExhausted,
}
impl std::fmt::Display for TransientKind {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter.write_str("SQLite busy timeout exhausted")
}
}
pub(crate) fn is_busy(source: &sqlx::Error) -> bool {
source
.as_database_error()
.and_then(|error| error.code())
.and_then(|code| code.parse::<i32>().ok())
.is_some_and(|code| code == 5 || code == 6 || code & 0xff == 5 || code & 0xff == 6)
}
/// Errors returned while enqueueing an event.
#[derive(Debug, Error)]
#[non_exhaustive]
pub enum EnqueueError {
/// The transaction was not opened with SQLite's writer lock.
#[error("enqueue requires a SQLite write transaction (BEGIN IMMEDIATE or a prior write)")]
WriteTransactionRequired,
/// The adapter's bounded busy policy is not representable by SQLite.
#[error("invalid SQLite busy configuration: {detail}")]
Configuration {
/// Diagnostic describing the invalid configuration.
detail: String,
},
/// Existing identity has different immutable event content.
#[error("idempotency conflict for existing row {existing_row_id:?}")]
IdempotencyConflict {
/// Existing Dovecote row that conflicted with the event.
existing_row_id: RowId,
},
/// Installed durable schema is incompatible with this adapter.
#[error("migration mismatch: {detail}")]
MigrationMismatch {
/// Diagnostic describing the incompatible schema.
detail: String,
},
/// Stored or returned data could not be represented safely.
#[error("serialization: {detail}")]
Serialization {
/// Diagnostic describing the invalid stored data.
detail: String,
},
/// A caller transaction remained blocked after the configured busy wait.
#[error("{operation}: busy lock exhausted by the caller transaction: {source}")]
BusyExhausted {
/// Operation being performed when the lock wait was exhausted.
operation: &'static str,
#[source]
/// Original underlying SQLite error.
source: sqlx::Error,
},
/// A non-busy SQLx operation failed.
#[error("{operation}: {source}")]
Sql {
/// Operation being performed when SQL failed.
operation: &'static str,
#[source]
/// Original underlying SQLite error.
source: sqlx::Error,
},
}
impl EnqueueError {
pub(crate) fn sql(operation: &'static str, source: sqlx::Error) -> Self {
if is_busy(&source) {
Self::BusyExhausted { operation, source }
} else {
Self::Sql { operation, source }
}
}
pub(crate) fn serialization(detail: impl Into<String>) -> Self {
Self::Serialization {
detail: detail.into(),
}
}
}
/// Errors returned by the migration-only importer.
#[derive(Debug, Error)]
#[non_exhaustive]
pub enum ImportError {
/// The transaction was not opened with SQLite's writer lock.
#[error("import requires a SQLite write transaction (BEGIN IMMEDIATE or a prior write)")]
WriteTransactionRequired,
/// Existing identity has different immutable event content.
#[error("immutable event identity conflict for existing row {existing_row_id:?}")]
IdentityConflict {
/// Existing Dovecote row that conflicted with the imported event.
existing_row_id: RowId,
},
/// Existing delivery state differs from the requested imported state.
#[error("imported delivery state conflict for existing row {existing_row_id:?}")]
ImportConflict {
/// Existing Dovecote row whose delivery state conflicted.
existing_row_id: RowId,
},
/// The adapter configuration is not valid for SQLite.
#[error("invalid SQLite configuration: {detail}")]
Configuration {
/// Diagnostic describing the invalid configuration.
detail: String,
},
/// Imported state failed Dovecote validation.
#[error("invalid imported delivery state: {source}")]
InvalidState {
#[source]
/// Original validation error.
source: dovecote::ValidationError,
},
/// Installed durable schema is incompatible with this adapter.
#[error("migration mismatch: {detail}")]
MigrationMismatch {
/// Diagnostic describing the incompatible schema.
detail: String,
},
/// Stored or returned data could not be represented safely.
#[error("serialization: {detail}")]
Serialization {
/// Diagnostic describing the invalid stored data.
detail: String,
},
/// A caller transaction remained blocked after the configured busy wait.
#[error("{operation}: busy lock exhausted by the caller transaction: {source}")]
BusyExhausted {
/// Operation being performed when the lock wait was exhausted.
operation: &'static str,
#[source]
/// Original underlying SQLite error.
source: sqlx::Error,
},
/// A non-busy SQLx operation failed.
#[error("{operation}: {source}")]
Sql {
/// Operation being performed when SQL failed.
operation: &'static str,
#[source]
/// Original underlying SQLite error.
source: sqlx::Error,
},
}
/// Errors returned by the migration-only delivery finalizer.
#[derive(Debug, Error)]
#[non_exhaustive]
pub enum FinalizeError {
/// The transaction was not opened with SQLite's writer lock.
#[error("finalization requires a SQLite write transaction (BEGIN IMMEDIATE or a prior write)")]
WriteTransactionRequired,
/// No event row exists for the requested delivery.
#[error("event row not found")]
NotFound,
/// The row is not a canonical imported pending delivery.
#[error("delivery row {row_id:?} is not a canonical imported pending delivery")]
StateConflict {
/// Dovecote row that was not in the expected state.
row_id: RowId,
},
/// Authoritative delivery time failed Dovecote validation.
#[error("invalid authoritative delivery timestamp: {source}")]
InvalidTimestamp {
#[source]
/// Original validation error.
source: dovecote::ValidationError,
},
/// Installed durable schema is incompatible with this adapter.
#[error("migration mismatch: {detail}")]
MigrationMismatch {
/// Diagnostic describing the incompatible schema.
detail: String,
},
/// Stored or returned data could not be represented safely.
#[error("serialization: {detail}")]
Serialization {
/// Diagnostic describing the invalid stored data.
detail: String,
},
/// A caller transaction remained blocked after the configured busy wait.
#[error("{operation}: busy lock exhausted by the caller transaction: {source}")]
BusyExhausted {
/// Operation being performed when the lock wait was exhausted.
operation: &'static str,
#[source]
/// Original underlying SQLite error.
source: sqlx::Error,
},
/// A non-busy SQLx operation failed.
#[error("{operation}: {source}")]
Sql {
/// Operation being performed when SQL failed.
operation: &'static str,
#[source]
/// Original underlying SQLite error.
source: sqlx::Error,
},
}
impl FinalizeError {
pub(crate) fn sql(operation: &'static str, source: sqlx::Error) -> Self {
if is_busy(&source) {
Self::BusyExhausted { operation, source }
} else {
Self::Sql { operation, source }
}
}
}
impl ImportError {
pub(crate) fn sql(operation: &'static str, source: sqlx::Error) -> Self {
if is_busy(&source) {
Self::BusyExhausted { operation, source }
} else {
Self::Sql { operation, source }
}
}
pub(crate) fn serialization(detail: impl Into<String>) -> Self {
Self::Serialization {
detail: detail.into(),
}
}
}
/// Errors returned while selecting and claiming a batch of events.
#[derive(Debug, Error)]
#[non_exhaustive]
pub enum ClaimError {
#[cfg(test)]
/// Test-only failpoint used to verify claim rollback.
#[error("test claim failpoint triggered after delivery updates")]
InjectedFailure,
/// The delivery attempt counter cannot be incremented.
#[error("attempt counter overflow for row {row_id:?}")]
CounterOverflow {
/// Dovecote row whose attempt counter overflowed.
row_id: RowId,
},
/// The operating system could not provide a fresh claim token.
#[error("operating-system entropy unavailable: {source}")]
EntropyUnavailable {
#[source]
/// Original entropy error.
source: getrandom::Error,
},
/// Stored or returned data could not be represented safely.
#[error("serialization: {detail}")]
Serialization {
/// Diagnostic describing the invalid stored data.
detail: String,
},
/// Installed durable schema is incompatible with this adapter.
#[error("migration mismatch: {detail}")]
MigrationMismatch {
/// Diagnostic describing the incompatible schema.
detail: String,
},
/// The adapter's bounded busy policy is not representable by SQLite.
#[error("invalid SQLite busy configuration: {detail}")]
Configuration {
/// Diagnostic describing the invalid configuration.
detail: String,
},
/// A claim could not finish after the configured busy retries.
#[error("{operation}: busy lock exhausted after bounded retries: {source}")]
BusyExhausted {
/// Operation being performed when retries were exhausted.
operation: &'static str,
#[source]
/// Original underlying SQLite error.
source: sqlx::Error,
},
/// A non-busy SQLx operation failed.
#[error("{operation}: {source}")]
Sql {
/// Operation being performed when SQL failed.
operation: &'static str,
#[source]
/// Original underlying SQLite error.
source: sqlx::Error,
},
}
impl ClaimError {
pub(crate) fn sql(operation: &'static str, source: sqlx::Error) -> Self {
Self::Sql { operation, source }
}
pub(crate) fn serialization(detail: impl Into<String>) -> Self {
Self::Serialization {
detail: detail.into(),
}
}
pub(crate) fn busy_source(&self) -> Option<&sqlx::Error> {
match self {
Self::Sql { source, .. } | Self::BusyExhausted { source, .. } if is_busy(source) => {
Some(source)
}
_ => None,
}
}
pub(crate) fn into_busy_exhausted(self) -> Self {
match self {
Self::Sql { operation, source } => Self::BusyExhausted { operation, source },
other => other,
}
}
}
/// Errors returned by claim-token-fenced delivery mutations.
#[derive(Debug, Error)]
#[non_exhaustive]
pub enum MutationError {
/// No event row exists for the requested delivery.
#[error("event row not found")]
NotFound,
/// The requested mutation is invalid for the current delivery state.
#[error("illegal delivery transition from {state:?}")]
IllegalTransition {
/// Current durable delivery state.
state: DeliveryState,
},
/// The supplied claim token no longer owns the delivery.
#[error("claim was lost")]
LostClaim,
/// Installed durable schema is incompatible with this adapter.
#[error("migration mismatch: {detail}")]
MigrationMismatch {
/// Diagnostic describing the incompatible schema.
detail: String,
},
/// The adapter's bounded busy policy is not representable by SQLite.
#[error("invalid SQLite busy configuration: {detail}")]
Configuration {
/// Diagnostic describing the invalid configuration.
detail: String,
},
/// Stored or returned data could not be represented safely.
#[error("serialization: {detail}")]
Serialization {
/// Diagnostic describing the invalid stored data.
detail: String,
},
/// A mutation could not finish after the configured busy retries.
#[error("{operation}: busy lock exhausted after bounded retries: {source}")]
BusyExhausted {
/// Operation being performed when retries were exhausted.
operation: &'static str,
#[source]
/// Original underlying SQLite error.
source: sqlx::Error,
},
/// A non-busy SQLx operation failed.
#[error("{operation}: {source}")]
Sql {
/// Operation being performed when SQL failed.
operation: &'static str,
#[source]
/// Original underlying SQLite error.
source: sqlx::Error,
},
}
impl MutationError {
pub(crate) fn sql(operation: &'static str, source: sqlx::Error) -> Self {
Self::Sql { operation, source }
}
pub(crate) fn serialization(detail: impl Into<String>) -> Self {
Self::Serialization {
detail: detail.into(),
}
}
pub(crate) fn into_busy_exhausted(self) -> Self {
match self {
Self::Sql { operation, source } => Self::BusyExhausted { operation, source },
other => other,
}
}
pub(crate) fn is_busy(&self) -> bool {
matches!(self, Self::Sql { source, .. } if is_busy(source))
}
}
/// Errors returned by live and snapshot paging.
#[derive(Debug, Error)]
#[non_exhaustive]
pub enum PageError {
/// The snapshot pager has already been finished or rolled back.
#[error("snapshot pager is closed")]
Closed,
/// Stored or returned data could not be represented safely.
#[error("serialization: {detail}")]
Serialization {
/// Diagnostic describing the invalid stored data.
detail: String,
},
/// A page operation could not finish after the configured busy retries.
#[error("{operation}: busy lock exhausted after bounded retries: {source}")]
BusyExhausted {
/// Operation being performed when retries were exhausted.
operation: &'static str,
#[source]
/// Original underlying SQLite error.
source: sqlx::Error,
},
/// A non-busy SQLx operation failed.
#[error("{operation}: {source}")]
Sql {
/// Operation being performed when SQL failed.
operation: &'static str,
#[source]
/// Original underlying SQLite error.
source: sqlx::Error,
},
}
impl PageError {
pub(crate) fn sql(operation: &'static str, source: sqlx::Error) -> Self {
if is_busy(&source) {
Self::BusyExhausted { operation, source }
} else {
Self::Sql { operation, source }
}
}
pub(crate) fn serialization(detail: impl Into<String>) -> Self {
Self::Serialization {
detail: detail.into(),
}
}
}
/// Errors returned while checking the installed schema.
#[derive(Debug, Error)]
#[non_exhaustive]
pub enum SchemaError {
/// Installed durable schema is incompatible with this adapter.
#[error("migration mismatch: {detail}")]
MigrationMismatch {
/// Diagnostic describing the incompatible schema.
detail: String,
},
/// A schema query could not finish after the configured busy wait.
#[error("{operation}: busy lock exhausted after bounded retries: {source}")]
BusyExhausted {
/// Operation being performed when the lock wait was exhausted.
operation: &'static str,
#[source]
/// Original underlying SQLite error.
source: sqlx::Error,
},
/// A non-busy SQLx operation failed.
#[error("{operation}: {source}")]
Sql {
/// Operation being performed when SQL failed.
operation: &'static str,
#[source]
/// Original underlying SQLite error.
source: sqlx::Error,
},
}
impl SchemaError {
pub(crate) fn sql(operation: &'static str, source: sqlx::Error) -> Self {
if is_busy(&source) {
Self::BusyExhausted { operation, source }
} else {
Self::Sql { operation, source }
}
}
}