diff --git a/bin/chat-cli/src/app.rs b/bin/chat-cli/src/app.rs index 30b9e24..7499181 100644 --- a/bin/chat-cli/src/app.rs +++ b/bin/chat-cli/src/app.rs @@ -15,6 +15,21 @@ pub struct DisplayMessage { pub from_self: bool, pub content: String, pub timestamp: u64, + pub message_id: Option, + #[serde(default)] + pub delivered_to: Vec, +} + +impl DisplayMessage { + fn new(from_self: bool, content: String) -> Self { + Self { + from_self, + content, + timestamp: now(), + message_id: None, + delivered_to: Vec::new(), + } + } } #[derive(Debug, Clone, Serialize, Deserialize)] @@ -189,11 +204,58 @@ where let Some(session) = self.state.chats.get_mut(&chat_id) else { return; }; - session.messages.push(DisplayMessage { - from_self: false, - content: String::from_utf8_lossy(&content).into_owned(), - timestamp: now(), - }); + session.messages.push(DisplayMessage::new( + false, + String::from_utf8_lossy(&content).into_owned(), + )); + } + Event::MessageAcked { + convo_id, + message_id, + acked_by, + } => { + let Some(session) = self.state.chats.get_mut(convo_id.as_ref()) else { + return; + }; + let Some(message) = session + .messages + .iter_mut() + .find(|m| m.message_id.as_deref() == Some(message_id.as_str())) + else { + return; // sent before this session, or not ours + }; + let peer = acked_by.map_or_else( + || "a member".to_string(), + |s| { + let id = s.account.unwrap_or(s.local_identity); + format!("{}…", &id.as_str()[..8.min(id.as_str().len())]) + }, + ); + if !message.delivered_to.contains(&peer) { + message.delivered_to.push(peer); + } + } + Event::MessageMissing { + convo_id, + sender_hint, + .. + } => { + let Some(session) = self.state.chats.get(convo_id.as_ref()) else { + return; + }; + // The hint is not authenticated (see `Event::MessageMissing`), + // so name the author loosely rather than as an established fact. + let author = sender_hint.map_or_else( + || "a member".to_string(), + |s| { + let id = s.account.unwrap_or(s.local_identity); + format!("{}…", &id.as_str()[..8.min(id.as_str().len())]) + }, + ); + self.status = format!( + "A message from {author} never arrived in '{}'.", + session.display_name() + ); } Event::InboundError { message } => { self.status = format!("Could not process incoming message: {message}"); @@ -209,16 +271,16 @@ where .clone() .ok_or_else(|| anyhow::anyhow!("No active chat. Use /connect or /switch first."))?; - self.client + let message_id = self + .client .send_message(&chat_id, content.as_bytes()) .map_err(|e| anyhow::anyhow!("{e:?}"))?; if let Some(session) = self.state.chats.get_mut(&chat_id) { - session.messages.push(DisplayMessage { - from_self: true, - content: content.to_string(), - timestamp: now(), - }); + let mut message = DisplayMessage::new(true, content.to_string()); + // Kept so `MessageAcked` can find this message again. + message.message_id = Some(message_id); + session.messages.push(message); } self.save_state()?; @@ -226,11 +288,8 @@ where } fn add_system_message(&mut self, content: &str) { - self.command_output.push(DisplayMessage { - from_self: true, - content: content.to_string(), - timestamp: now(), - }); + self.command_output + .push(DisplayMessage::new(true, content.to_string())); } pub fn handle_command(&mut self, cmd: &str) -> Result> { @@ -273,7 +332,8 @@ where .client .create_direct_conversation(args) .map_err(|e| anyhow::anyhow!("{e:?}"))?; - self.client + let message_id = self + .client .send_message(&chat_id, initial.as_bytes()) .map_err(|e| anyhow::anyhow!("{e:?}"))?; @@ -283,11 +343,9 @@ where nickname: None, messages: Vec::new(), }; - session.messages.push(DisplayMessage { - from_self: true, - content: initial, - timestamp: now(), - }); + let mut message = DisplayMessage::new(true, initial); + message.message_id = Some(message_id); + session.messages.push(message); self.state.chats.insert(chat_id.clone(), session); self.set_active_chat(Some(chat_id)); self.save_state()?; diff --git a/bin/chat-cli/src/ui.rs b/bin/chat-cli/src/ui.rs index 595c48e..4914dc6 100644 --- a/bin/chat-cli/src/ui.rs +++ b/bin/chat-cli/src/ui.rs @@ -156,6 +156,18 @@ where remaining = tail; } + // Delivery receipts for our own sends: the peers whose later + // messages showed they hold this one. + if !msg.delivered_to.is_empty() { + items.push(ListItem::new(Line::from(vec![ + Span::raw(indent.clone()), + Span::styled( + format!("↳ delivered to {}", msg.delivered_to.join(", ")), + Style::default().fg(Color::DarkGray), + ), + ]))); + } + items }) .collect(); diff --git a/core/conversations/src/causal_history.rs b/core/conversations/src/causal_history.rs index 7afc7b1..ebdf1e2 100644 --- a/core/conversations/src/causal_history.rs +++ b/core/conversations/src/causal_history.rs @@ -9,12 +9,19 @@ //! - assign a deterministic message ID + Lamport timestamp to outbound msgs //! - attach a bounded causal-history frontier to each outbound message //! - on receive, detect referenced-but-unseen message IDs (gaps) +//! - on receive, detect references to *our own* messages (acknowledgements) //! //! Out of scope here: bloom-filter acknowledgements, //! resend / outgoing buffer, incoming reorder buffer, Store-based recovery. //! This is detection only — an out-of-order message is still delivered to //! the application, but the gap it implies is reported. //! +//! The same references also show who received our messages: a peer that names +//! one of ours must have had it. Nothing is sent back — the acknowledgement +//! rides on whatever the peer says next — so a silent peer never acknowledges, +//! and neither does one that speaks after our message has dropped out of its +//! [`CAUSAL_HISTORY_LEN`]-entry frontier. +//! //! State is in-memory and session-scoped, matching the crate's current //! in-memory MLS state. @@ -72,6 +79,21 @@ pub struct MissingMessage { pub frontier: Frontier, } +/// A peer acknowledging one of our messages: it named that message in the +/// causal history of a message it sent, so it held ours at the time. +/// +/// Evidence of *delivery to a peer's client*, not of a human reading it. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct DeliveryAck { + pub conversation_id: String, + /// The message of ours the peer acknowledged. + pub message_id: String, + /// The acknowledging peer's `sender_id`, verbatim off the wire — + /// self-asserted like [`Frontier::sender_id`], not bound to the MLS + /// identity that sent the payload. + pub acked_by: String, +} + /// Per-conversation causal state. #[derive(Debug, Default)] struct ConvoState { @@ -84,6 +106,11 @@ struct ConvoState { frontiers: VecDeque, /// Missing IDs already reported, so a gap is surfaced exactly once. reported_missing: HashSet, + /// IDs of messages we authored: a reference to one is an acknowledgement. + own: HashSet, + /// Which peers have acknowledged each of our messages, so each is + /// surfaced exactly once. + acked_by: HashMap>, } impl ConvoState { @@ -102,6 +129,9 @@ struct Inner { convos: HashMap, /// Detected gaps, drained by the client (future #97 event bus). missing: Vec, + /// Detected acknowledgements of our own messages, drained alongside + /// `missing`. + acks: Vec, } /// Session-scoped causal-history store shared by every `GroupV1Convo` @@ -143,7 +173,9 @@ impl CausalHistoryStore { .collect(); // Our own message joins the seen-set so it appears in our future - // causal history (and, later, so we can ack peers' references to it). + // causal history, and the own-set so a peer referencing it back is + // recognised as acknowledging this send. + state.own.insert(message_id.clone()); state.record_seen(frontier); ReliablePayload { @@ -166,7 +198,11 @@ impl CausalHistoryStore { payload: &ReliablePayload, ) -> Vec { let mut inner = self.inner.borrow_mut(); - let Inner { convos, missing } = &mut *inner; + let Inner { + convos, + missing, + acks, + } = &mut *inner; let state = convos.entry(conversation_id.to_owned()).or_default(); // Lamport merge: the next local send will be strictly greater than @@ -175,6 +211,23 @@ impl CausalHistoryStore { let mut detected = Vec::new(); for entry in &payload.causal_history { + // The sender named one of ours, so it has it. Reported once per + // peer per message, and never for the message's own author. + if state.own.contains(&entry.message_id) + && payload.sender_id != entry.sender_id + && state + .acked_by + .entry(entry.message_id.clone()) + .or_default() + .insert(payload.sender_id.clone()) + { + acks.push(DeliveryAck { + conversation_id: conversation_id.to_owned(), + message_id: entry.message_id.clone(), + acked_by: payload.sender_id.clone(), + }); + } + let frontier = Frontier::new(entry.sender_id.clone(), entry.message_id.clone()); if !state.seen.contains(&frontier) && state.reported_missing.insert(frontier.clone()) { let m = MissingMessage { @@ -201,6 +254,11 @@ impl CausalHistoryStore { pub fn take_missing(&self) -> Vec { std::mem::take(&mut self.inner.borrow_mut().missing) } + + /// Drain all acknowledgements of our own messages detected so far. + pub fn take_acks(&self) -> Vec { + std::mem::take(&mut self.inner.borrow_mut().acks) + } } /// Deterministic, collision-resistant message ID. @@ -293,6 +351,87 @@ mod tests { assert_eq!(missing[0].frontier.sender_id(), "alice"); } + /// Bob replies after receiving Alice's message, so his causal history + /// names it — that reference is the acknowledgement. + #[test] + fn a_peer_referencing_our_message_acknowledges_it() { + let alice = CausalHistoryStore::new(); + let bob = CausalHistoryStore::new(); + + let a1 = payload(&alice, "c", "alice", b"hello"); + bob.on_receive("c", &a1); + let b1 = payload(&bob, "c", "bob", b"hi back"); + alice.on_receive("c", &b1); + + assert_eq!( + alice.take_acks(), + vec![DeliveryAck { + conversation_id: "c".to_owned(), + message_id: a1.message_id.clone(), + acked_by: "bob".to_owned(), + }] + ); + // Draining clears the report. + assert!(alice.take_acks().is_empty()); + } + + /// Every member that replies acknowledges separately, which is what lets an + /// application list the peers that hold a message. + #[test] + fn each_peer_acknowledges_separately() { + let alice = CausalHistoryStore::new(); + let bob = CausalHistoryStore::new(); + let carol = CausalHistoryStore::new(); + + let a1 = payload(&alice, "c", "alice", b"hello all"); + bob.on_receive("c", &a1); + carol.on_receive("c", &a1); + alice.on_receive("c", &payload(&bob, "c", "bob", b"bob here")); + alice.on_receive("c", &payload(&carol, "c", "carol", b"carol here")); + + let holders: Vec = alice + .take_acks() + .into_iter() + .filter(|a| a.message_id == a1.message_id) + .map(|a| a.acked_by) + .collect(); + assert_eq!(holders, vec!["bob".to_owned(), "carol".to_owned()]); + } + + /// Bob keeps naming the message in later sends; the application is told + /// once. + #[test] + fn a_peer_acknowledges_a_message_only_once() { + let alice = CausalHistoryStore::new(); + let bob = CausalHistoryStore::new(); + + let a1 = payload(&alice, "c", "alice", b"hello"); + bob.on_receive("c", &a1); + alice.on_receive("c", &payload(&bob, "c", "bob", b"first reply")); + alice.take_acks(); + alice.on_receive("c", &payload(&bob, "c", "bob", b"second reply")); + + assert!( + alice.take_acks().is_empty(), + "a peer's acknowledgement of one message is reported once" + ); + } + + /// Carol's reply names Bob's message, not ours — nothing for us to report. + #[test] + fn a_reference_to_someone_elses_message_is_not_our_acknowledgement() { + let alice = CausalHistoryStore::new(); + let bob = CausalHistoryStore::new(); + let carol = CausalHistoryStore::new(); + + let b1 = payload(&bob, "c", "bob", b"bob speaks"); + carol.on_receive("c", &b1); + // Alice observes Carol's reply, which references Bob's message only. + alice.on_receive("c", &payload(&carol, "c", "carol", b"carol replies")); + + assert!(alice.take_acks().is_empty()); + } + #[test] fn a_gap_is_reported_only_once() { let sender = CausalHistoryStore::new(); diff --git a/core/conversations/src/conversation.rs b/core/conversations/src/conversation.rs index 2d016f3..7a9f1a1 100644 --- a/core/conversations/src/conversation.rs +++ b/core/conversations/src/conversation.rs @@ -16,10 +16,19 @@ use shared_traits::IdentIdRef; pub type ConversationId = String; pub type ConversationIdRef<'a> = &'a str; +/// Identifies one message within a conversation, as carried by the +/// causal-history envelope. Handed back by a send so a caller can match later +/// observations — acknowledgements, gaps — to the message that produced them. +pub type MessageId = String; + /// Behaviour shared by every conversation kind. pub(crate) trait Convo: Identified + Send { - fn send_content(&mut self, cx: &mut ServiceContext, content: &[u8]) - -> Result<(), ChatError>; + /// Encrypt and publish `content`, returning the id assigned to it. + fn send_content( + &mut self, + cx: &mut ServiceContext, + content: &[u8], + ) -> Result; /// Decrypts and processes an incoming encrypted frame. /// diff --git a/core/conversations/src/conversation/direct_v1.rs b/core/conversations/src/conversation/direct_v1.rs index b1cf04a..d9d191c 100644 --- a/core/conversations/src/conversation/direct_v1.rs +++ b/core/conversations/src/conversation/direct_v1.rs @@ -2,7 +2,7 @@ use chat_proto::logoschat::encryption::EncryptedPayload; use shared_traits::IdentIdRef; use crate::{ - ChatError, ExternalServices, + ChatError, ExternalServices, MessageId, conversation::{ConversationIdRef, Convo, GroupConvo, GroupV1Convo, Identified}, service_context::ServiceContext, }; @@ -43,7 +43,7 @@ where &mut self, cx: &mut ServiceContext, content: &[u8], - ) -> Result<(), super::ChatError> { + ) -> Result { self.inner_group.send_content(cx, content) } diff --git a/core/conversations/src/conversation/group_v1.rs b/core/conversations/src/conversation/group_v1.rs index 665a9f4..6c8fd22 100644 --- a/core/conversations/src/conversation/group_v1.rs +++ b/core/conversations/src/conversation/group_v1.rs @@ -12,7 +12,7 @@ use shared_traits::IdentIdRef; use std::collections::VecDeque; use tracing::debug; -use crate::conversation::ConversationIdRef; +use crate::conversation::{ConversationIdRef, MessageId}; use crate::inbox_v2::MlsProvider; use crate::service_context::{ExternalServices, ServiceContext}; @@ -177,7 +177,7 @@ impl GroupV1Convo { &mut self, content: &[u8], cx: &mut ServiceContext, - ) -> Result<(), ChatError> { + ) -> Result { let sender_id = cx.mls_identity.id().as_str(); let reliable = cx.causal.on_send(&self.convo_id, sender_id, content); let wire = reliable.encode_to_vec(); @@ -188,7 +188,8 @@ impl GroupV1Convo { .unwrap(); let msg_bytes = mls_message_out.to_bytes().unwrap(); - self.send_payload(cx, msg_bytes) + self.send_payload(cx, msg_bytes)?; + Ok(reliable.message_id) } // Publish outboubound payloads to the DeliveryService @@ -233,7 +234,7 @@ impl Convo for GroupV1Convo { &mut self, cx: &mut ServiceContext, content: &[u8], - ) -> Result<(), ChatError> { + ) -> Result { self.send_message(content, cx) } diff --git a/core/conversations/src/conversation/group_v2.rs b/core/conversations/src/conversation/group_v2.rs index 5a4172e..d37b2be 100644 --- a/core/conversations/src/conversation/group_v2.rs +++ b/core/conversations/src/conversation/group_v2.rs @@ -10,6 +10,7 @@ use crate::{Content, WakeupService}; use alloy::signers::local::PrivateKeySigner; use blake2::{Blake2b, Digest, digest::consts::U6}; use chat_proto::logoschat::encryption::{EncryptedPayload, Plaintext, encrypted_payload}; +use chat_proto::logoschat::reliability::ReliablePayload; use de_mls::protos::de_mls::messages::v1::{ AppMessage as AppMessageProto, MemberWelcome, app_message, }; @@ -31,7 +32,7 @@ use std::time::{Duration, SystemTime, UNIX_EPOCH}; use tracing::{info, instrument}; use crate::IdentityProvider; -use crate::conversation::{ConversationIdRef, ExternalServices, ServiceContext}; +use crate::conversation::{ConversationIdRef, ExternalServices, MessageId, ServiceContext}; use crate::{ ConvoOutcome, DeliveryService, RegistrationService, conversation::{ChatError, Convo, GroupConvo, Identified}, @@ -304,14 +305,20 @@ where &mut self, service_ctx: &mut super::ServiceContext, content: &[u8], - ) -> Result<(), ChatError> { + ) -> Result { + let reliable = service_ctx.causal.on_send( + &self.convo_id, + service_ctx.mls_identity.id().as_str(), + content, + ); + self.conversation.send_message( &service_ctx.mls_provider, &service_ctx.mls_identity, - content.to_vec(), + reliable.encode_to_vec(), )?; self.after_op(service_ctx)?; - Ok(()) + Ok(reliable.message_id) } #[instrument(name = "groupv2.handle_frame", skip_all, fields(user_id = %service_ctx.mls_identity.display_name()))] @@ -341,7 +348,7 @@ where self.conversation .poll(&service_ctx.mls_provider, &service_ctx.mls_identity); let events = self.after_op(service_ctx)?; // route + publish + re-arm, returns events - Ok(self.outcome_from_events(&events)) + self.outcome_from_events(service_ctx, &events) } #[instrument(name = "groupv2.wakeup", skip_all, fields(user_id = %ctx.mls_identity.display_name()))] @@ -355,7 +362,7 @@ where tracing::warn!(convo = %self.convo_id, "conversation requested teardown"); } let events = self.after_op(ctx)?; // publish what poll produced + re-arm alarm - Ok(self.outcome_from_events(&events)) + self.outcome_from_events(ctx, &events) } fn members(&self) -> Result>, ChatError> { @@ -491,27 +498,47 @@ impl GroupV2Convo { Ok(events) } - fn outcome_from_events(&self, events: &[ConversationEvent]) -> ConvoOutcome { - let content = events.iter().find_map(|evt| match evt { - ConversationEvent::ConversationMessage(AppMessageProto { - payload: Some(app_message::Payload::ConversationMessage(cm)), - }) => Some(Content { - bytes: cm.message.clone(), - encoded_credential: cm.sender.clone(), - }), - _ => None, - }); + /// Turn drained de-mls events into a [`ConvoOutcome`], unwrapping the + /// message from its causal-history envelope. + /// + /// An outcome holds one message and de-mls emits at most one per frame, so + /// the first wins. A second would be dropped without being recorded as + /// seen, leaving a later reference to report it missing. + fn outcome_from_events( + &self, + service_ctx: &ServiceContext, + events: &[ConversationEvent], + ) -> Result { + let content = events + .iter() + .find_map(|evt| match evt { + ConversationEvent::ConversationMessage(AppMessageProto { + payload: Some(app_message::Payload::ConversationMessage(cm)), + }) => Some(cm), + _ => None, + }) + .map(|cm| -> Result { + let reliable = + ReliablePayload::decode(cm.message.as_slice()).map_err(ChatError::generic)?; + service_ctx.causal.on_receive(&self.convo_id, &reliable); + Ok(Content { + bytes: reliable.content.to_vec(), + encoded_credential: cm.sender.clone(), + }) + }) + .transpose()?; + let members_changed = events.iter().any(|evt| { matches!( evt, ConversationEvent::CommitApplied(_) | ConversationEvent::WelcomeReady { .. } ) }); - ConvoOutcome { + Ok(ConvoOutcome { convo_id: self.convo_id.clone(), content, members_changed, - } + }) } } diff --git a/core/conversations/src/core.rs b/core/conversations/src/core.rs index b696594..f6003f1 100644 --- a/core/conversations/src/core.rs +++ b/core/conversations/src/core.rs @@ -1,6 +1,6 @@ -use crate::causal_history::{CausalHistoryStore, MissingMessage}; +use crate::causal_history::{CausalHistoryStore, DeliveryAck, MissingMessage}; use crate::conversation::{ - ConversationIdRef, DirectV1Convo, GroupV1Convo, GroupV2Convo, Identified, + ConversationIdRef, DirectV1Convo, GroupV1Convo, GroupV2Convo, Identified, MessageId, }; use crate::service_context::{ExternalServices, ServiceContext}; use crate::types::ConvoMetadata; @@ -329,8 +329,16 @@ impl<'a, S: ExternalServices + 'static> Core { self.services.causal.take_missing() } - /// Encrypt and publish `content` to an existing conversation. - pub fn send_content(&mut self, convo_id: &str, content: &[u8]) -> Result<(), ChatError> { + /// Drain the acknowledgements observed since the last call: peers that + /// referenced one of our messages, and so demonstrably hold it. + pub fn take_acks(&self) -> Vec { + self.services.causal.take_acks() + } + + /// Encrypt and publish `content` to an existing conversation, returning the + /// id assigned to the message so later acknowledgements can be matched to + /// it. + pub fn send_content(&mut self, convo_id: &str, content: &[u8]) -> Result { if self.cached_convos.contains_key(convo_id) { let convo = self .cached_convos @@ -508,7 +516,7 @@ impl Convo for ConvoTypeOwned { &mut self, cx: &mut ServiceContext, content: &[u8], - ) -> Result<(), ChatError> { + ) -> Result { match self { ConvoTypeOwned::Group(group_convo) => group_convo.send_content(cx, content), ConvoTypeOwned::Direct(convo) => convo.send_content(cx, content), diff --git a/core/conversations/src/lib.rs b/core/conversations/src/lib.rs index 4dadf8d..ca319fb 100644 --- a/core/conversations/src/lib.rs +++ b/core/conversations/src/lib.rs @@ -10,10 +10,10 @@ mod service_traits; mod types; mod utils; -pub use causal_history::{Frontier, MissingMessage}; +pub use causal_history::{DeliveryAck, Frontier, MissingMessage}; pub use chat_sqlite::ChatStorage; pub use chat_sqlite::StorageConfig; -pub use conversation::GroupV2Clock; +pub use conversation::{GroupV2Clock, MessageId}; pub use core::{ConversationId, Core}; /// Timing/policy for GroupV2 conversations (de-mls's per-conversation config). /// Defaults to the de-mls library defaults; inject via diff --git a/core/integration_tests_core/src/test_client.rs b/core/integration_tests_core/src/test_client.rs index bf95b15..2c5a0e1 100644 --- a/core/integration_tests_core/src/test_client.rs +++ b/core/integration_tests_core/src/test_client.rs @@ -84,6 +84,13 @@ impl TestClient { outcomes } + /// Poll and discard every payload waiting for this client — simulates + /// frames the transport never delivered. + pub fn drop_pending_payloads(&mut self) { + let ds = self.inner.ds(); + while ds.poll().is_some() {} + } + pub fn received_messages(&self) -> &[ReceivedMessage>] { &self.received_messages } diff --git a/core/integration_tests_core/tests/test_group_v2.rs b/core/integration_tests_core/tests/test_group_v2.rs index acfe88a..ef672d3 100644 --- a/core/integration_tests_core/tests/test_group_v2.rs +++ b/core/integration_tests_core/tests/test_group_v2.rs @@ -1,4 +1,5 @@ use integration_tests_core::TestHarness; +use libchat::{DeliveryAck, MissingMessage}; use tracing::info; #[test] @@ -310,3 +311,145 @@ fn direct_v1_then_group_v2_reuses_key_package() { "raya did not join the group after a direct chat" ); } + +/// End-to-end causal-history gap detection on GroupV2. +/// +/// Saro and Raya share a GroupV2 conversation. Saro sends three messages; the +/// second never reaches Raya. The third carries the second in its causal +/// history, so Raya must detect and report the gap. +#[test] +fn missing_group_v2_message_is_detected() { + let mut harness = TestHarness::<2>::new(|_, _| {}); + + let participants = &[&harness.raya().addr()]; + let convo_id = harness + .saro() + .create_group_convo_v2(participants, "", "") + .expect("saro create group"); + + // Carry the invite through (commit, WelcomeReady, inbox routing, welcome + // accept) until Raya has joined. + harness.process_until_label("raya joins", |h| h.raya().convo_count() == 1); + + // M1 is delivered normally. + harness + .saro() + .send_content(&convo_id, b"first") + .expect("saro send m1"); + harness.process_until_label("raya gets m1", |h| h.raya().check(&convo_id, b"first")); + assert!( + harness.raya().take_missing_messages().is_empty(), + "no gap expected while every message is delivered" + ); + + // M2 is published but never reaches Raya. The settle above drained Raya's + // queue, and `send_content` publishes synchronously, so discarding what is + // pending now drops that message and nothing of the group protocol. + harness + .saro() + .send_content(&convo_id, b"second") + .expect("saro send m2"); + harness.raya().drop_pending_payloads(); + + // M3 is delivered; its causal history references the dropped M2. + harness + .saro() + .send_content(&convo_id, b"third") + .expect("saro send m3"); + harness.process_until_label("raya gets m3", |h| h.raya().check(&convo_id, b"third")); + + let missing: Vec = harness.raya().take_missing_messages(); + assert_eq!(missing.len(), 1, "exactly one message should be missing"); + assert_eq!(missing[0].conversation_id, convo_id); + assert!( + !missing[0].frontier.message_id().is_empty(), + "the missing message must be identified" + ); + // The causal sender hint carries the MLS identity id ("saro") — the same + // value de-mls stamps as the message's authenticated member id — not the + // signer id the inbox and registry key on. + assert_eq!( + missing[0].frontier.sender_id(), + "saro", + "missing-message sender hint should attribute to Saro" + ); + + // Draining clears the report; a reported gap is not surfaced again. + assert!(harness.raya().take_missing_messages().is_empty()); +} + +/// End-to-end acknowledgement detection on GroupV2. +/// +/// Saro sends a message; Raya and Pax reply. Each reply carries Saro's message +/// in its causal history, so Saro learns both peers hold it — without either +/// sending anything back on purpose. +#[test] +fn replies_acknowledge_the_message_they_were_sent_after() { + let mut harness = TestHarness::<3>::new(|_, _| {}); + + let participants = &[&harness.raya().addr(), &harness.pax().addr()]; + let convo_id = harness + .saro() + .create_group_convo_v2(participants, "", "") + .expect("saro create group"); + + harness.process_until_label("peers join", |h| { + h.raya().convo_count() == 1 && h.pax().convo_count() == 1 + }); + + let message_id = harness + .saro() + .send_content(&convo_id, b"anyone there?") + .expect("saro send"); + harness.process_until_label("peers get the message", |h| { + h.raya().check(&convo_id, b"anyone there?") && h.pax().check(&convo_id, b"anyone there?") + }); + assert!( + harness.saro().take_acks().is_empty(), + "holding a message is only observable once the peer sends" + ); + + // Each reply names Saro's message in its causal history. + harness + .raya() + .send_content(&convo_id, b"raya here") + .expect("raya reply"); + harness + .pax() + .send_content(&convo_id, b"pax here") + .expect("pax reply"); + harness.process_until_label("saro gets both replies", |h| { + h.saro().check(&convo_id, b"raya here") && h.saro().check(&convo_id, b"pax here") + }); + + let acks: Vec = harness.saro().take_acks(); + let mut holders: Vec<&str> = acks + .iter() + .filter(|a| a.conversation_id == convo_id && a.message_id == message_id) + .map(|a| a.acked_by.as_str()) + .collect(); + holders.sort_unstable(); + assert_eq!( + holders, + vec!["pax", "raya"], + "both peers that replied should be reported as holding the message" + ); + + // Draining clears the reports, and neither peer acknowledges twice. + assert!(harness.saro().take_acks().is_empty()); + harness + .raya() + .send_content(&convo_id, b"raya again") + .expect("raya second reply"); + harness.process_until_label("saro gets the second reply", |h| { + h.saro().check(&convo_id, b"raya again") + }); + assert!( + harness + .saro() + .take_acks() + .iter() + .all(|a| a.message_id != message_id), + "a peer acknowledges one message only once" + ); +} diff --git a/crates/generic-chat/src/client.rs b/crates/generic-chat/src/client.rs index f4e0df4..27006ad 100644 --- a/crates/generic-chat/src/client.rs +++ b/crates/generic-chat/src/client.rs @@ -6,8 +6,9 @@ use components::{ThreadedWakeupService, WakeupEvent}; use crossbeam_channel::{Receiver, Sender, select}; use crypto::Ed25519VerifyingKey; use libchat::{ - ConversationId, ConvoMetadata, ConvoOutcome, Core, DeliveryService, GroupV2Config, IdentId, - IdentIdRef, InboxOutcome, PayloadOutcome, RegistrationService, + ConversationId, ConvoMetadata, ConvoOutcome, Core, DeliveryAck, DeliveryService, GroupV2Config, + IdentId, IdentIdRef, InboxOutcome, MessageId, MissingMessage, PayloadOutcome, + RegistrationService, }; use logos_account::{AccountDirectory, resolve_device_ids}; use parking_lot::Mutex; @@ -275,7 +276,14 @@ where /// Encrypt and send `content` to an existing conversation. The core /// publishes the outbound envelope. - pub fn send_message(&mut self, convo_id: &str, content: &[u8]) -> Result<(), ClientError> { + /// + /// Returns the message's id, which later [`Event::MessageAcked`] events + /// carry — hold onto it to show which peers have the message. + pub fn send_message( + &mut self, + convo_id: &str, + content: &[u8], + ) -> Result { self.core .lock() .send_content(convo_id, content) @@ -347,7 +355,7 @@ fn worker_loop( }; let events = { let mut core = core.lock(); - match core.handle_payload(&bytes) { + let mut events = match core.handle_payload(&bytes) { Ok(outcome) => events_from_inbound(outcome, &directory), Err(e) => { tracing::warn!("inbound handle_payload failed: {e:?}"); @@ -355,7 +363,10 @@ fn worker_loop( message: e.to_string(), }] } - } + }; + events.extend(delivery_ack_events(core.take_acks(), &directory)); + events.extend(missing_events(core.take_missing_messages(), &directory)); + events }; for event in events { if event_tx.send(event).is_err() { @@ -368,12 +379,18 @@ fn worker_loop( return; // wakeup service's sender dropped }; // A wakeup can drive the steward's own commit, so it yields events too. - let events = match core.lock().wakeup(&convo_id) { - Ok(outcome) => events_from_inbound(outcome, &directory), - Err(e) => { - tracing::warn!("wakeup failed: {e:?}"); - Vec::new() - } + let events = { + let mut core = core.lock(); + let mut events = match core.wakeup(&convo_id) { + Ok(outcome) => events_from_inbound(outcome, &directory), + Err(e) => { + tracing::warn!("wakeup failed: {e:?}"); + Vec::new() + } + }; + events.extend(delivery_ack_events(core.take_acks(), &directory)); + events.extend(missing_events(core.take_missing_messages(), &directory)); + events }; for event in events { if event_tx.send(event).is_err() { @@ -398,6 +415,57 @@ fn events_from_inbound(result: PayloadOutcome, directory: &impl AccountDirectory } } +/// Map the acknowledgements the core observed while processing one payload onto +/// [`Event::MessageAcked`], one per peer per message. +/// +/// Drained from the same place as [`missing_events`]: the causal history of the +/// message just processed is what carried the acknowledgement. +fn delivery_ack_events(acks: Vec, directory: &impl AccountDirectory) -> Vec { + acks.into_iter() + .map(|a| Event::MessageAcked { + convo_id: Arc::from(a.conversation_id), + message_id: a.message_id, + acked_by: sender_hint(directory, &a.acked_by), + }) + .collect() +} + +/// Map the causal-history gaps the core detected while processing one payload +/// onto [`Event::MessageMissing`]. +/// +/// Drained right after each drive of the core, so a gap arrives with the batch +/// of events for the message that revealed it — and after them, so a gap on a +/// conversation this payload just started still follows its +/// [`Event::ConversationStarted`]. +fn missing_events(missing: Vec, directory: &impl AccountDirectory) -> Vec { + missing + .into_iter() + .map(|m| Event::MessageMissing { + convo_id: Arc::from(m.conversation_id), + message_id: m.frontier.message_id().to_owned(), + sender_hint: sender_hint(directory, m.frontier.sender_id()), + }) + .collect() +} + +/// Resolve a participant a causal-history observation named — the author of a +/// message we never saw, or the peer acknowledging one of ours. +/// +/// Same credential decoding as a delivered message's sender, but the claim is +/// self-asserted rather than authenticated, so an unconfirmable account yields +/// the device alone rather than dropping the observation. `None` when the value +/// is not a credential at all. +fn sender_hint(directory: &impl AccountDirectory, encoded: &str) -> Option { + let (device, claim) = parse_credential(directory, encoded.as_bytes()).ok()?; + Some(MessageSender { + account: match claim { + AccountClaim::Verified(account) => Some(account), + AccountClaim::None | AccountClaim::Unverified(_) => None, + }, + local_identity: device, + }) +} + /// Interpret a hex account address as an Ed25519 account verifying key. fn account_key_from_hex(addr: &str) -> Option { let bytes: [u8; 32] = hex::decode(addr).ok()?.try_into().ok()?; @@ -601,10 +669,11 @@ mod sender_check_tests { use logos_account::{DeviceSet, SignedDeviceBundle}; use super::{ - GroupMember, MessageSender, SenderError, decode_sender, dedup_members, member_key, - roster_member, + Event, GroupMember, MessageSender, SenderError, decode_sender, dedup_members, + delivery_ack_events, member_key, missing_events, roster_member, }; use crate::delegate::DelegateCredential; + use libchat::{DeliveryAck, Frontier, MissingMessage}; /// In-test account → device directory. Holds device id sets keyed by the hex /// account key, and can be made to fail to simulate a directory outage. @@ -913,4 +982,112 @@ mod sender_check_tests { vec![committed] ); } + + /// A gap reported by the causal history, as the core hands it over: the + /// sender hint travels in the same encoding a message's credential does. + fn gap(sender_hint: &str) -> MissingMessage { + MissingMessage { + conversation_id: "convo".to_owned(), + frontier: Frontier::new(sender_hint.to_owned(), "msg-id".to_owned()), + } + } + + fn hex_cred(cred: DelegateCredential) -> String { + hex::encode(cred.serialize()) + } + + /// Unwrap the single `MessageMissing` a one-gap batch produces. + fn only_missing(events: Vec) -> (String, Option) { + match <[Event; 1]>::try_from(events) + .expect("one gap produces one event") + .into_iter() + .next() + .unwrap() + { + Event::MessageMissing { + convo_id, + message_id, + sender_hint, + } => { + assert_eq!(&*convo_id, "convo"); + (message_id, sender_hint) + } + other => panic!("expected MessageMissing, got {other:?}"), + } + } + + /// An account claim the directory contradicts drops a *delivered* message, + /// but a gap is still worth reporting: the hint keeps the device and + /// forgoes the account, since nothing about an unseen message is verifiable + /// anyway. + #[test] + fn missing_message_hint_keeps_the_device_when_the_account_claim_fails() { + let account = key(); + let endorsed = key(); + let spoofer = key(); + let dir = FakeDir::with_devices(&account, &[&endorsed]); + let cred = DelegateCredential::associated(&spoofer, &hex::encode(account.as_ref())); + + let (_, sender) = only_missing(missing_events(vec![gap(&hex_cred(cred))], &dir)); + assert_eq!( + sender, + Some(MessageSender { + account: None, + local_identity: local_id(&spoofer), + }) + ); + } + + /// A hint that is not a credential at all still reports the gap — the + /// message id is the part the application needs. + #[test] + fn missing_message_without_a_resolvable_hint_is_still_reported() { + let (message_id, sender) = + only_missing(missing_events(vec![gap("saro")], &FakeDir::default())); + assert_eq!(message_id, "msg-id"); + assert_eq!(sender, None); + } + + /// One acknowledgement per peer per message, each naming the peer an + /// application would list against the message. + #[test] + fn acks_name_the_peers_that_hold_the_message() { + let account = key(); + let device = key(); + let dir = FakeDir::with_devices(&account, &[&device]); + let peer = DelegateCredential::associated(&device, &hex::encode(account.as_ref())); + + let events = delivery_ack_events( + vec![DeliveryAck { + conversation_id: "convo".to_owned(), + message_id: "msg-id".to_owned(), + acked_by: hex_cred(peer), + }], + &dir, + ); + + match <[Event; 1]>::try_from(events) + .expect("one ack produces one event") + .into_iter() + .next() + .unwrap() + { + Event::MessageAcked { + convo_id, + message_id, + acked_by, + } => { + assert_eq!(&*convo_id, "convo"); + assert_eq!(message_id, "msg-id"); + assert_eq!( + acked_by, + Some(MessageSender { + account: Some(local_id(&account)), + local_identity: local_id(&device), + }) + ); + } + other => panic!("expected MessageAcked, got {other:?}"), + } + } } diff --git a/crates/generic-chat/src/event.rs b/crates/generic-chat/src/event.rs index 4618e02..5e2e4a8 100644 --- a/crates/generic-chat/src/event.rs +++ b/crates/generic-chat/src/event.rs @@ -38,6 +38,36 @@ pub enum Event { content: Vec, sender: MessageSender, }, + /// A peer acknowledged a message this client sent: it referenced that + /// message in the causal history of a message of its own, so it held ours + /// when it sent. `message_id` is the id the send returned. + /// + /// Evidence of delivery to the peer's client, not of a human reading it. + /// The acknowledgement is passive — nothing is sent back on purpose — so a + /// peer that never sends never acknowledges, and an application should + /// treat the absence of one as "not confirmed" rather than "not delivered". + /// + /// `acked_by` is resolved from the peer's self-asserted `sender_id` and is + /// **not authenticated**; see [`Self::MessageMissing`]'s `sender_hint`. + /// `None` when it could not be resolved to a device. + MessageAcked { + convo_id: Arc, + message_id: String, + acked_by: Option, + }, + /// A message this client never received, revealed by the causal history of + /// one that did arrive. Detection only — nothing is fetched or replayed, + /// and the gap is reported once. + /// + /// `sender_hint` is the author the *referencing* peer named, resolved the + /// same way as [`Self::MessageReceived`]'s sender but **not authenticated**: + /// nothing about a message we never saw can be verified, so treat it as a + /// display hint. `None` when the hint could not be resolved to a device. + MessageMissing { + convo_id: Arc, + message_id: String, + sender_hint: Option, + }, /// A commit changed a conversation's membership. ConversationMembersChanged { convo_id: Arc, diff --git a/crates/generic-chat/src/lib.rs b/crates/generic-chat/src/lib.rs index 91d1318..7156ca7 100644 --- a/crates/generic-chat/src/lib.rs +++ b/crates/generic-chat/src/lib.rs @@ -15,7 +15,8 @@ pub use event::{Event, MessageSender}; // Re-export types callers need to interact with ChatClient. pub use libchat::{ AddressedEnvelope, ChatStore, ConversationClass, ConversationId, ConvoMetadata, - DeliveryService, GroupV2Config, IdentityProvider, RegistrationService, StorageConfig, + DeliveryService, GroupV2Config, IdentityProvider, MessageId, RegistrationService, + StorageConfig, }; // The directory trait bounds ChatClient's registry parameter, so callers // writing code generic over ChatClient need it too. diff --git a/crates/generic-chat/tests/group_v2.rs b/crates/generic-chat/tests/group_v2.rs index 9ace11e..447f70a 100644 --- a/crates/generic-chat/tests/group_v2.rs +++ b/crates/generic-chat/tests/group_v2.rs @@ -479,3 +479,63 @@ fn group_metadata_defaults_to_empty() { assert_eq!(meta.name, ""); assert_eq!(meta.desc, ""); } + +/// The peers that hold a sent message surface as `MessageAcked` events keyed by +/// the id the send returned — what an application needs to list the peers that +/// hold a message. The acknowledgement is passive: Raya and Pax only send +/// ordinary replies, never a receipt. +#[test] +fn a_sent_message_is_acknowledged_by_the_peers_that_reply() { + let bus = MessageBus::default(); + let reg = EphemeralRegistry::new(); + + let (mut saro, saro_events, saro_addr) = create_test_client(bus.clone(), reg.clone()); + let (mut raya, raya_events, raya_addr) = create_test_client(bus.clone(), reg.clone()); + let (mut pax, pax_events, pax_addr) = create_test_client(bus.clone(), reg.clone()); + + let convo_id = saro + .create_group_conversation(&[&raya_addr, &pax_addr], unnamed_group()) + .expect("saro create group"); + wait_for_group_started(&raya_events, "raya ConversationStarted"); + wait_for_group_started(&pax_events, "pax ConversationStarted"); + wait_for_members(&mut saro, &convo_id, &[&saro_addr, &raya_addr, &pax_addr]); + + let message_id = saro + .send_message(&convo_id, b"anyone there?") + .expect("saro send"); + wait_for_message(&raya_events, b"anyone there?"); + wait_for_message(&pax_events, b"anyone there?"); + + // Ordinary replies; their causal history carries the acknowledgement. + raya.send_message(&convo_id, b"raya here") + .expect("raya reply"); + pax.send_message(&convo_id, b"pax here").expect("pax reply"); + + let mut holders = Vec::new(); + while holders.len() < 2 { + let peer = wait_for_event( + &saro_events, + "saro MessageAcked", + Duration::from_secs(10), + |e| match e { + Event::MessageAcked { + convo_id: id, + message_id: acked, + acked_by, + } if **id == *convo_id && *acked == message_id => Some( + acked_by + .as_ref() + .and_then(|a| a.account.as_ref()) + .map(|a| a.as_str().to_string()), + ), + _ => None, + }, + ); + holders.push(peer.expect("the acknowledging peer's account should be directory-verified")); + } + holders.sort(); + + let mut expected = vec![raya_addr.clone(), pax_addr.clone()]; + expected.sort(); + assert_eq!(holders, expected, "both replying peers should be listed"); +} diff --git a/crates/generic-chat/tests/saro_and_raya.rs b/crates/generic-chat/tests/saro_and_raya.rs index 4c7f0d5..09f3ae7 100644 --- a/crates/generic-chat/tests/saro_and_raya.rs +++ b/crates/generic-chat/tests/saro_and_raya.rs @@ -54,6 +54,25 @@ where f(event).unwrap_or_else(|other| panic!("expected {label}, got {other:?}")) } +/// [`expect_event`] for a back-and-forth exchange, skipping acknowledgements. +/// +/// Each reply acknowledges the message it was sent after, so `MessageAcked` +/// lands at points a test driving one direction at a time does not control. +fn expect_event_ignoring_acks(events: &Receiver, label: &str, mut f: F) -> T +where + F: FnMut(Event) -> Result, +{ + loop { + let event = events + .recv_timeout(Duration::from_secs(5)) + .unwrap_or_else(|_| panic!("timed out waiting for {label}")); + if matches!(event, Event::MessageAcked { .. }) { + continue; + } + return f(event).unwrap_or_else(|other| panic!("expected {label}, got {other:?}")); + } +} + #[test] fn direct_v1_integration() { let bus = MessageBus::default(); @@ -255,7 +274,7 @@ fn saro_raya_message_exchange() { for i in 0u8..5 { let msg = format!("msg {i}"); saro.send_message(&saro_convo_id, msg.as_bytes()).unwrap(); - expect_event( + expect_event_ignoring_acks( &raya_events, &format!("MessageReceived(msg {i})"), |e| match e { @@ -269,7 +288,7 @@ fn saro_raya_message_exchange() { let reply = format!("reply {i}"); raya.send_message(&raya_convo_id, reply.as_bytes()).unwrap(); - expect_event( + expect_event_ignoring_acks( &saro_events, &format!("MessageReceived(reply {i})"), |e| match e {