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
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
//! Typed errors at the PostgreSQL adapter boundary.
use dovecote::{DeliveryState, RowId};
use thiserror::Error;
/// PostgreSQL SQLSTATE categories for failures callers may retry as a whole
/// operation. The original SQLx error remains available as the source.
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
#[non_exhaustive]
pub enum TransientKind {
/// Serialization failure (`40001`).
SerializationFailure,
/// Deadlock detected (`40P01`).
DeadlockDetected,
/// Statement/query cancellation or lock timeout (`57014`/`55P03`).
StatementOrLockTimeout,
}
impl std::fmt::Display for TransientKind {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let label = match self {
Self::SerializationFailure => "serialization failure",
Self::DeadlockDetected => "deadlock detected",
Self::StatementOrLockTimeout => "statement or lock timeout",
};
formatter.write_str(label)
}
}
impl TransientKind {
pub(crate) fn from_sqlx(source: &sqlx::Error) -> Option<Self> {
Self::from_sqlstate(source.as_database_error()?.code()?.as_ref())
}
pub(crate) fn from_sqlstate(code: &str) -> Option<Self> {
match code {
"40001" => Some(Self::SerializationFailure),
"40P01" => Some(Self::DeadlockDetected),
"57014" | "55P03" => Some(Self::StatementOrLockTimeout),
_ => None,
}
}
}
#[cfg(test)]
mod tests {
use super::TransientKind;
#[test]
fn postgres_transient_sqlstates_have_typed_categories() {
assert_eq!(
TransientKind::from_sqlstate("40001"),
Some(TransientKind::SerializationFailure)
);
assert_eq!(
TransientKind::from_sqlstate("40P01"),
Some(TransientKind::DeadlockDetected)
);
for code in ["57014", "55P03"] {
assert_eq!(
TransientKind::from_sqlstate(code),
Some(TransientKind::StatementOrLockTimeout)
);
}
assert_eq!(TransientKind::from_sqlstate("23505"), None);
}
}
/// Errors returned while enqueueing an event.
#[derive(Debug, Error)]
#[non_exhaustive]
pub enum EnqueueError {
/// The immutable event identity already exists with different content.
#[error("idempotency conflict for existing row {existing_row_id:?}")]
IdempotencyConflict {
/// The existing event row whose identity conflicted.
existing_row_id: RowId,
},
#[error("migration mismatch: {detail}")]
/// The installed schema does not satisfy this adapter's migration contract.
MigrationMismatch {
/// Details identifying the incompatible schema contract.
detail: String,
},
/// A database value could not be reconstructed as a valid domain value.
#[error("serialization: {detail}")]
Serialization {
/// Details describing the invalid stored value.
detail: String,
},
#[error("{operation}: {source}")]
/// A non-transient SQL operation failed.
Sql {
/// The adapter operation that failed.
operation: &'static str,
/// The underlying SQLx error.
#[source]
source: sqlx::Error,
},
#[error("{operation}: {kind}: {source}")]
/// A SQL operation failed with a retryable PostgreSQL condition.
Transient {
/// The adapter operation that failed.
operation: &'static str,
/// The retryable PostgreSQL failure category.
kind: TransientKind,
/// The underlying SQLx error.
#[source]
source: sqlx::Error,
},
}
impl EnqueueError {
pub(crate) fn sql(operation: &'static str, source: sqlx::Error) -> Self {
match TransientKind::from_sqlx(&source) {
Some(kind) => Self::Transient {
operation,
kind,
source,
},
None => 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 immutable event identity conflicts with an existing row.
#[error("immutable event identity conflict for existing row {existing_row_id:?}")]
IdentityConflict {
/// The existing event row whose identity conflicted.
existing_row_id: RowId,
},
#[error("imported delivery state conflict for existing row {existing_row_id:?}")]
/// The imported delivery state conflicts with an existing row.
ImportConflict {
/// The existing event row whose delivery state conflicted.
existing_row_id: RowId,
},
#[error("invalid imported delivery state: {source}")]
/// The supplied legacy delivery state is invalid.
InvalidState {
/// The validation failure from the domain state.
#[source]
source: dovecote::ValidationError,
},
#[error("migration mismatch: {detail}")]
/// The installed schema does not satisfy this adapter's migration contract.
MigrationMismatch {
/// Details identifying the incompatible schema contract.
detail: String,
},
#[error("serialization: {detail}")]
/// A database value could not be reconstructed as a valid domain value.
Serialization {
/// Details describing the invalid stored value.
detail: String,
},
#[error("{operation}: {source}")]
/// A non-transient SQL operation failed.
Sql {
/// The adapter operation that failed.
operation: &'static str,
/// The underlying SQLx error.
#[source]
source: sqlx::Error,
},
#[error("{operation}: {kind}: {source}")]
/// A SQL operation failed with a retryable PostgreSQL condition.
Transient {
/// The adapter operation that failed.
operation: &'static str,
/// The retryable PostgreSQL failure category.
kind: TransientKind,
/// The underlying SQLx error.
#[source]
source: sqlx::Error,
},
}
/// Errors returned by the migration-only delivery finalizer.
#[derive(Debug, Error)]
#[non_exhaustive]
pub enum FinalizeError {
/// The requested event row does not exist.
#[error("event row not found")]
NotFound,
#[error("delivery row {row_id:?} is not a canonical imported pending delivery")]
/// The delivery row is not in the canonical imported pending state.
StateConflict {
/// The event row whose delivery state conflicted.
row_id: RowId,
},
#[error("invalid authoritative delivery timestamp: {source}")]
/// The supplied authoritative timestamp is invalid.
InvalidTimestamp {
/// The timestamp validation failure from the domain type.
#[source]
source: dovecote::ValidationError,
},
#[error("migration mismatch: {detail}")]
/// The installed schema does not satisfy this adapter's migration contract.
MigrationMismatch {
/// Details identifying the incompatible schema contract.
detail: String,
},
#[error("serialization: {detail}")]
/// A database value could not be reconstructed as a valid domain value.
Serialization {
/// Details describing the invalid stored value.
detail: String,
},
#[error("{operation}: {source}")]
/// A non-transient SQL operation failed.
Sql {
/// The adapter operation that failed.
operation: &'static str,
/// The underlying SQLx error.
#[source]
source: sqlx::Error,
},
#[error("{operation}: {kind}: {source}")]
/// A SQL operation failed with a retryable PostgreSQL condition.
Transient {
/// The adapter operation that failed.
operation: &'static str,
/// The retryable PostgreSQL failure category.
kind: TransientKind,
/// The underlying SQLx error.
#[source]
source: sqlx::Error,
},
}
impl FinalizeError {
pub(crate) fn sql(operation: &'static str, source: sqlx::Error) -> Self {
match TransientKind::from_sqlx(&source) {
Some(kind) => Self::Transient {
operation,
kind,
source,
},
None => Self::Sql { operation, source },
}
}
}
impl ImportError {
pub(crate) fn sql(operation: &'static str, source: sqlx::Error) -> Self {
match TransientKind::from_sqlx(&source) {
Some(kind) => Self::Transient {
operation,
kind,
source,
},
None => 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 {
/// The delivery attempt counter cannot be incremented safely.
#[error("attempt counter overflow for row {row_id:?}")]
CounterOverflow {
/// The event row whose attempt count overflowed.
row_id: RowId,
},
#[error("operating-system entropy unavailable: {source}")]
/// The operating system could not provide claim-token entropy.
EntropyUnavailable {
/// The underlying entropy-provider error.
#[source]
source: getrandom::Error,
},
#[error("serialization: {detail}")]
/// A database value could not be reconstructed as a valid domain value.
Serialization {
/// Details describing the invalid stored value.
detail: String,
},
#[error("migration mismatch: {detail}")]
/// The installed schema does not satisfy this adapter's migration contract.
MigrationMismatch {
/// Details identifying the incompatible schema contract.
detail: String,
},
#[error("{operation}: {source}")]
/// A non-transient SQL operation failed.
Sql {
/// The adapter operation that failed.
operation: &'static str,
/// The underlying SQLx error.
#[source]
source: sqlx::Error,
},
#[error("{operation}: {kind}: {source}")]
/// A SQL operation failed with a retryable PostgreSQL condition.
Transient {
/// The adapter operation that failed.
operation: &'static str,
/// The retryable PostgreSQL failure category.
kind: TransientKind,
/// The underlying SQLx error.
#[source]
source: sqlx::Error,
},
}
impl ClaimError {
pub(crate) fn sql(operation: &'static str, source: sqlx::Error) -> Self {
match TransientKind::from_sqlx(&source) {
Some(kind) => Self::Transient {
operation,
kind,
source,
},
None => Self::Sql { operation, source },
}
}
pub(crate) fn serialization(detail: impl Into<String>) -> Self {
Self::Serialization {
detail: detail.into(),
}
}
}
/// Errors returned by fenced post-claim mutations.
#[derive(Debug, Error)]
#[non_exhaustive]
pub enum MutationError {
/// The requested event row does not exist.
#[error("event row not found")]
NotFound,
#[error("illegal delivery transition from {state:?}")]
/// The current delivery state cannot perform this mutation.
IllegalTransition {
/// The delivery state that rejected the mutation.
state: DeliveryState,
},
#[error("claim was lost")]
/// The claim token is stale or the lease has expired.
LostClaim,
#[error("migration mismatch: {detail}")]
/// The installed schema does not satisfy this adapter's migration contract.
MigrationMismatch {
/// Details identifying the incompatible schema contract.
detail: String,
},
#[error("serialization: {detail}")]
/// A database value could not be reconstructed as a valid domain value.
Serialization {
/// Details describing the invalid stored value.
detail: String,
},
#[error("{operation}: {source}")]
/// A non-transient SQL operation failed.
Sql {
/// The adapter operation that failed.
operation: &'static str,
/// The underlying SQLx error.
#[source]
source: sqlx::Error,
},
#[error("{operation}: {kind}: {source}")]
/// A SQL operation failed with a retryable PostgreSQL condition.
Transient {
/// The adapter operation that failed.
operation: &'static str,
/// The retryable PostgreSQL failure category.
kind: TransientKind,
/// The underlying SQLx error.
#[source]
source: sqlx::Error,
},
}
/// Errors returned by live and snapshot paging.
#[derive(Debug, Error)]
#[non_exhaustive]
pub enum PageError {
/// A database value could not be reconstructed as a valid domain value.
#[error("serialization: {detail}")]
Serialization {
/// Details describing the invalid stored value.
detail: String,
},
#[error("{operation}: {source}")]
/// A non-transient SQL operation failed.
Sql {
/// The adapter operation that failed.
operation: &'static str,
/// The underlying SQLx error.
#[source]
source: sqlx::Error,
},
#[error("{operation}: {kind}: {source}")]
/// A SQL operation failed with a retryable PostgreSQL condition.
Transient {
/// The adapter operation that failed.
operation: &'static str,
/// The retryable PostgreSQL failure category.
kind: TransientKind,
/// The underlying SQLx error.
#[source]
source: sqlx::Error,
},
}
impl PageError {
pub(crate) fn sql(operation: &'static str, source: sqlx::Error) -> Self {
match TransientKind::from_sqlx(&source) {
Some(kind) => Self::Transient {
operation,
kind,
source,
},
None => Self::Sql { operation, source },
}
}
pub(crate) fn serialization(detail: impl Into<String>) -> Self {
Self::Serialization {
detail: detail.into(),
}
}
}
impl MutationError {
pub(crate) fn sql(operation: &'static str, source: sqlx::Error) -> Self {
match TransientKind::from_sqlx(&source) {
Some(kind) => Self::Transient {
operation,
kind,
source,
},
None => 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 {
/// The installed schema does not satisfy this adapter's migration contract.
#[error("migration mismatch: {detail}")]
MigrationMismatch {
/// Details identifying the incompatible schema contract.
detail: String,
},
#[error("{operation}: {source}")]
/// A non-transient SQL operation failed.
Sql {
/// The adapter operation that failed.
operation: &'static str,
/// The underlying SQLx error.
#[source]
source: sqlx::Error,
},
#[error("{operation}: {kind}: {source}")]
/// A SQL operation failed with a retryable PostgreSQL condition.
Transient {
/// The adapter operation that failed.
operation: &'static str,
/// The retryable PostgreSQL failure category.
kind: TransientKind,
/// The underlying SQLx error.
#[source]
source: sqlx::Error,
},
}
impl SchemaError {
pub(crate) fn sql(operation: &'static str, source: sqlx::Error) -> Self {
match TransientKind::from_sqlx(&source) {
Some(kind) => Self::Transient {
operation,
kind,
source,
},
None => Self::Sql { operation, source },
}
}
}