Enqueue a failed trigger event.
(&mut self, params: DlqEnqueueParams)
| 140 | |
| 141 | /// Enqueue a failed trigger event. |
| 142 | pub fn enqueue(&mut self, params: DlqEnqueueParams) -> crate::Result<u64> { |
| 143 | let entry_id = self.next_entry_id; |
| 144 | self.next_entry_id += 1; |
| 145 | |
| 146 | let now = SystemTime::now() |
| 147 | .duration_since(UNIX_EPOCH) |
| 148 | .unwrap_or_default() |
| 149 | .as_millis() as u64; |
| 150 | |
| 151 | let entry = TriggerDlqEntry { |
| 152 | entry_id, |
| 153 | tenant_id: params.tenant_id, |
| 154 | source_collection: params.source_collection, |
| 155 | row_id: params.row_id, |
| 156 | operation: params.operation, |
| 157 | trigger_name: params.trigger_name, |
| 158 | error: params.error, |
| 159 | retry_count: params.retry_count, |
| 160 | source_lsn: params.source_lsn, |
| 161 | source_sequence: params.source_sequence, |
| 162 | created_at: now, |
| 163 | resolved: false, |
| 164 | }; |
| 165 | |
| 166 | // Evict oldest if at capacity. |
| 167 | while self.entries.len() >= self.max_entries { |
| 168 | if let Some(evicted) = self.entries.pop_front() { |
| 169 | self.delete_from_redb(evicted.entry_id); |
| 170 | warn!( |
| 171 | entry_id = evicted.entry_id, |
| 172 | trigger = %evicted.trigger_name, |
| 173 | "trigger DLQ evicted oldest entry (at capacity)" |
| 174 | ); |
| 175 | } |
| 176 | } |
| 177 | |
| 178 | // Persist to redb. |
| 179 | self.write_to_redb(&entry)?; |
| 180 | |
| 181 | debug!( |
| 182 | entry_id, |
| 183 | trigger = %entry.trigger_name, |
| 184 | "trigger event sent to DLQ" |
| 185 | ); |
| 186 | self.entries.push_back(entry); |
| 187 | Ok(entry_id) |
| 188 | } |
| 189 | |
| 190 | /// List all unresolved DLQ entries. |
| 191 | pub fn list_unresolved(&self) -> Vec<&TriggerDlqEntry> { |