Skip to main content

khive_runtime/curation/
outbound_messages.rs

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    /// Claim `external_id` on an outbound `message` note through the
9    /// ADR-124-sanctioned store-level owner path, bypassing the
10    /// caller-facing owner-established-property refusal in
11    /// [`Self::update_note_with_embedding_report`] (and its crate-internal prepare path). This is deliberately
12    /// NOT exposed through any registered verb (ADR-124's stated bound): it is
13    /// reachable only from pack/runtime code that owns outbox bookkeeping for
14    /// the `message` note kind.
15    ///
16    /// Refuses (returns `Err`, never writes) unless the live row is a
17    /// `message` note, `properties.direction == "outbound"`, and
18    /// `properties.external_id` is currently absent or empty. The claim is
19    /// committed against that exact snapshot and advances its timestamp so
20    /// competing claims and delivery-outcome CAS writes cannot overwrite it.
21    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    /// Non-wire outbox scan for the channel delivery loops.
76    ///
77    /// Fetches live `message` notes matching the SQL-side pending predicate
78    /// newest-first (`created_at DESC, id ASC`), bounded by an internal page
79    /// cap. Direction, `delivered_at`, terminal `delivery` state, the optional
80    /// `to_actor` channel prefix, and `next_attempt_at` are filtered by SQLite
81    /// before the page bound. Pending means `delivered_at`
82    /// is absent or null, `properties.delivery` carries no terminal state
83    /// (`"delivered"` / `"failed"`), and a valid `next_attempt_at` is absent
84    /// or due (ADR-122 §1). Malformed legacy deadlines fail open so a bad
85    /// property cannot strand mail forever.
86    ///
87    /// The channel prefix has to be in the statement, not applied to the
88    /// fetched page: every actor-to-actor outbound row matches the pending
89    /// predicate forever (nothing marks those delivered). A full `name:`
90    /// channel prefix also supplies an indexed bucket equality, followed by
91    /// an indexed deadline bound; arbitrary partial prefixes retain the
92    /// recipient range. The final newest-first sort preserves delivery order.
93    /// This lives on the runtime rather than going through the wire registry
94    /// for the same reason as
95    /// [`Self::claim_outbound_message_external_id`]: the delivery loop must
96    /// scan the backend that actually holds comm's notes, and under a
97    /// `[packs.comm]` backend assignment that is not the backend serving the
98    /// generic kg verbs.
99    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    /// Channel-specific scan for the delivery pass. An explicit slug belongs
115    /// only to its named adapter; a missing slug is eligible only when that
116    /// kind has exactly one configured adapter. Querying these disjoint
117    /// partitions before paging prevents held unknown or ambiguous rows from
118    /// crowding an eligible row out of the finite outbox scan window.
119    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        // The former Rust predicate compared `timestamp_micros()`, so a
168        // deadline within the current microsecond counted as due. Preserve
169        // that boundary when comparing the SQL function's nanosecond keys.
170        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    /// Load a live outbound `message` note, returning `InvalidInput`
251    /// otherwise. Guard shared by the delivery-outcome markers: they take
252    /// caller-supplied UUIDs, and the generic note patch path would happily
253    /// stamp delivery properties onto any note kind.
254    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        // Delivery outcomes preserve transport-owned route fields from the loaded snapshot.
305        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        // Storage assigns the persisted revision; the pre-write snapshot
313        // still carries the old version even when this CAS succeeds.
314        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        // Owner operations preserve any transport evidence on the snapshot.
337        // Re-running the public full-row guard would reject that existing evidence.
338        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        // Return the storage-assigned revision, just as the claim path does.
346        store
347            .get_note(id)
348            .await?
349            .ok_or_else(|| RuntimeError::NotFound(format!("note {id}")))
350    }
351
352    /// True once a `delivery` outcome has been terminally recorded
353    /// (`"delivered"` or `"failed"`). Shared by every delivery-outcome
354    /// marker: concurrent outbox workers (two daemon processes overlapping
355    /// during a restart, per ADR-122 §4/Consequences) can both load the same
356    /// pending note before either writes, so a marker must re-check the
357    /// freshly loaded snapshot rather than trust the compare-and-swap alone
358    /// -- the CAS only rejects a write against a snapshot that has since
359    /// changed, not a write that starts from an up-to-date terminal snapshot
360    /// and would otherwise happily overwrite it with a different outcome.
361    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    /// Record a transient transport failure while leaving the message
369    /// pending. The retry deadline is derived from the incremented persisted
370    /// attempt count, using `base_delay * 2^(attempt - 1)` capped at
371    /// `max_delay`.
372    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    /// Schedule an external-id claim retry only for a still-unclaimed
458    /// outbound snapshot. Existing claims and terminal outcomes are unchanged.
459    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    /// Mark an outbound `message` note delivered by merging the ADR-122 §1
503    /// terminal-outcome properties (`delivery = "delivered"`, `delivered_at`,
504    /// and `transport_message_id` when the transport minted one), and clearing
505    /// `delivery_attempts` / `next_attempt_at`. The compare-and-swap protects
506    /// unrelated properties from a concurrent full-row overwrite;
507    /// `delivered_at` remains deliberately caller-patchable, pinned by
508    /// `generic_update_can_still_patch_delivered_at_on_message_note`.
509    /// Refuses (`InvalidInput`) unless `id` names a live outbound `message`
510    /// note. Non-wire companion to
511    /// [`Self::list_undelivered_outbound_messages`] so the delivery loop
512    /// writes the backend that holds the note.
513    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    /// Record a permanent delivery failure on an outbound `message` note:
557    /// `delivery = "failed"`, `failed_at`, `last_error` (ADR-122 §2 — an
558    /// allowlist rejection must be recorded, not skipped, or the row stays
559    /// pending forever while the caller saw `ok: true`). Any retry counter and
560    /// deadline are cleared because the outcome is terminal. Refuses
561    /// (`InvalidInput`) unless `id` names a live outbound `message` note.
562    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    /// Owner-only visible hold for an outbound email whose stored Message-ID
601    /// cannot be bound to its own row and configured sending domain.
602    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    /// Bounded maintenance scan for holds whose keyed diagnostic still needs
653    /// confirmation. These rows never enter the ordinary send selection.
654    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    /// One keyed operator-visible observation, atomically annotated to the
743    /// offending message. A replay cannot create a second diagnostic.
744    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        // A keyed replay can return an existing note without running the
776        // creation-only annotation work. Verify the link before marking the
777        // message's diagnostic as recorded.
778        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    /// Park a deterministic external-id claim refusal only while the exact
803    /// current outbound snapshot remains unclaimed. An existing claim or a
804    /// terminal delivery outcome is returned unchanged. Non-wire owner API:
805    /// a second worker must not turn another worker's successful claim into
806    /// a permanent delivery failure.
807    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}