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
12pub 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}