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
//! `reliar-outbox` is the storage-agnostic transactional outbox: the [`OutboxStore`]/
//! [`OutboxDeadLetters`] capability traits (plus `reliar_core::Publisher`, re-exported here for
//! convenience), the request and result types that cross their boundary, a pure [`RetryPolicy`],
//! the feature's [`OutboxSettings`], and the [`OutboxMetrics`] hook.
//!
//! **The object names the guarantee — there is no facade type joining the two** (ADR 0036
//! amendment B). A caller enqueues durably by calling [`OutboxEnqueue::enqueue`] on the
//! provider store, in its own transaction: the message becomes visible when that transaction
//! commits and is published later by an [`OutboxDispatcher`], at-least-once, with the duplicate
//! windows below. A caller that instead wants to send now, with no Reliar durability at all, calls
//! the transport's own [`Publisher::publish`] directly — `reliar-outbox` never wraps it. Nothing
//! decides between the two at runtime, and no setting can.
//!
//! This slice ships the traits and types a provider builds against, a host configures, the
//! [`OutboxDispatcher`] worker loop, and the `test-support` fakes a test drives without a
//! database.
//!
//! Enable the `test-support` feature for `InMemoryOutboxStore`, `RecordingPublisher`,
//! `ScriptedPublisher` and `RecordingMetrics` — one shared set of fakes reused by provider
//! crates, examples and `tests/system`.
// Plain code spans above, not intra-doc links: those types only exist when `test-support` is
// enabled, and `cargo doc` on default features must not break trying to resolve a link to an
// item that is not compiled in.
//!
// Gated: this example is only compilable with `test-support` (`InMemoryOutboxStore`,
// `InMemoryTransaction`, `RecordingPublisher`) — without the feature it becomes `ignore` so
// `cargo test -p reliar-outbox` (no `--all-features`) still compiles.
//! # use reliar_core::{ContentType, Envelope, Message, Publisher as _, Serializer};
//! # use reliar_outbox::{InMemoryOutboxStore, InMemoryTransaction, OutboxEnqueue, RecordingPublisher};
//! #
//! # #[derive(serde::Serialize, serde::Deserialize)]
//! # struct OrderCreated;
//! # impl Message for OrderCreated {
//! # const TYPE: &'static str = "orders.created";
//! # const VERSION: u16 = 1;
//! # }
//! #
//! # // A minimal `Serializer` fixture: `reliar-outbox` names no wire format of its own, and
//! # // holds none on this path — the caller serializes. Hidden: irrelevant to the pattern below.
//! # struct RawJson;
//! # impl Serializer for RawJson {
//! # type Error = serde_json::Error;
//! # fn content_type(&self) -> &ContentType { &ContentType::JSON }
//! # fn serialize<T: Message>(&self, body: &T) -> Result<bytes::Bytes, Self::Error> {
//! # serde_json::to_vec(body).map(bytes::Bytes::from)
//! # }
//! # fn deserialize<T: Message>(&self, bytes: &[u8]) -> Result<T, Self::Error> {
//! # serde_json::from_slice(bytes)
//! # }
//! # }
//! #
//! # #[tokio::main(flavor = "current_thread")]
//! # async fn main() -> Result<(), Box<dyn std::error::Error>> {
//! let store = InMemoryOutboxStore::default();
//! let publisher = RecordingPublisher::default();
//!
//! // The durable path: a provider store implements `OutboxEnqueue` directly — no facade type in
//! // between. A bare message becomes an envelope with default metadata and a freshly rooted
//! // conversation; enqueued in the caller's own transaction, published later by an
//! // `OutboxDispatcher`. Use `enqueue_envelope` with `Envelope::builder(..)` instead when an id
//! // must propagate from an inbound request.
//! let mut tx = InMemoryTransaction;
//! store.enqueue(&mut tx, OrderCreated).await?;
//!
//! // The bypass path: straight to the transport, no transaction needed, no Reliar guarantee.
//! // The caller serializes once, exactly as it would for a bare `NatsPublisher`.
//! let envelope = Envelope::builder(OrderCreated).build();
//! let bytes = RawJson.serialize(&envelope.body)?;
//! let mut serialized = envelope.map_body(|_| bytes);
//! serialized.metadata.delivery.content_type = RawJson.content_type().clone();
//! publisher.publish(&serialized).await?;
//! # Ok(())
//! # }
//! ```
//!
//! # Guarantees
//!
//! - **Durable at-least-once publication. Never exactly-once.** Duplicate delivery is expected
//! and must be handled by an idempotent consumer. Three distinct windows produce a duplicate,
//! and all three are unavoidable:
//! 1. **The crash window:** a publish reaches the broker, the worker crashes before
//! `complete` persists, the lease expires, and another worker republishes the same message.
//! 2. **The slow-batch window:** no crash at all — a worker claims a large batch under
//! a lease shorter than the batch takes to drain, the lease expires while the worker is
//! still healthily publishing, a second worker reclaims and republishes the tail, and the
//! first worker's later `complete`/`fail` is rejected by the `locked_by` guard.
//! 3. **The drain window:** on cancellation, `run()` drains in-flight publishes for at
//! most `DispatcherSettings::drain_timeout`; a publish still unresolved at the timeout is
//! released rather than awaited further, and its outcome — success or failure — is the same
//! duplicate risk as the other two windows, just triggered by shutdown instead of a lease.
//! - **No ordering by default.** [`Ordering::Unordered`] (the default) guarantees **nothing**
//! about order — not globally, not per `conversation_id`, not per aggregate, not
//! approximately. `SKIP LOCKED`, concurrent publishing, per-message backoff and multiple
//! workers each reorder freely (ADR 0013). [`Ordering::PerKey`] is a configuration error in
//! this release — see [`Ordering::validate`].
//! - **Pure retry.** [`RetryPolicy`] is I/O-free and clock-free: it returns a [`core::time::Duration`],
//! never a timestamp. The store applies it as `available_at = now() + delay` in SQL, so a
//! worker's clock skew can never hot-loop a row or park it in the future (ADR 0009).
//! - **The library never reads the environment implicitly.** Only [`OutboxSettings::from_env`]
//! touches `std::env`, and only when called (ADR 0019).
//! - **Calling the transport publisher directly bypasses the outbox.** A call to a transport's
//! [`Publisher::publish`] carries **none** of the above: one attempt, no retry, no backoff, no
//! dead state, no duplicate window — only as much retry as the transport itself performs, and
//! no relationship to any transaction the caller has open. Use [`OutboxEnqueue::enqueue`]/
//! [`OutboxEnqueue::enqueue_envelope`] for the durable path instead (ADR 0036 amendment B).
//! See `docs/guides/outbox-enqueue-and-publish.md` for the full comparison.
pub use ;
pub use OutboxEnqueue;
pub use ConfigError;
pub use ;
pub use Ordering;
pub use ;
/// Re-exported from `reliar-core` (ADR 0032): a store author's or a publisher's `Classify`
/// bound, a publish/store failure's `FailureKind`, the `Publisher` capability trait, and the
/// shared `SettingsError` all live in core now. New code should name `reliar_core::` directly;
/// this re-export keeps existing `use reliar_outbox::{…}` imports one line.
pub use ;
pub use ;
pub use ;
pub use ;
pub use ;
pub use WorkerId;
// The README quickstart drives the in-memory fake, so it compiles only with `test-support`.