| 217 | } |
| 218 | |
| 219 | fn write_to_redb(&self, entry: &TriggerDlqEntry) -> crate::Result<()> { |
| 220 | let bytes = zerompk::to_msgpack_vec(entry).map_err(|e| crate::Error::Serialization { |
| 221 | format: "msgpack".into(), |
| 222 | detail: format!("trigger DLQ entry: {e}"), |
| 223 | })?; |
| 224 | let txn = self.db.begin_write().map_err(|e| crate::Error::Storage { |
| 225 | engine: "event_plane".into(), |
| 226 | detail: format!("begin_write: {e}"), |
| 227 | })?; |
| 228 | { |
| 229 | let mut table = txn |
| 230 | .open_table(TRIGGER_DLQ) |
| 231 | .map_err(|e| crate::Error::Storage { |
| 232 | engine: "event_plane".into(), |
| 233 | detail: format!("open_table: {e}"), |
| 234 | })?; |
| 235 | table |
| 236 | .insert(entry.entry_id, bytes.as_slice()) |
| 237 | .map_err(|e| crate::Error::Storage { |
| 238 | engine: "event_plane".into(), |
| 239 | detail: format!("insert: {e}"), |
| 240 | })?; |
| 241 | } |
| 242 | txn.commit().map_err(|e| crate::Error::Storage { |
| 243 | engine: "event_plane".into(), |
| 244 | detail: format!("commit: {e}"), |
| 245 | })?; |
| 246 | Ok(()) |
| 247 | } |
| 248 | |
| 249 | fn delete_from_redb(&self, entry_id: u64) { |
| 250 | if let Ok(txn) = self.db.begin_write() { |