1use super::{
2 stale_note_snapshot_error, EdgeListFilter, EdgeRelation, FilterOp, KhiveRuntime,
3 NamespaceToken, NoteFilter, OutboxSlugFilter, PageRequest, PropertyFilter, RuntimeError,
4 RuntimeResult, SqlValue, Uuid, Value,
5};
6
7impl KhiveRuntime {
8 pub async fn claim_outbound_message_external_id(
22 &self,
23 token: &NamespaceToken,
24 id: Uuid,
25 external_id: String,
26 ) -> RuntimeResult<khive_storage::note::Note> {
27 let store = self.notes(token)?;
28 let note = store
29 .get_note(id)
30 .await?
31 .ok_or_else(|| RuntimeError::NotFound(format!("note {id}")))?;
32 if note.kind != "message" {
33 return Err(RuntimeError::InvalidInput(format!(
34 "external_id can only be claimed on a `message` note; note {id} is a `{}`",
35 note.kind
36 )));
37 }
38 let props = note.properties.as_ref().and_then(|v| v.as_object());
39 let direction = props
40 .and_then(|p| p.get("direction"))
41 .and_then(|v| v.as_str());
42 if direction != Some("outbound") {
43 return Err(RuntimeError::InvalidInput(format!(
44 "external_id can only be claimed on an outbound message note; note {id} has \
45 direction {:?}",
46 direction
47 )));
48 }
49 if note.deleted_at.is_some()
50 || Self::outbound_delivery_is_terminal(props)
51 || props
52 .and_then(|p| p.get("delivered_at"))
53 .is_some_and(|value| !value.is_null())
54 {
55 return Err(RuntimeError::InvalidInput(format!(
56 "note {id} is not pending outbound delivery"
57 )));
58 }
59 let existing = props
60 .and_then(|p| p.get("external_id"))
61 .and_then(|v| v.as_str());
62 if existing.is_some_and(|v| !v.is_empty()) {
63 return Err(RuntimeError::InvalidInput(format!(
64 "note {id} already has an external_id claimed"
65 )));
66 }
67 let mut properties = props
68 .cloned()
69 .expect("outbound direction requires an object");
70 properties.insert("external_id".to_string(), Value::String(external_id));
71 self.replace_outbound_message_properties_as_owner(token, note, properties)
72 .await
73 }
74
75 pub async fn list_undelivered_outbound_messages(
100 &self,
101 token: &NamespaceToken,
102 to_prefix: Option<&str>,
103 limit: u32,
104 ) -> RuntimeResult<Vec<khive_storage::note::Note>> {
105 self.list_undelivered_outbound_messages_scoped(
106 token,
107 to_prefix,
108 limit,
109 OutboxSlugFilter::Any,
110 )
111 .await
112 }
113
114 pub async fn list_undelivered_outbound_messages_for_channel(
120 &self,
121 token: &NamespaceToken,
122 to_prefix: &str,
123 channel_slug: &str,
124 include_legacy: bool,
125 limit: u32,
126 ) -> RuntimeResult<Vec<khive_storage::note::Note>> {
127 let mut notes = self
128 .list_undelivered_outbound_messages_scoped(
129 token,
130 Some(to_prefix),
131 limit,
132 OutboxSlugFilter::Exact(channel_slug),
133 )
134 .await?;
135 if include_legacy {
136 notes.extend(
137 self.list_undelivered_outbound_messages_scoped(
138 token,
139 Some(to_prefix),
140 limit,
141 OutboxSlugFilter::Missing,
142 )
143 .await?,
144 );
145 notes.sort_by(|a, b| {
146 b.created_at
147 .cmp(&a.created_at)
148 .then_with(|| a.id.cmp(&b.id))
149 });
150 notes.truncate(limit as usize);
151 }
152 Ok(notes)
153 }
154
155 async fn list_undelivered_outbound_messages_scoped(
156 &self,
157 token: &NamespaceToken,
158 to_prefix: Option<&str>,
159 limit: u32,
160 slug_filter: OutboxSlugFilter<'_>,
161 ) -> RuntimeResult<Vec<khive_storage::note::Note>> {
162 const MAX_PAGE_TOTAL: u32 = 10_000;
163 if limit == 0 {
164 return Ok(Vec::new());
165 }
166 let now_micros = chrono::Utc::now().timestamp_micros();
167 let due_through = chrono::DateTime::<chrono::Utc>::from_timestamp_micros(now_micros)
171 .expect("current UTC time fits a chrono timestamp")
172 + chrono::Duration::nanoseconds(999);
173 let mut property_filters = vec![
174 PropertyFilter {
175 json_path: "$.direction".to_string(),
176 op: FilterOp::Eq,
177 value: SqlValue::Text("outbound".to_string()),
178 },
179 PropertyFilter {
180 json_path: "$.delivered_at".to_string(),
181 op: FilterOp::JsonTypeMissingOrNullIndexed,
182 value: SqlValue::Null,
183 },
184 PropertyFilter {
185 json_path: "$.delivery".to_string(),
186 op: FilterOp::NotInOrMissing(vec![
187 SqlValue::Text("delivered".to_string()),
188 SqlValue::Text("failed".to_string()),
189 ]),
190 value: SqlValue::Null,
191 },
192 PropertyFilter {
193 json_path: "$.delivery_hold".to_string(),
194 op: FilterOp::JsonTypeMissingOrNullIndexed,
195 value: SqlValue::Null,
196 },
197 PropertyFilter {
198 json_path: "$.next_attempt_at".to_string(),
199 op: FilterOp::Rfc3339LteOrInvalid,
200 value: SqlValue::Timestamp(due_through),
201 },
202 ];
203 if let Some(prefix) = to_prefix {
204 let op = if prefix
205 .strip_suffix(':')
206 .is_some_and(|head| !head.is_empty() && !head.contains(':'))
207 {
208 FilterOp::TextColonPrefixBucketIndexed
209 } else {
210 FilterOp::TextStartsWithIndexed
211 };
212 property_filters.push(PropertyFilter {
213 json_path: "$.to_actor".to_string(),
214 op,
215 value: SqlValue::Text(prefix.to_string()),
216 });
217 }
218 match slug_filter {
219 OutboxSlugFilter::Any => {}
220 OutboxSlugFilter::Exact(slug) => property_filters.push(PropertyFilter {
221 json_path: "$.channel_slug".to_string(),
222 op: FilterOp::Eq,
223 value: SqlValue::Text(slug.to_string()),
224 }),
225 OutboxSlugFilter::Missing => property_filters.push(PropertyFilter {
226 json_path: "$.channel_slug".to_string(),
227 op: FilterOp::JsonTypeMissing,
228 value: SqlValue::Null,
229 }),
230 }
231 let filter = NoteFilter {
232 kind: Some("message".to_string()),
233 property_filters,
234 ..Default::default()
235 };
236 let page = self
237 .notes(token)?
238 .query_notes_filtered_count_free(
239 token.namespace().as_str(),
240 &filter,
241 PageRequest {
242 limit: limit.min(MAX_PAGE_TOTAL),
243 offset: 0,
244 },
245 )
246 .await?;
247 Ok(page.items)
248 }
249
250 async fn outbound_message(
255 &self,
256 token: &NamespaceToken,
257 id: Uuid,
258 ) -> RuntimeResult<khive_storage::note::Note> {
259 let note = self
260 .notes(token)?
261 .get_note(id)
262 .await?
263 .ok_or_else(|| RuntimeError::NotFound(format!("note {id}")))?;
264 if note.kind != "message" || note.deleted_at.is_some() {
265 return Err(RuntimeError::InvalidInput(format!(
266 "note {id} is not a live message note (kind {})",
267 note.kind
268 )));
269 }
270 let outbound = note
271 .properties
272 .as_ref()
273 .and_then(|v| v.as_object())
274 .and_then(|p| p.get("direction"))
275 .and_then(|v| v.as_str())
276 == Some("outbound");
277 if !outbound {
278 return Err(RuntimeError::InvalidInput(format!(
279 "note {id} is not an outbound message"
280 )));
281 }
282 Ok(note)
283 }
284
285 async fn replace_outbound_message_properties(
286 &self,
287 token: &NamespaceToken,
288 mut snapshot: khive_storage::note::Note,
289 properties: serde_json::Map<String, Value>,
290 ) -> RuntimeResult<khive_storage::note::Note> {
291 let expected_updated_at = snapshot.updated_at;
292 let expected_deleted_at = snapshot.deleted_at;
293 let id = snapshot.id;
294 snapshot.properties = Some(Value::Object(properties));
295 crate::secret_gate::reject_reserved_secret_gate_property(snapshot.properties.as_ref())?;
296 snapshot.updated_at = chrono::Utc::now().timestamp_micros().max(
297 expected_updated_at.checked_add(1).ok_or_else(|| {
298 RuntimeError::Internal(format!(
299 "note {id} updated_at is already at i64::MAX and cannot advance"
300 ))
301 })?,
302 );
303
304 let store = self.raw_notes(token)?;
306 let persisted = store
307 .replace_note_if_unchanged(snapshot, expected_updated_at, expected_deleted_at)
308 .await?;
309 if !persisted {
310 return Err(stale_note_snapshot_error(id));
311 }
312 store
315 .get_note(id)
316 .await?
317 .ok_or_else(|| RuntimeError::NotFound(format!("note {id}")))
318 }
319
320 async fn replace_outbound_message_properties_as_owner(
321 &self,
322 token: &NamespaceToken,
323 mut snapshot: khive_storage::note::Note,
324 properties: serde_json::Map<String, Value>,
325 ) -> RuntimeResult<khive_storage::note::Note> {
326 let expected_updated_at = snapshot.updated_at;
327 let expected_deleted_at = snapshot.deleted_at;
328 let id = snapshot.id;
329 snapshot.properties = Some(Value::Object(properties));
330 crate::secret_gate::reject_reserved_secret_gate_property(snapshot.properties.as_ref())?;
331 snapshot.updated_at = chrono::Utc::now().timestamp_micros().max(
332 expected_updated_at.checked_add(1).ok_or_else(|| {
333 RuntimeError::Internal(format!("note {id} updated_at cannot advance"))
334 })?,
335 );
336 let store = self.raw_notes(token)?;
339 if !store
340 .replace_note_if_unchanged(snapshot, expected_updated_at, expected_deleted_at)
341 .await?
342 {
343 return Err(stale_note_snapshot_error(id));
344 }
345 store
347 .get_note(id)
348 .await?
349 .ok_or_else(|| RuntimeError::NotFound(format!("note {id}")))
350 }
351
352 fn outbound_delivery_is_terminal(props: Option<&serde_json::Map<String, Value>>) -> bool {
362 props
363 .and_then(|properties| properties.get("delivery"))
364 .and_then(Value::as_str)
365 .is_some_and(|state| state == "delivered" || state == "failed")
366 }
367
368 pub async fn mark_outbound_message_transient_failure(
373 &self,
374 token: &NamespaceToken,
375 id: Uuid,
376 attempted_at: chrono::DateTime<chrono::Utc>,
377 last_error: String,
378 base_delay: std::time::Duration,
379 max_delay: std::time::Duration,
380 ) -> RuntimeResult<khive_storage::note::Note> {
381 if base_delay.is_zero() || max_delay < base_delay {
382 return Err(RuntimeError::InvalidInput(
383 "outbound retry delays require a non-zero base no greater than the ceiling"
384 .to_string(),
385 ));
386 }
387
388 crate::secret_gate::check_json_at(
389 &serde_json::json!({
390 "last_error": &last_error,
391 }),
392 "message",
393 "last_error",
394 )?;
395
396 let snapshot = self.outbound_message(token, id).await?;
397 let props = snapshot.properties.as_ref().and_then(Value::as_object);
398 if Self::outbound_delivery_is_terminal(props) {
399 return Err(RuntimeError::InvalidInput(format!(
400 "outbound message {id} already has a terminal delivery outcome"
401 )));
402 }
403
404 let properties = Self::outbound_retry_properties(
405 props,
406 attempted_at,
407 last_error,
408 base_delay,
409 max_delay,
410 )?;
411 self.replace_outbound_message_properties(token, snapshot, properties)
412 .await
413 }
414
415 fn outbound_retry_properties(
416 props: Option<&serde_json::Map<String, Value>>,
417 attempted_at: chrono::DateTime<chrono::Utc>,
418 last_error: String,
419 base_delay: std::time::Duration,
420 max_delay: std::time::Duration,
421 ) -> RuntimeResult<serde_json::Map<String, Value>> {
422 let attempts = props
423 .and_then(|properties| properties.get("delivery_attempts"))
424 .and_then(Value::as_u64)
425 .unwrap_or(0)
426 .saturating_add(1);
427 let exponent = attempts.saturating_sub(1).min(127) as u32;
428 let delay_nanos = base_delay
429 .as_nanos()
430 .saturating_mul(1u128 << exponent)
431 .min(max_delay.as_nanos());
432 let delay = std::time::Duration::new(
433 (delay_nanos / 1_000_000_000) as u64,
434 (delay_nanos % 1_000_000_000) as u32,
435 );
436 let chrono_delay = chrono::TimeDelta::from_std(delay).map_err(|_| {
437 RuntimeError::InvalidInput("outbound retry ceiling exceeds RFC 3339 range".to_string())
438 })?;
439 let next_attempt_at = attempted_at
440 .checked_add_signed(chrono_delay)
441 .ok_or_else(|| {
442 RuntimeError::InvalidInput(
443 "outbound retry deadline exceeds RFC 3339 range".to_string(),
444 )
445 })?;
446
447 let mut properties = props.cloned().unwrap_or_default();
448 properties.insert("delivery_attempts".to_string(), Value::from(attempts));
449 properties.insert(
450 "next_attempt_at".to_string(),
451 Value::String(next_attempt_at.to_rfc3339()),
452 );
453 properties.insert("last_error".to_string(), Value::String(last_error));
454 Ok(properties)
455 }
456
457 pub async fn mark_outbound_message_claim_transient_failure(
460 &self,
461 token: &NamespaceToken,
462 id: Uuid,
463 attempted_at: chrono::DateTime<chrono::Utc>,
464 last_error: String,
465 base_delay: std::time::Duration,
466 max_delay: std::time::Duration,
467 ) -> RuntimeResult<khive_storage::note::Note> {
468 let snapshot = self.outbound_message(token, id).await?;
469 let props = snapshot.properties.as_ref().and_then(Value::as_object);
470 let has_claim = props
471 .and_then(|properties| properties.get("external_id"))
472 .and_then(Value::as_str)
473 .is_some_and(|value| !value.is_empty());
474 let has_delivery = props
475 .and_then(|properties| properties.get("delivered_at"))
476 .is_some_and(|value| !value.is_null());
477 if has_claim || has_delivery || Self::outbound_delivery_is_terminal(props) {
478 return Ok(snapshot);
479 }
480 if base_delay.is_zero() || max_delay < base_delay {
481 return Err(RuntimeError::InvalidInput(
482 "outbound retry delays require a non-zero base no greater than the ceiling"
483 .to_string(),
484 ));
485 }
486 crate::secret_gate::check_json_at(
487 &serde_json::json!({"last_error": &last_error}),
488 "message",
489 "last_error",
490 )?;
491 let properties = Self::outbound_retry_properties(
492 props,
493 attempted_at,
494 last_error,
495 base_delay,
496 max_delay,
497 )?;
498 self.replace_outbound_message_properties_as_owner(token, snapshot, properties)
499 .await
500 }
501
502 pub async fn mark_outbound_message_delivered(
514 &self,
515 token: &NamespaceToken,
516 id: Uuid,
517 delivered_at: String,
518 transport_message_id: Option<String>,
519 ) -> RuntimeResult<khive_storage::note::Note> {
520 crate::secret_gate::check_json_at(
521 &serde_json::json!({
522 "delivered_at": &delivered_at,
523 "transport_message_id": &transport_message_id,
524 }),
525 "message",
526 "delivered",
527 )?;
528 let snapshot = self.outbound_message(token, id).await?;
529 if Self::outbound_delivery_is_terminal(
530 snapshot.properties.as_ref().and_then(Value::as_object),
531 ) {
532 return Err(RuntimeError::InvalidInput(format!(
533 "outbound message {id} already has a terminal delivery outcome"
534 )));
535 }
536 let mut props = snapshot
537 .properties
538 .as_ref()
539 .and_then(Value::as_object)
540 .cloned()
541 .unwrap_or_default();
542 props.remove("delivery_attempts");
543 props.remove("next_attempt_at");
544 props.insert("delivery".into(), Value::String("delivered".into()));
545 props.insert("delivered_at".into(), Value::String(delivered_at));
546 if let Some(transport_message_id) = transport_message_id {
547 props.insert(
548 "transport_message_id".into(),
549 Value::String(transport_message_id),
550 );
551 }
552 self.replace_outbound_message_properties(token, snapshot, props)
553 .await
554 }
555
556 pub async fn mark_outbound_message_failed(
563 &self,
564 token: &NamespaceToken,
565 id: Uuid,
566 failed_at: String,
567 last_error: String,
568 ) -> RuntimeResult<khive_storage::note::Note> {
569 crate::secret_gate::check_json_at(
570 &serde_json::json!({
571 "failed_at": &failed_at,
572 "last_error": &last_error,
573 }),
574 "message",
575 "failed",
576 )?;
577 let snapshot = self.outbound_message(token, id).await?;
578 if Self::outbound_delivery_is_terminal(
579 snapshot.properties.as_ref().and_then(Value::as_object),
580 ) {
581 return Err(RuntimeError::InvalidInput(format!(
582 "outbound message {id} already has a terminal delivery outcome"
583 )));
584 }
585 let mut props = snapshot
586 .properties
587 .as_ref()
588 .and_then(Value::as_object)
589 .cloned()
590 .unwrap_or_default();
591 props.remove("delivery_attempts");
592 props.remove("next_attempt_at");
593 props.insert("delivery".into(), Value::String("failed".into()));
594 props.insert("failed_at".into(), Value::String(failed_at));
595 props.insert("last_error".into(), Value::String(last_error));
596 self.replace_outbound_message_properties(token, snapshot, props)
597 .await
598 }
599
600 pub async fn hold_outbound_message_external_id_unverifiable(
603 &self,
604 token: &NamespaceToken,
605 id: Uuid,
606 reason: String,
607 ) -> RuntimeResult<khive_storage::note::Note> {
608 crate::secret_gate::check_at(&reason, "message", "delivery_hold_reason")?;
609 let snapshot = self.outbound_message(token, id).await?;
610 let props = snapshot.properties.as_ref().and_then(Value::as_object);
611 if props
612 .and_then(|p| p.get("delivery_hold"))
613 .and_then(Value::as_str)
614 == Some("external_id_unverifiable")
615 {
616 return Ok(snapshot);
617 }
618 if Self::outbound_delivery_is_terminal(props)
619 || props
620 .and_then(|p| p.get("delivered_at"))
621 .is_some_and(|v| !v.is_null())
622 {
623 return Err(RuntimeError::InvalidInput(format!(
624 "outbound message {id} is no longer pending delivery"
625 )));
626 }
627 if props
628 .and_then(|p| p.get("external_id"))
629 .and_then(Value::as_str)
630 .is_none_or(|s| s.is_empty())
631 {
632 return Err(RuntimeError::InvalidInput(format!(
633 "outbound message {id} has no nonempty external_id to hold"
634 )));
635 }
636 let mut properties = props.cloned().unwrap_or_default();
637 properties.remove("delivery_attempts");
638 properties.remove("next_attempt_at");
639 properties.insert(
640 "delivery_hold".into(),
641 Value::String("external_id_unverifiable".into()),
642 );
643 properties.insert("delivery_hold_reason".into(), Value::String(reason));
644 properties.insert(
645 "delivery_hold_at".into(),
646 Value::String(chrono::Utc::now().to_rfc3339()),
647 );
648 self.replace_outbound_message_properties_as_owner(token, snapshot, properties)
649 .await
650 }
651
652 pub async fn list_outbound_external_id_holds_missing_diagnostic(
655 &self,
656 token: &NamespaceToken,
657 limit: u32,
658 ) -> RuntimeResult<Vec<khive_storage::note::Note>> {
659 let filter = NoteFilter {
660 kind: Some("message".to_string()),
661 property_filters: vec![
662 PropertyFilter {
663 json_path: "$.direction".into(),
664 op: FilterOp::Eq,
665 value: SqlValue::Text("outbound".into()),
666 },
667 PropertyFilter {
668 json_path: "$.to_actor".into(),
669 op: FilterOp::TextStartsWithIndexed,
670 value: SqlValue::Text("email:".into()),
671 },
672 PropertyFilter {
673 json_path: "$.delivery_hold".into(),
674 op: FilterOp::Eq,
675 value: SqlValue::Text("external_id_unverifiable".into()),
676 },
677 PropertyFilter {
678 json_path: "$.external_id_diagnostic_note_id".into(),
679 op: FilterOp::JsonTypeMissingOrNullIndexed,
680 value: SqlValue::Null,
681 },
682 ],
683 ..Default::default()
684 };
685 Ok(self
686 .notes(token)?
687 .query_notes_filtered_count_free(
688 token.namespace().as_str(),
689 &filter,
690 PageRequest {
691 limit: limit.min(200),
692 offset: 0,
693 },
694 )
695 .await?
696 .items)
697 }
698
699 pub async fn mark_outbound_external_id_diagnostic_recorded(
700 &self,
701 token: &NamespaceToken,
702 id: Uuid,
703 diagnostic_id: Uuid,
704 ) -> RuntimeResult<khive_storage::note::Note> {
705 let snapshot = self.outbound_message(token, id).await?;
706 let props = snapshot.properties.as_ref().and_then(Value::as_object);
707 let diagnostic_id_text = diagnostic_id.to_string();
708 if props
709 .and_then(|p| p.get("external_id_diagnostic_note_id"))
710 .and_then(Value::as_str)
711 == Some(diagnostic_id_text.as_str())
712 {
713 return Ok(snapshot);
714 }
715 if props
716 .and_then(|p| p.get("external_id_diagnostic_note_id"))
717 .and_then(Value::as_str)
718 .is_some()
719 {
720 return Err(RuntimeError::InvalidInput(format!(
721 "outbound message {id} already names a different diagnostic"
722 )));
723 }
724 if props
725 .and_then(|p| p.get("delivery_hold"))
726 .and_then(Value::as_str)
727 != Some("external_id_unverifiable")
728 {
729 return Err(RuntimeError::InvalidInput(format!(
730 "outbound message {id} has no external_id_unverifiable hold"
731 )));
732 }
733 let mut properties = props.cloned().unwrap_or_default();
734 properties.insert(
735 "external_id_diagnostic_note_id".into(),
736 Value::String(diagnostic_id_text),
737 );
738 self.replace_outbound_message_properties_as_owner(token, snapshot, properties)
739 .await
740 }
741
742 pub async fn record_outbound_external_id_diagnostic(
745 &self,
746 token: &NamespaceToken,
747 id: Uuid,
748 reason: &str,
749 ) -> RuntimeResult<khive_storage::note::Note> {
750 let key = format!("outbound-email-external-id-unverifiable:{id}");
751 let content = format!("Outbound email message {id} is held: {reason}");
752 let properties = serde_json::json!({
753 "diagnostic_code": "external_id_unverifiable",
754 "offending_message_id": id.to_string(),
755 });
756 let (note, _) = self
757 .create_note_with_options(
758 token,
759 "observation",
760 Some("Outbound email Message-ID unverifiable"),
761 &content,
762 None,
763 None,
764 None,
765 Some(properties),
766 vec![id],
767 None,
768 crate::note_write::NoteWriteOptions {
769 key: Some(key.clone()),
770 embed: Some(false),
771 ..Default::default()
772 },
773 )
774 .await?;
775 let linked = self
779 .list_edges(
780 token,
781 EdgeListFilter {
782 source_id: Some(note.id),
783 target_id: Some(id),
784 relations: vec![EdgeRelation::Annotates],
785 ..Default::default()
786 },
787 1,
788 0,
789 )
790 .await?
791 .len()
792 == 1;
793 if linked {
794 Ok(note)
795 } else {
796 Err(RuntimeError::InvalidInput(format!(
797 "outbound message {id} diagnostic has no annotation link"
798 )))
799 }
800 }
801
802 pub async fn mark_outbound_message_claim_failed(
808 &self,
809 token: &NamespaceToken,
810 id: Uuid,
811 failed_at: String,
812 last_error: String,
813 ) -> RuntimeResult<khive_storage::note::Note> {
814 let snapshot = self.outbound_message(token, id).await?;
815 self.mark_outbound_message_claim_failed_from_snapshot(
816 token, snapshot, failed_at, last_error,
817 )
818 .await
819 }
820
821 pub(super) async fn mark_outbound_message_claim_failed_from_snapshot(
822 &self,
823 token: &NamespaceToken,
824 snapshot: khive_storage::note::Note,
825 failed_at: String,
826 last_error: String,
827 ) -> RuntimeResult<khive_storage::note::Note> {
828 let props = snapshot.properties.as_ref().and_then(Value::as_object);
829 let has_claim = props
830 .and_then(|properties| properties.get("external_id"))
831 .and_then(Value::as_str)
832 .is_some_and(|value| !value.is_empty());
833 let has_delivery = props
834 .and_then(|properties| properties.get("delivered_at"))
835 .is_some_and(|value| !value.is_null());
836 if has_claim || has_delivery || Self::outbound_delivery_is_terminal(props) {
837 return Ok(snapshot);
838 }
839 crate::secret_gate::check_json_at(
840 &serde_json::json!({ "failed_at": &failed_at, "last_error": &last_error }),
841 "message",
842 "failed",
843 )?;
844 let mut properties = props.cloned().unwrap_or_default();
845 properties.remove("delivery_attempts");
846 properties.remove("next_attempt_at");
847 properties.insert("delivery".into(), Value::String("failed".into()));
848 properties.insert("failed_at".into(), Value::String(failed_at));
849 properties.insert("last_error".into(), Value::String(last_error));
850 self.replace_outbound_message_properties_as_owner(token, snapshot, properties)
851 .await
852 }
853}