use mail4agent_api::{
Ack, Address, Directory, DirectoryEntry, InboxPage, MailError, Message, MessageId, Participant,
ParticipantId, RoomEntry, RoomId, SendRequest, SendResponse, SessionCard, SessionEntry, SessionId,
UnreadCount, MESSAGE_ID_HEX_LEN, MESSAGE_ID_PREFIX,
};
use sha2::{Digest, Sha256};
use subtle::ConstantTimeEq;
use crate::store::{InsertMessageOutcome, MailStore, ParticipantRecord, SecretDigest, SessionRecord, StoreError};
pub const SECRET_HEX_LEN: usize = 64;
#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
pub struct ParticipantPermissions {
pub may_send: bool,
pub may_read: bool,
pub operator: bool,
}
pub type LivenessCheck<'a> = &'a dyn Fn(u32, u64) -> bool;
pub struct MailboxEngine<S> {
store: S,
}
impl<S: MailStore> MailboxEngine<S> {
pub fn new(store: S) -> Self {
Self { store }
}
pub fn authenticate(&self, presented_secret: &str) -> Result<ParticipantId, MailError> {
let digest = sha256_digest(presented_secret.as_bytes());
let found = self
.store
.find_participant_by_digest(&digest)
.map_err(|err| store_unavailable("find_participant_by_digest", err))?;
let Some((id, record)) = found else {
return Err(MailError::PermissionDenied { need: "mail:authenticate".to_string() });
};
if bool::from(record.secret_digest.ct_eq(&digest)) {
Ok(id)
} else {
Err(MailError::PermissionDenied { need: "mail:authenticate".to_string() })
}
}
pub fn register_participant(
&mut self,
id: ParticipantId,
label: Option<String>,
permissions: ParticipantPermissions,
) -> Result<String, MailError> {
Participant { id: id.clone(), label: label.clone() }.validate()?;
let existing = self.store.get_participant(&id).map_err(|err| store_unavailable("get_participant", err))?;
if existing.is_some() {
return Err(MailError::Malformed {
field: "participant".to_string(),
reason: format!("participant \"{id}\" is already registered"),
});
}
let secret = generate_secret();
let record = ParticipantRecord {
label,
secret_digest: sha256_digest(secret.as_bytes()),
may_send: permissions.may_send,
may_read: permissions.may_read,
operator: permissions.operator,
listener_url: None,
};
self.store.register_participant(id, record).map_err(|err| store_unavailable("register_participant", err))?;
Ok(secret)
}
pub fn deregister_participant(&mut self, id: &ParticipantId) -> Result<(), MailError> {
self.require_participant(id)?;
self.store.deregister_participant(id).map_err(|err| store_unavailable("deregister_participant", err))
}
pub fn rotate_participant_secret(&mut self, id: &ParticipantId) -> Result<String, MailError> {
self.require_participant(id)?;
let secret = generate_secret();
self.store
.set_participant_secret_digest(id, sha256_digest(secret.as_bytes()))
.map_err(|err| store_unavailable("set_participant_secret_digest", err))?;
Ok(secret)
}
pub fn revoke_participant_secret(&mut self, id: &ParticipantId) -> Result<(), MailError> {
self.rotate_participant_secret(id).map(|_secret| ())
}
pub fn set_listener(&mut self, id: &ParticipantId, url: String) -> Result<(), MailError> {
validate_listener_url(&url)?;
self.require_participant(id)?;
self.store.set_listener_url(id, Some(url)).map_err(|err| store_unavailable("set_listener_url", err))
}
pub fn remove_listener(&mut self, id: &ParticipantId) -> Result<(), MailError> {
self.require_participant(id)?;
self.store.set_listener_url(id, None).map_err(|err| store_unavailable("set_listener_url", err))
}
pub fn create_room(&mut self, id: RoomId, now_unix_ms: u64) -> Result<(), MailError> {
let existing = self.store.get_room(&id).map_err(|err| store_unavailable("get_room", err))?;
if existing.is_some() {
return Err(MailError::Malformed {
field: "room".to_string(),
reason: format!("room \"{id}\" already exists"),
});
}
self.store.create_room(id, now_unix_ms).map_err(|err| store_unavailable("create_room", err))
}
pub fn add_room_member(&mut self, room: &RoomId, participant: ParticipantId) -> Result<(), MailError> {
let room_record = self.store.get_room(room).map_err(|err| store_unavailable("get_room", err))?;
if room_record.is_none() {
return Err(MailError::UnknownRoom { room: room.clone() });
}
let participant_record =
self.store.get_participant(&participant).map_err(|err| store_unavailable("get_participant", err))?;
if participant_record.is_none() {
return Err(MailError::UnknownParticipant { participant });
}
self.store.add_room_member(room, participant).map_err(|err| store_unavailable("add_room_member", err))
}
pub fn remove_room_member(&mut self, room: &RoomId, participant: &ParticipantId) -> Result<(), MailError> {
let room_record = self.store.get_room(room).map_err(|err| store_unavailable("get_room", err))?;
if room_record.is_none() {
return Err(MailError::UnknownRoom { room: room.clone() });
}
self.store.remove_room_member(room, participant).map_err(|err| store_unavailable("remove_room_member", err))
}
pub fn ensure_session(
&mut self,
account: ParticipantId,
session_id: SessionId,
card: SessionCard,
now_unix_ms: u64,
) -> Result<SessionId, MailError> {
card.validate()?;
session_id.validate()?;
self.require_participant(&account)?;
let existing = self.store.get_session(&session_id).map_err(|err| store_unavailable("get_session", err))?;
let declared = match &existing {
Some(record) if record.account == account => record.card.declared.clone(),
Some(record) => {
return Err(MailError::SessionAccountMismatch {
session: session_id,
expected: record.account.clone(),
presented: account,
});
}
None => Default::default(),
};
let merged = SessionCard { attested: card.attested, corroborated: card.corroborated, declared };
merged.validate()?;
let record = SessionRecord { account, card: merged, last_seen_unix_ms: now_unix_ms };
self.store.upsert_session(session_id.clone(), record).map_err(|err| store_unavailable("upsert_session", err))?;
Ok(session_id)
}
pub fn set_declared(
&mut self,
session: &SessionId,
working_on: Option<String>,
role: Option<String>,
parent: Option<SessionId>,
) -> Result<(), MailError> {
session.validate()?;
let mut record = self
.store
.get_session(session)
.map_err(|err| store_unavailable("get_session", err))?
.ok_or_else(|| MailError::UnknownSession { session: session.clone() })?;
let declared = mail4agent_api::SessionDeclared { working_on, role, parent };
declared.validate()?;
record.card.declared = declared;
self.store.upsert_session(session.clone(), record).map_err(|err| store_unavailable("upsert_session", err))
}
pub fn send(&mut self, sender: &Address, request: SendRequest, now_unix_ms: u64) -> Result<SendResponse, MailError> {
request.validate()?;
let (_, sender_record) = self.resolve_identity(sender)?;
if !sender_record.may_send {
return Err(MailError::PermissionDenied { need: "mail:send".to_string() });
}
match &request.to {
Address::Direct { participant } => {
let exists =
self.store.get_participant(participant).map_err(|err| store_unavailable("get_participant", err))?;
if exists.is_none() {
return Err(MailError::UnknownParticipant { participant: participant.clone() });
}
}
Address::Session { participant, session } => {
let session_record =
self.store.get_session(session).map_err(|err| store_unavailable("get_session", err))?;
match session_record {
Some(record) if &record.account == participant => {}
Some(record) => {
return Err(MailError::SessionAccountMismatch {
session: session.clone(),
expected: record.account,
presented: participant.clone(),
});
}
None => return Err(MailError::UnknownSession { session: session.clone() }),
}
}
Address::Room { room } => {
let exists = self.store.get_room(room).map_err(|err| store_unavailable("get_room", err))?;
if exists.is_none() {
return Err(MailError::UnknownRoom { room: room.clone() });
}
}
}
let idempotency = request.idempotency_key.clone().map(|key| (sender.clone(), key));
let message_id = derive_message_id(sender, &request, now_unix_ms);
let message = Message {
message_id: message_id.clone(),
from: sender.clone(),
to: request.to,
subject: request.subject,
body: request.body,
reply_to: request.reply_to,
correlation: request.correlation,
refs: request.refs,
created_at_unix_ms: now_unix_ms,
};
message.validate()?;
let outcome =
self.store.insert_message(message, idempotency).map_err(|err| store_unavailable("insert_message", err))?;
let message_id = match outcome {
InsertMessageOutcome::Inserted => message_id,
InsertMessageOutcome::Deduplicated { message_id } => message_id,
};
Ok(SendResponse { message_id, from: sender.clone() })
}
pub fn inbox(&self, reader: &Address, since_unix_ms: u64, limit: u16) -> Result<InboxPage, MailError> {
let (account, record) = self.resolve_identity(reader)?;
self.require_read_permission(&record)?;
self.build_inbox(reader, &account, since_unix_ms, limit)
}
pub fn inbox_of(
&self,
caller: &Address,
target: &Address,
since_unix_ms: u64,
limit: u16,
) -> Result<InboxPage, MailError> {
let (_, caller_record) = self.resolve_identity(caller)?;
if caller != target && !caller_record.operator {
return Err(MailError::PermissionDenied { need: "mail:operator".to_string() });
}
let (target_account, target_record) = self.resolve_identity(target)?;
if caller == target {
self.require_read_permission(&target_record)?;
}
self.build_inbox(target, &target_account, since_unix_ms, limit)
}
pub fn ack(&mut self, reader: &Address, message_id: &MessageId, now_unix_ms: u64) -> Result<Ack, MailError> {
let (account, record) = self.resolve_identity(reader)?;
self.require_read_permission(&record)?;
let message = self
.store
.get_message(message_id)
.map_err(|err| store_unavailable("get_message", err))?
.ok_or_else(|| MailError::UnknownMessage { message_id: message_id.clone() })?;
if !self.is_readable(&message, reader, &account, &record)? {
return Err(MailError::NotAddressedToYou { message_id: message_id.clone() });
}
let ack = Ack { message_id: message_id.clone(), reader: reader.clone(), acked_at_unix_ms: now_unix_ms };
ack.validate()?;
self.store.record_ack(ack).map_err(|err| store_unavailable("record_ack", err))
}
pub fn message_get(&self, reader: &Address, message_id: &MessageId) -> Result<Message, MailError> {
let (account, record) = self.resolve_identity(reader)?;
self.require_read_permission(&record)?;
let message = self
.store
.get_message(message_id)
.map_err(|err| store_unavailable("get_message", err))?
.ok_or_else(|| MailError::UnknownMessage { message_id: message_id.clone() })?;
if !self.is_readable(&message, reader, &account, &record)? {
return Err(MailError::NotAddressedToYou { message_id: message_id.clone() });
}
Ok(message)
}
pub fn unread_count_of(&self, caller: &Address, target: &Address) -> Result<UnreadCount, MailError> {
let (_, caller_record) = self.resolve_identity(caller)?;
if caller != target && !caller_record.operator {
return Err(MailError::PermissionDenied { need: "mail:operator".to_string() });
}
let (target_account, target_record) = self.resolve_identity(target)?;
if caller == target {
self.require_read_permission(&target_record)?;
}
let unread = self.count_unread(target, &target_account)?;
Ok(UnreadCount { target: target.clone(), unread })
}
pub fn directory(&self, caller: &Address, is_alive: LivenessCheck<'_>) -> Result<Directory, MailError> {
let (account, record) = self.resolve_identity(caller)?;
self.require_read_permission(&record)?;
let participants = self.store.list_participants().map_err(|err| store_unavailable("list_participants", err))?;
let mut entries = Vec::with_capacity(participants.len());
for summary in participants {
let sessions = self
.store
.sessions_of(&summary.id)
.map_err(|err| store_unavailable("sessions_of", err))?
.into_iter()
.map(|(id, record)| SessionEntry {
live: is_alive(record.card.attested.pid, record.card.attested.started_at_unix_ms),
last_seen_unix_ms: record.last_seen_unix_ms,
card: record.card,
id,
})
.collect();
entries.push(DirectoryEntry { id: summary.id, label: summary.label, sessions });
}
let rooms = self
.store
.list_rooms()
.map_err(|err| store_unavailable("list_rooms", err))?
.into_iter()
.map(|summary| RoomEntry { member: summary.members.contains(&account), id: summary.id })
.collect();
Ok(Directory { participants: entries, rooms })
}
fn require_participant(&self, id: &ParticipantId) -> Result<ParticipantRecord, MailError> {
self.store
.get_participant(id)
.map_err(|err| store_unavailable("get_participant", err))?
.ok_or_else(|| MailError::UnknownParticipant { participant: id.clone() })
}
fn resolve_identity(&self, address: &Address) -> Result<(ParticipantId, ParticipantRecord), MailError> {
address.validate()?;
match address {
Address::Direct { participant } => {
let record = self.require_participant(participant)?;
Ok((participant.clone(), record))
}
Address::Session { participant, session } => {
let session_record = self
.store
.get_session(session)
.map_err(|err| store_unavailable("get_session", err))?
.ok_or_else(|| MailError::UnknownSession { session: session.clone() })?;
if &session_record.account != participant {
return Err(MailError::SessionAccountMismatch {
session: session.clone(),
expected: session_record.account,
presented: participant.clone(),
});
}
let record = self.require_participant(participant)?;
Ok((participant.clone(), record))
}
Address::Room { room } => Err(MailError::Malformed {
field: "address".to_string(),
reason: format!("a room (\"{room}\") cannot act as a participant identity"),
}),
}
}
fn require_read_permission(&self, record: &ParticipantRecord) -> Result<(), MailError> {
if record.may_read || record.operator {
Ok(())
} else {
Err(MailError::PermissionDenied { need: "mail:read".to_string() })
}
}
fn is_readable(
&self,
message: &Message,
reader: &Address,
account: &ParticipantId,
record: &ParticipantRecord,
) -> Result<bool, MailError> {
if record.operator {
return Ok(true);
}
match &message.to {
Address::Direct { participant } => Ok(participant == account),
Address::Session { .. } => Ok(reader == &message.to),
Address::Room { room } => {
let room_record = self.store.get_room(room).map_err(|err| store_unavailable("get_room", err))?;
Ok(room_record.is_some_and(|room_record| room_record.members.contains(account)))
}
}
}
fn own_messages(&self, identity: &Address, account: &ParticipantId, since_unix_ms: u64) -> Result<Vec<Message>, MailError> {
let mut messages = self
.store
.messages_to_since(identity, since_unix_ms)
.map_err(|err| store_unavailable("messages_to_since", err))?;
if matches!(identity, Address::Session { .. }) {
let account_address = Address::Direct { participant: account.clone() };
let account_messages = self
.store
.messages_to_since(&account_address, since_unix_ms)
.map_err(|err| store_unavailable("messages_to_since", err))?;
messages.extend(account_messages);
}
let rooms = self.store.rooms_containing(account).map_err(|err| store_unavailable("rooms_containing", err))?;
for room in rooms {
let room_messages = self
.store
.room_messages_since(&room, since_unix_ms)
.map_err(|err| store_unavailable("room_messages_since", err))?;
messages.extend(room_messages);
}
Ok(messages)
}
fn build_inbox(
&self,
identity: &Address,
account: &ParticipantId,
since_unix_ms: u64,
limit: u16,
) -> Result<InboxPage, MailError> {
let mut messages = self.own_messages(identity, account, since_unix_ms)?;
messages.sort_by(|a, b| a.created_at_unix_ms.cmp(&b.created_at_unix_ms).then_with(|| a.message_id.cmp(&b.message_id)));
messages.truncate(usize::from(limit));
let unread = self.count_unread(identity, account)?;
Ok(InboxPage { messages, unread })
}
fn count_unread(&self, identity: &Address, account: &ParticipantId) -> Result<u32, MailError> {
let messages = self.own_messages(identity, account, 0)?;
let mut unread = 0u32;
for message in messages {
if &message.from == identity {
continue;
}
let ack = self
.store
.get_ack(&message.message_id, identity)
.map_err(|err| store_unavailable("get_ack", err))?;
if ack.is_none() {
unread += 1;
}
}
Ok(unread)
}
}
fn store_unavailable(operation: &'static str, err: StoreError) -> MailError {
tracing::error!(operation, error = %err, "mail store operation failed");
MailError::StoreUnavailable { operation: operation.to_string() }
}
const LISTENER_URL_MAX_BYTES: usize = 512;
fn validate_listener_url(url: &str) -> Result<(), MailError> {
let malformed = |reason: &str| MailError::Malformed { field: "url".to_string(), reason: reason.to_string() };
if url.len() > LISTENER_URL_MAX_BYTES {
return Err(MailError::TooLarge {
field: "url".to_string(),
limit: LISTENER_URL_MAX_BYTES,
actual: url.len(),
});
}
if url.chars().any(char::is_control) {
return Err(malformed("must not contain control characters"));
}
const LOOPBACK_REFUSAL: &str = "must be an http://127.0.0.1:* or http://localhost:* URL -- this mailbox is \
a local service and never turns a message into an outbound call anywhere else";
let Some(after_scheme) = url.strip_prefix("http://") else {
return Err(malformed(LOOPBACK_REFUSAL));
};
let authority = after_scheme.split(['/', '?', '#']).next().unwrap_or("");
if authority.contains('@') {
return Err(malformed("must not carry userinfo (\"user:pass@\") in a loopback listener URL"));
}
let host = authority.split(':').next().unwrap_or("");
if !host.eq_ignore_ascii_case("127.0.0.1") && !host.eq_ignore_ascii_case("localhost") {
return Err(malformed(LOOPBACK_REFUSAL));
}
Ok(())
}
fn generate_secret() -> String {
let mut bytes = [0u8; SECRET_HEX_LEN / 2];
getrandom::getrandom(&mut bytes).expect("OS random source unavailable: cannot mint a participant secret without it");
hex::encode(bytes)
}
fn sha256_digest(data: &[u8]) -> SecretDigest {
let mut hasher = Sha256::new();
hasher.update(data);
let mut digest = [0u8; 32];
digest.copy_from_slice(&hasher.finalize());
digest
}
fn derive_message_id(sender: &Address, request: &SendRequest, now_unix_ms: u64) -> MessageId {
let mut hasher = Sha256::new();
hasher.update(sender.to_string().as_bytes());
hasher.update([0u8]);
match &request.to {
Address::Direct { participant } => {
hasher.update(b"direct:");
hasher.update(participant.as_str().as_bytes());
}
Address::Session { participant, session } => {
hasher.update(b"session:");
hasher.update(participant.as_str().as_bytes());
hasher.update([0u8]);
hasher.update(session.as_str().as_bytes());
}
Address::Room { room } => {
hasher.update(b"room:");
hasher.update(room.as_str().as_bytes());
}
}
hasher.update([0u8]);
hasher.update(request.subject.as_bytes());
hasher.update([0u8]);
hasher.update(request.body.as_bytes());
hasher.update([0u8]);
if let Some(reply_to) = &request.reply_to {
hasher.update(reply_to.as_str().as_bytes());
}
hasher.update([0u8]);
if let Some(correlation) = &request.correlation {
hasher.update(correlation.as_bytes());
}
hasher.update([0u8]);
for reference in &request.refs {
hasher.update(reference.kind.as_bytes());
hasher.update([0u8]);
hasher.update(reference.locator.as_bytes());
hasher.update([0u8]);
if let Some(digest) = &reference.digest {
hasher.update(digest.as_bytes());
}
hasher.update([0u8]);
}
hasher.update(now_unix_ms.to_le_bytes());
hasher.update([0u8]);
match &request.idempotency_key {
Some(key) => hasher.update(key.as_bytes()),
None => {
let mut nonce = [0u8; 16];
getrandom::getrandom(&mut nonce).expect("OS random source unavailable: cannot mint a message id without it");
hasher.update(nonce);
}
}
let digest = hasher.finalize();
let hex_digest = hex::encode(digest);
let body = &hex_digest[..MESSAGE_ID_HEX_LEN];
MessageId::new(format!("{MESSAGE_ID_PREFIX}{body}")).expect("derived message id always matches MessageId's own shape")
}