1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
//! Trusted in-process recipient ingest. No registry verb dispatches this API.
//!
//! The authenticated node-loop issuer is not yet present, so this entry point
//! is unavailable to callers outside the runtime crate.
//!
//! ```compile_fail
//! use khive_runtime::KhiveRuntime;
//! fn main() { let _ = KhiveRuntime::ingest_verified_recipient; }
//! ```
//!
//! ```compile_fail
//! use khive_runtime::comm_recipient::LocalRecipientBinding;
//! fn main() { let _ = std::mem::size_of::<LocalRecipientBinding>(); }
//! ```
//!
//! ```compile_fail
//! use khive_runtime::comm_recipient::VerifiedInboundContent;
//! fn main() { let _ = std::mem::size_of::<VerifiedInboundContent>(); }
//! ```
#![allow(dead_code)] // The verified node-loop caller is supplied by #3537.
use crate::{KhiveRuntime, NamespaceToken, RuntimeError, RuntimeResult};
use khive_channel::InboundReceiptTicket;
use khive_db::stores::note::recipient::{
AckJournalEntry, QuarantineRecord, RecipientCommit, RecipientTransportStore,
};
pub use khive_db::stores::note::recipient::{
AcknowledgementRetirementReason, QuarantineReason, RecipientCommitResult, RecipientDisposition,
};
use khive_storage::Note;
use serde_json::json;
use uuid::Uuid;
/// An acknowledgement returned by the runtime's durable journal read.
///
/// Callers can inspect the binding and retry bookkeeping, but cannot manufacture
/// an entry from uncommitted delivery data. This type stores no signed bytes.
///
/// ```compile_fail
/// use khive_db::stores::note::recipient::AckJournalEntry;
/// use khive_runtime::comm_recipient::DueAcknowledgementEntry;
/// fn forge(journal: AckJournalEntry) {
/// let _ = DueAcknowledgementEntry { journal };
/// }
/// ```
#[derive(Debug)]
pub struct DueAcknowledgementEntry {
journal: AckJournalEntry,
}
impl DueAcknowledgementEntry {
/// The binding recorded when the delivery was committed or replayed.
pub fn binding(&self) -> &serde_json::Value {
&self.journal.binding
}
/// The disposition first committed for this logical message.
pub fn disposition(&self) -> RecipientDisposition {
self.journal.disposition
}
/// The exact delivery attempt this acknowledgement finishes.
pub fn delivery_attempt_id(&self) -> Uuid {
self.journal.delivery_attempt_id
}
/// The number of failed acknowledgement tries recorded durably.
pub fn attempt_count(&self) -> u64 {
self.journal.attempt_count
}
/// The earliest retry time, in microseconds since the Unix epoch.
pub fn not_before(&self) -> Option<i64> {
self.journal.not_before
}
}
impl KhiveRuntime {
/// Read due acknowledgements from comm's assigned runtime, oldest first.
/// This is the only construction path for [`DueAcknowledgementEntry`].
pub async fn due_acknowledgements(
&self,
now: i64,
limit: usize,
) -> RuntimeResult<Vec<DueAcknowledgementEntry>> {
let store = RecipientTransportStore::new(self.backend().pool_arc());
Ok(store
.list_due_acknowledgements(now, limit)
.await?
.into_iter()
.map(|journal| DueAcknowledgementEntry { journal })
.collect())
}
/// Finish a pending journal entry. Returns false for an absent or terminal
/// attempt, including an entry already finished by an earlier call.
pub async fn finish_acknowledgement(&self, delivery_attempt_id: Uuid) -> RuntimeResult<bool> {
let store = RecipientTransportStore::new(self.backend().pool_arc());
Ok(store.finish_acknowledgement(delivery_attempt_id).await?)
}
/// Record one failed try and its earliest retry time durably. An absent or
/// terminal attempt is unchanged and returns false.
pub async fn record_acknowledgement_failed_try(
&self,
delivery_attempt_id: Uuid,
not_before: i64,
) -> RuntimeResult<bool> {
let store = RecipientTransportStore::new(self.backend().pool_arc());
Ok(store
.record_acknowledgement_failed_try(delivery_attempt_id, not_before)
.await?)
}
/// Retire a pending entry after a permanent transport refusal. The message
/// and its committed receipt remain intact. An absent or terminal attempt
/// is unchanged and returns false.
pub async fn retire_acknowledgement(
&self,
delivery_attempt_id: Uuid,
reason: AcknowledgementRetirementReason,
) -> RuntimeResult<bool> {
let store = RecipientTransportStore::new(self.backend().pool_arc());
Ok(store
.retire_acknowledgement(delivery_attempt_id, reason)
.await?)
}
}
/// Local enrollment authority, supplied by the trusted slug owner. Never fill
/// this from plaintext or a wire request. Actor labels are local, not agent UUIDs.
pub(crate) struct LocalRecipientBinding {
pub actor: String,
pub realm: String,
pub slug: String,
pub agent_id: String,
pub device_id: Uuid,
pub key_epoch: u64,
}
/// The only purpose values a protocol-v1 sender may declare.
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(crate) enum DeclaredMessageKind {
Announce,
Report,
Ask,
}
impl DeclaredMessageKind {
fn as_str(self) -> &'static str {
match self {
Self::Announce => "announce",
Self::Report => "report",
Self::Ask => "ask",
}
}
}
/// Content has deliberately no sender, recipient, namespace or actor field.
/// Identity comes exclusively from local binding and an authenticated ticket.
pub(crate) enum VerifiedInboundContent {
Message {
content: String,
subject: Option<String>,
kind: Option<DeclaredMessageKind>,
in_reply_to: Option<Uuid>,
correlation: Option<String>,
sent_at: String,
},
Quarantine {
reason: QuarantineReason,
/// The parsed plaintext object, present exactly when the pair policy
/// refused a valid message (ADR-105 A.8 Receiving, step 4).
parsed_plaintext: Option<serde_json::Value>,
},
}
fn invalid(message: &str) -> RuntimeError {
RuntimeError::InvalidInput(message.into())
}
fn canonical_id(id: &str) -> bool {
Uuid::parse_str(id)
.ok()
.is_some_and(|u| u.to_string() == id)
}
impl KhiveRuntime {
/// Commit an authenticated delivery on comm's assigned runtime. The caller
/// MUST authenticate/decrypt and verify enrollment before creating the ticket.
/// This method performs no cryptography; it consumes the in-memory ticket,
/// checks local binding, and commits note/replay/ack/quarantine atomically.
///
/// For quarantine, delivery_item is the exact original JSON object, bounded
/// at 98,304 bytes by the protocol request limit. No blob store is involved.
/// A failure commits nothing and returns no recipient disposition. Retryable
/// storage errors must not be converted into quarantine by the caller.
/// A SecretDetected refusal is deterministic: redelivering the same payload
/// is refused again.
pub(crate) async fn ingest_verified_recipient(
&self,
token: &NamespaceToken,
local: &LocalRecipientBinding,
ticket: InboundReceiptTicket,
payload: VerifiedInboundContent,
delivery_item: Vec<u8>,
) -> RuntimeResult<RecipientCommitResult> {
let binding = ticket.binding();
if !canonical_id(&binding.sender_agent_id)
|| !canonical_id(&binding.recipient_agent_id)
|| !canonical_id(&local.agent_id)
{
return Err(invalid("verified receipt agents must be canonical UUIDs"));
}
if binding.protocol_version != 1
|| binding.recipient_agent_id != local.agent_id
|| binding.recipient_device_id != local.device_id
|| binding.recipient_key_epoch != local.key_epoch
{
return Err(invalid(
"verified delivery does not match local recipient binding",
));
}
if [
binding.recipient_key_epoch,
binding.contact_generation,
ticket.sender_key_epoch(),
]
.iter()
.any(|n| *n == 0 || *n > u32::MAX as u64)
{
return Err(invalid("invalid ticket epoch or generation"));
}
if local.actor.trim().is_empty()
|| local.slug.is_empty()
|| local.realm.is_empty()
|| local.realm.len() > 64
|| !local
.realm
.bytes()
.all(|b| b.is_ascii_lowercase() || b.is_ascii_digit() || b"._-".contains(&b))
{
return Err(invalid("invalid local recipient route"));
}
// ADR-105 A.8 step 4: a successfully opened replay is answered from
// its durable claim before interpreting the new attempt's plaintext.
let binding_value = serde_json::to_value(binding)
.map_err(|error| invalid(&format!("invalid receipt binding: {error}")))?;
let recipient_store = RecipientTransportStore::new(self.backend().pool_arc());
if let Some(replay) = recipient_store
.ack_if_replayed(binding_value.clone(), &local.actor)
.await?
{
return Ok(replay);
}
self.validate_note_kind("message")?;
let from = format!("khive1:{}/{}", local.realm, binding.sender_agent_id);
let received_at = chrono::Utc::now().to_rfc3339();
let mut in_reply_to = None;
let mut correlation = None;
let (note, disposition, quarantine) = match payload {
VerifiedInboundContent::Message {
content,
subject,
kind,
in_reply_to: parent,
correlation: message_correlation,
sent_at,
} => {
let parsed_sent_at = chrono::DateTime::parse_from_rfc3339(&sent_at)
.ok()
.filter(|stamp| stamp.offset().local_minus_utc() == 0);
if let Some(stamp) = parsed_sent_at {
// An empty body is a valid A.5 message and is stored like any other.
crate::secret_gate::check_at(&content, "note", "content")?;
if let Some(subject) = &subject {
crate::secret_gate::check_at(subject, "note", "name")?;
}
in_reply_to = parent;
correlation = message_correlation;
let mut note = Note::new(token.namespace().as_str(), "message", content);
note.name = subject.clone();
let mut props = json!({
"comm_schema_version": 1,
"from": from,
"from_actor": from,
"to": local.actor,
"to_actor": local.actor,
"direction": "inbound",
"read": false,
"received_at": received_at,
"channel_kind": "khive",
"channel_slug": local.slug,
"logical_message_id": binding.logical_message_id,
"sent_at": stamp
.with_timezone(&chrono::Utc)
.to_rfc3339_opts(chrono::SecondsFormat::AutoSi, true),
"message_kind": kind.map_or("unspecified", DeclaredMessageKind::as_str),
});
if let Some(subject) = subject {
props["subject"] = json!(subject);
}
if let Some(parent) = in_reply_to {
props["in_reply_to"] = json!(parent);
}
crate::secret_gate::check_json_at(&props, "note", "properties")?;
note.properties = Some(props);
(Some(note), RecipientDisposition::Stored, None)
} else {
(
None,
RecipientDisposition::Quarantined,
Some(QuarantineRecord {
reason: QuarantineReason::InvalidPlaintext,
delivery_item,
parsed_plaintext: None,
}),
)
}
}
VerifiedInboundContent::Quarantine {
reason,
parsed_plaintext,
} => {
// Retained plaintext passes the same write-time secret gate as a stored message.
// Refuse before the quarantine/replay/ack transaction can write.
if let Some(plaintext) = &parsed_plaintext {
crate::secret_gate::check_json_at(plaintext, "quarantine", "parsed_plaintext")?;
}
(
None,
RecipientDisposition::Quarantined,
Some(QuarantineRecord {
reason,
delivery_item,
parsed_plaintext,
}),
)
}
};
let result = recipient_store
.commit(RecipientCommit {
note,
recipient_actor: local.actor.clone(),
binding: binding_value,
sender_agent_id: binding.sender_agent_id.clone(),
logical_message_id: binding.logical_message_id,
delivery_attempt_id: binding.delivery_attempt_id,
disposition,
quarantine,
in_reply_to,
correlation,
})
.await?;
if let Some(note) = &result.note {
// Like ordinary ingest, indexing is best-effort after the durable commit.
if let Ok(fts) = self.text_for_notes(token) {
if let Err(error) = fts
.upsert_document(crate::curation::note_fts_document(note))
.await
{
tracing::warn!(note_id=%note.id,error=%error,"verified ingest FTS indexing failed");
}
}
for model in self.registered_embedding_model_names() {
match self
.embed_document_with_model_outcome_for_token(
token,
&model,
crate::curation::note_embedding_text_ref(note),
)
.await
{
Ok(outcome) => {
if let Err(error) = self
.publish_note_vector_revision(token, note, &model, &outcome.vector)
.await
{
tracing::warn!(note_id=%note.id,error=%error,"verified ingest vector indexing failed");
}
}
Err(error) => {
tracing::warn!(note_id=%note.id,error=%error,"verified ingest embedding failed")
}
}
}
}
Ok(result)
}
}
#[cfg(test)]
mod tests;
#[cfg(test)]
mod acknowledgement_journal_tests;