Skip to main content

distributed/outbox/
table.rs

1use std::time::{SystemTime, UNIX_EPOCH};
2
3use crate::outbox::{OutboxMessage, OutboxMessageStatus};
4use crate::table::{
5    ColumnType, ExpectedVersion, PrimaryKey, RowKey, RowValue, RowValues, RowWriteMode,
6    TableColumn, TableIndex, TableModel, TableMutation, TableRowMutation, TableSchema,
7    TableStoreError, TableWritePlan,
8};
9
10pub const OUTBOX_MESSAGES_TABLE: &str = "outbox_messages";
11
12/// Schema for the durable outbox delivery table.
13pub fn outbox_message_schema() -> TableSchema {
14    TableSchema {
15        model_name: "OutboxMessage".into(),
16        table_name: OUTBOX_MESSAGES_TABLE.into(),
17        columns: vec![
18            table_column("message_id", ColumnType::Text, false),
19            table_column("event_type", ColumnType::Text, false),
20            table_column("payload", ColumnType::Bytes, false),
21            table_column("payload_codec", ColumnType::Text, false),
22            table_column("payload_codec_version", ColumnType::UnsignedInteger, false),
23            table_column("destination", ColumnType::Text, true),
24            table_column("metadata", ColumnType::Json, false),
25            table_column("status", ColumnType::Text, false),
26            table_column("created_at", ColumnType::Timestamp, false),
27            table_column("next_available_at", ColumnType::Timestamp, false),
28            timestamp_column_with_default("updated_at", "CURRENT_TIMESTAMP"),
29            table_column("claimed_by", ColumnType::Text, true),
30            table_column("claimed_until", ColumnType::Timestamp, true),
31            table_column("attempts", ColumnType::UnsignedInteger, false),
32            table_column("last_error", ColumnType::Text, true),
33            table_column("published_at", ColumnType::Timestamp, true),
34            table_column("failed_at", ColumnType::Timestamp, true),
35            table_column("source_aggregate_type", ColumnType::Text, true),
36            table_column("source_aggregate_id", ColumnType::Text, true),
37            table_column("source_sequence", ColumnType::UnsignedInteger, true),
38            table_column("correlation_id", ColumnType::Text, true),
39            table_column("causation_id", ColumnType::Text, true),
40        ],
41        primary_key: PrimaryKey::new(["message_id"]),
42        version_column: None,
43        foreign_keys: Vec::new(),
44        indexes: vec![
45            named_index(
46                "outbox_messages_claimable_idx",
47                ["status", "next_available_at", "claimed_until", "created_at"],
48            ),
49            named_index(
50                "outbox_messages_source_idx",
51                [
52                    "source_aggregate_type",
53                    "source_aggregate_id",
54                    "source_sequence",
55                ],
56            ),
57            named_index("outbox_messages_destination_idx", ["destination", "status"]),
58        ],
59        relationships: Vec::new(),
60    }
61}
62
63pub fn outbox_message_insert_plan(
64    message: &OutboxMessage,
65) -> Result<TableWritePlan, TableStoreError> {
66    let mutation = TableMutation::UpsertRow(TableRowMutation {
67        schema: outbox_message_schema(),
68        key: outbox_message_key(message.id()),
69        values: outbox_message_row_values(message)?,
70        expected_version: ExpectedVersion::NotExists,
71        mode: RowWriteMode::Insert,
72    });
73    Ok(TableWritePlan::new(vec![mutation]))
74}
75
76pub(crate) fn validate_outbox_message_table_write(
77    message: &OutboxMessage,
78) -> Result<(), TableStoreError> {
79    let plan = outbox_message_insert_plan(message)?;
80    plan.validate()
81}
82
83pub fn outbox_message_key(message_id: &str) -> RowKey {
84    RowKey::new([("message_id", RowValue::String(message_id.to_string()))])
85}
86
87pub fn outbox_message_row_values(message: &OutboxMessage) -> Result<RowValues, TableStoreError> {
88    let mut row = RowValues::new();
89    row.insert("message_id", RowValue::String(message.id().to_string()));
90    row.insert("event_type", RowValue::String(message.event_type.clone()));
91    row.insert("payload", RowValue::Bytes(message.payload.clone()));
92    row.insert(
93        "payload_codec",
94        RowValue::String(message.payload_codec.clone()),
95    );
96    row.insert(
97        "payload_codec_version",
98        RowValue::U64(message.payload_codec_version as u64),
99    );
100    row.insert("destination", optional_string(&message.destination));
101    row.insert_serde("metadata", &message.metadata)?;
102    row.insert("status", RowValue::String(message.status.as_str().into()));
103    let created_at = system_time_epoch_secs(message.created_at)?;
104    row.insert("created_at", RowValue::U64(created_at));
105    row.insert("next_available_at", RowValue::U64(created_at));
106    row.insert("claimed_by", optional_string(&message.worker_id));
107    row.insert(
108        "claimed_until",
109        optional_time_epoch_secs(message.leased_until)?,
110    );
111    row.insert("attempts", RowValue::U64(message.attempts as u64));
112    row.insert("last_error", optional_string(&message.last_error));
113    row.insert("published_at", RowValue::Null);
114    row.insert("failed_at", status_failed_at(message)?);
115    row.insert(
116        "source_aggregate_type",
117        optional_string(&message.source_aggregate_type),
118    );
119    row.insert(
120        "source_aggregate_id",
121        optional_string(&message.source_aggregate_id),
122    );
123    row.insert("source_sequence", optional_u64(message.source_sequence));
124    row.insert(
125        "correlation_id",
126        optional_str(message.metadata.get("correlation_id").map(String::as_str)),
127    );
128    row.insert(
129        "causation_id",
130        optional_str(message.metadata.get("causation_id").map(String::as_str)),
131    );
132    Ok(row)
133}
134
135impl TableModel for OutboxMessage {
136    fn table_schema() -> TableSchema {
137        outbox_message_schema()
138    }
139
140    fn table_key(&self) -> Result<RowKey, TableStoreError> {
141        Ok(outbox_message_key(self.id()))
142    }
143
144    fn to_table_row(&self) -> Result<RowValues, TableStoreError> {
145        outbox_message_row_values(self)
146    }
147}
148
149fn table_column(name: &str, column_type: ColumnType, nullable: bool) -> TableColumn {
150    let mut column = TableColumn::new(name, name, column_type);
151    column.nullable = nullable;
152    column
153}
154
155fn timestamp_column_with_default(name: &str, default: &str) -> TableColumn {
156    let mut column = table_column(name, ColumnType::Timestamp, false);
157    column.has_default = true;
158    column.default = Some(default.into());
159    column
160}
161
162fn named_index(name: &str, columns: impl IntoIterator<Item = &'static str>) -> TableIndex {
163    let mut index = TableIndex::new(columns);
164    index.name = Some(name.into());
165    index
166}
167
168fn optional_string(value: &Option<String>) -> RowValue {
169    optional_str(value.as_deref())
170}
171
172fn optional_str(value: Option<&str>) -> RowValue {
173    value
174        .map(|value| RowValue::String(value.to_string()))
175        .unwrap_or(RowValue::Null)
176}
177
178fn optional_u64(value: Option<u64>) -> RowValue {
179    value.map(RowValue::U64).unwrap_or(RowValue::Null)
180}
181
182fn optional_time_epoch_secs(value: Option<SystemTime>) -> Result<RowValue, TableStoreError> {
183    value
184        .map(system_time_epoch_secs)
185        .transpose()
186        .map(optional_u64)
187}
188
189fn status_failed_at(message: &OutboxMessage) -> Result<RowValue, TableStoreError> {
190    if message.status == OutboxMessageStatus::Failed {
191        Ok(RowValue::U64(system_time_epoch_secs(SystemTime::now())?))
192    } else {
193        Ok(RowValue::Null)
194    }
195}
196
197fn system_time_epoch_secs(value: SystemTime) -> Result<u64, TableStoreError> {
198    value
199        .duration_since(UNIX_EPOCH)
200        .map(|duration| duration.as_secs())
201        .map_err(|err| TableStoreError::Metadata(err.to_string()))
202}
203
204#[cfg(test)]
205mod tests {
206    use super::*;
207
208    #[test]
209    fn outbox_insert_plan_uses_table_row_mutation() {
210        let message = OutboxMessage::create("msg-1", "Event", b"{}".to_vec()).unwrap();
211
212        let plan = outbox_message_insert_plan(&message).unwrap();
213
214        assert_eq!(plan.mutations.len(), 1);
215        let TableMutation::UpsertRow(mutation) = &plan.mutations[0] else {
216            panic!("outbox insert should lower to a table row mutation");
217        };
218        assert_eq!(mutation.schema.table_name, OUTBOX_MESSAGES_TABLE);
219        assert_eq!(mutation.key, outbox_message_key("msg-1"));
220        assert_eq!(mutation.expected_version, ExpectedVersion::NotExists);
221        assert_eq!(mutation.mode, RowWriteMode::Insert);
222        assert_eq!(
223            mutation.values.get("event_type"),
224            Some(&RowValue::String("Event".into()))
225        );
226    }
227}