From daa486f7e7faecb4821cc214bf661bb426aa9fba Mon Sep 17 00:00:00 2001 From: John Coffey Date: Mon, 28 Sep 2026 21:45:09 -0700 Subject: [PATCH] Journaling: search, read and export over JMAP, and the chain check Phase 4 of the journaling spec. - inbuxa:JournalEntry/query and /get (sysJournalSearch): filter by time, sender, recipient, either, direction, subject words, Message-ID and journal, newest first; the whole report only when asked for. - inbuxa:JournalExport/set (sysJournalExport): a reason is required; a ZIP of the matching reports with manifest.csv, exceptions.csv and manifest.sha256, up to 10,000 reports and 1 GB. - inbuxa:JournalVerification/set (sysJournalGet): chains and reports rechecked. - Every search, listing, read, export and check is written to the audit log before anything is returned, with existing actions only. - Catalog entries for the three objects; spec as-built notes. journal_tests: administrators can't search; a Compliance Officer searches, lists, reads a report, exports (reason required) and checks the chain; the officer can't change journals; each of those is in the audit log. --- crates/features/src/journal/entries.rs | 180 ++++ .../src/object/inbuxa_journal_entry.rs | 340 ++++++++ crates/jmap-proto/src/object/mod.rs | 1 + crates/jmap-proto/src/references/eval.rs | 3 + crates/jmap-proto/src/references/resolve.rs | 7 + crates/jmap-proto/src/request/method.rs | 21 +- crates/jmap-proto/src/request/mod.rs | 4 + crates/jmap-proto/src/request/parser.rs | 28 + crates/jmap-proto/src/response/mod.rs | 21 + crates/jmap/src/api/auth.rs | 20 + crates/jmap/src/api/request.rs | 32 + crates/jmap/src/changes/get.rs | 3 + crates/jmap/src/inbuxa/journal_entry.rs | 791 ++++++++++++++++++ crates/jmap/src/inbuxa/mod.rs | 1 + docs/spec/features/journaling.md | 21 + resources/privacy/catalog.toml | 30 + tests/src/system/journal.rs | 224 ++++- 17 files changed, 1725 insertions(+), 2 deletions(-) create mode 100644 crates/jmap-proto/src/object/inbuxa_journal_entry.rs create mode 100644 crates/jmap/src/inbuxa/journal_entry.rs diff --git a/crates/features/src/journal/entries.rs b/crates/features/src/journal/entries.rs index ebdf3c9..2327768 100644 --- a/crates/features/src/journal/entries.rs +++ b/crates/features/src/journal/entries.rs @@ -358,6 +358,120 @@ pub async fn list( Ok(out) } +/// Most results one search page returns. +pub const MAX_QUERY_LIMIT: usize = 500; + +/// A search of the journal (JR-15): conditions that must all hold. +#[derive(Debug, Clone, Default, PartialEq, Eq, SerdeSerialize)] +#[serde(rename_all = "camelCase")] +pub struct Filter { + /// From this time on, in seconds. + #[serde(skip_serializing_if = "Option::is_none")] + pub after: Option, + /// Before this time, in seconds. + #[serde(skip_serializing_if = "Option::is_none")] + pub before: Option, + /// Part of the sender's address, ignoring case. + #[serde(skip_serializing_if = "Option::is_none")] + pub sender: Option, + /// Part of any recipient's address, ignoring case. + #[serde(skip_serializing_if = "Option::is_none")] + pub recipient: Option, + /// Part of the sender's or any recipient's address. + #[serde(skip_serializing_if = "Option::is_none")] + pub address: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub direction: Option, + /// Words that must all appear in the subject, ignoring case. + #[serde(skip_serializing_if = "Option::is_none")] + pub text: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub message_id: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub journal_id: Option, +} + +impl Filter { + pub fn matches(&self, entry: &Entry) -> bool { + let has = |value: &str, part: &str| value.to_lowercase().contains(&part.to_lowercase()); + self.after.is_none_or(|after| entry.at >= after) + && self.before.is_none_or(|before| entry.at < before) + && self.sender.as_deref().is_none_or(|s| has(&entry.sender, s)) + && self + .recipient + .as_deref() + .is_none_or(|r| entry.recipients.iter().any(|a| has(a, r))) + && self + .address + .as_deref() + .is_none_or(|a| has(&entry.sender, a) || entry.recipients.iter().any(|r| has(r, a))) + && self + .direction + .is_none_or(|d| d == Direction::Any || d == entry.direction) + && self.text.as_deref().is_none_or(|text| { + let subject = entry.subject.to_lowercase(); + text.to_lowercase() + .split_whitespace() + .all(|word| subject.contains(word)) + }) + && self.message_id.as_deref().is_none_or(|id| { + entry.message_id.trim_matches(['<', '>']) == id.trim_matches(['<', '>']) + }) + && self.journal_id.is_none_or(|j| entry.journals.contains(&j)) + } +} + +/// Entries matching `filter`, newest first: a page from `position`, up to +/// `limit`, and, when asked, how many match in all. +pub async fn query( + data: &Store, + filter: &Filter, + position: usize, + limit: usize, + count_all: bool, +) -> trc::Result<(Vec, usize)> { + let after = filter.after.unwrap_or(0); + let before = filter.before.unwrap_or(u64::MAX); + let mut ids = Vec::new(); + data.iterate( + IterateParams::new( + key(KIND_TIME, &[after, 0, 0]), + key(KIND_TIME, &[before.saturating_sub(1), u64::MAX, u64::MAX]), + ) + .descending() + .no_values(), + |key, _| { + if let Some(parts) = parse_key(key, KIND_TIME, 3) { + ids.push(EntryId { + node: parts[1], + seq: parts[2], + }); + } + Ok(true) + }, + ) + .await + .caused_by(trc::location!())?; + let mut page = Vec::new(); + let mut total = 0; + for id in ids { + let Some(entry) = get(data, id).await? else { + continue; + }; + if !filter.matches(&entry) { + continue; + } + if total >= position && page.len() < limit { + page.push(id); + } + total += 1; + if !count_all && page.len() >= limit { + break; + } + } + Ok((page, total)) +} + /// What a purge did. #[derive(Debug, Clone, Default, PartialEq, Eq)] pub struct Purged { @@ -672,6 +786,72 @@ mod tests { assert_eq!(parse_key(&any.key, KIND_TIME, 3), None); } + #[test] + fn filters_match() { + let entry = Entry { + queue_id: 1, + at: 100, + direction: Direction::Outgoing, + sender: "Alice@example.com".into(), + authenticated: true, + recipients: vec!["pay@bank.example".into()], + subject: "Q3 figures, final".into(), + message_id: "".into(), + accounts: vec![3], + tenants: vec![], + journals: vec![2], + held: false, + blob: String::new(), + size: 0, + sha256: String::new(), + expires_at: 0, + }; + let yes = |f: Filter| assert!(f.matches(&entry), "{f:?}"); + let no = |f: Filter| assert!(!f.matches(&entry), "{f:?}"); + yes(Filter::default()); + yes(Filter { + sender: Some("alice@".into()), + ..Default::default() + }); + yes(Filter { + address: Some("BANK".into()), + ..Default::default() + }); + yes(Filter { + text: Some("final q3".into()), + ..Default::default() + }); + yes(Filter { + message_id: Some("abc@example.com".into()), + ..Default::default() + }); + yes(Filter { + direction: Some(Direction::Any), + ..Default::default() + }); + no(Filter { + direction: Some(Direction::Incoming), + ..Default::default() + }); + no(Filter { + recipient: Some("alice".into()), + ..Default::default() + }); + no(Filter { + before: Some(100), + ..Default::default() + }); + yes(Filter { + after: Some(100), + journal_id: Some(2), + ..Default::default() + }); + no(Filter { + journal_id: Some(5), + ..Default::default() + }); + } + #[test] fn hex_round_trips() { let bytes = [0u8, 1, 0xab, 0xff]; diff --git a/crates/jmap-proto/src/object/inbuxa_journal_entry.rs b/crates/jmap-proto/src/object/inbuxa_journal_entry.rs new file mode 100644 index 0000000..69ea0f9 --- /dev/null +++ b/crates/jmap-proto/src/object/inbuxa_journal_entry.rs @@ -0,0 +1,340 @@ +/* + * SPDX-FileCopyrightText: 2026 Coffey Labs + * + * SPDX-License-Identifier: AGPL-3.0-only + */ + +//! The journal's JMAP objects under `urn:inbuxa:jmap` (journaling spec, +//! JR-6, JR-15 to JR-17): +//! +//! - `inbuxa:JournalEntry/get` and `/query`: what was journaled, read-only. +//! `report` (the whole journal report) comes only when asked for. +//! - `inbuxa:JournalExport/set`: create one to get a ZIP of the reports a +//! filter matches. +//! - `inbuxa:JournalVerification/set`: create one to recheck every chain. +//! +//! They share one set of properties. Nested values (an export's filter, a +//! verification's chains) are plain JSON objects. + +use crate::{ + object::{AnyId, JmapObject, JmapObjectId}, + request::deserialize::DeserializeArguments, +}; +use jmap_tools::{Element, Key, Property}; +use std::{borrow::Cow, str::FromStr}; +use types::id::Id; + +#[derive(Debug, Clone, Default)] +pub struct JournalEntry; + +#[derive(Debug, Clone, Default)] +pub struct JournalExport; + +#[derive(Debug, Clone, Default)] +pub struct JournalVerification; + +#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)] +pub enum JournalEntryProperty { + Id, + ReceivedAt, + Direction, + Sender, + Authenticated, + Recipients, + Subject, + MessageId, + JournalIds, + Held, + Size, + Sha256, + ExpiresAt, + Report, + Filter, + Reason, + BlobId, + Count, + Verified, + Chains, +} + +#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)] +pub enum JournalEntryValue { + Id(Id), +} + +impl Property for JournalEntryProperty { + fn try_parse(parent: Option<&Key<'_, Self>>, value: &str) -> Option { + // Keys inside a filter or a chain report stay plain keys + match parent { + None => JournalEntryProperty::parse(value), + Some(_) => None, + } + } + + fn to_cow(&self) -> Cow<'static, str> { + match self { + JournalEntryProperty::Id => "id", + JournalEntryProperty::ReceivedAt => "receivedAt", + JournalEntryProperty::Direction => "direction", + JournalEntryProperty::Sender => "sender", + JournalEntryProperty::Authenticated => "authenticated", + JournalEntryProperty::Recipients => "recipients", + JournalEntryProperty::Subject => "subject", + JournalEntryProperty::MessageId => "messageId", + JournalEntryProperty::JournalIds => "journalIds", + JournalEntryProperty::Held => "held", + JournalEntryProperty::Size => "size", + JournalEntryProperty::Sha256 => "sha256", + JournalEntryProperty::ExpiresAt => "expiresAt", + JournalEntryProperty::Report => "report", + JournalEntryProperty::Filter => "filter", + JournalEntryProperty::Reason => "reason", + JournalEntryProperty::BlobId => "blobId", + JournalEntryProperty::Count => "count", + JournalEntryProperty::Verified => "verified", + JournalEntryProperty::Chains => "chains", + } + .into() + } +} + +impl JournalEntryProperty { + fn parse(value: &str) -> Option { + hashify::tiny_map!(value.as_bytes(), + b"id" => JournalEntryProperty::Id, + b"receivedAt" => JournalEntryProperty::ReceivedAt, + b"direction" => JournalEntryProperty::Direction, + b"sender" => JournalEntryProperty::Sender, + b"authenticated" => JournalEntryProperty::Authenticated, + b"recipients" => JournalEntryProperty::Recipients, + b"subject" => JournalEntryProperty::Subject, + b"messageId" => JournalEntryProperty::MessageId, + b"journalIds" => JournalEntryProperty::JournalIds, + b"held" => JournalEntryProperty::Held, + b"size" => JournalEntryProperty::Size, + b"sha256" => JournalEntryProperty::Sha256, + b"expiresAt" => JournalEntryProperty::ExpiresAt, + b"report" => JournalEntryProperty::Report, + b"filter" => JournalEntryProperty::Filter, + b"reason" => JournalEntryProperty::Reason, + b"blobId" => JournalEntryProperty::BlobId, + b"count" => JournalEntryProperty::Count, + b"verified" => JournalEntryProperty::Verified, + b"chains" => JournalEntryProperty::Chains, + ) + } +} + +impl FromStr for JournalEntryProperty { + type Err = (); + + fn from_str(s: &str) -> Result { + JournalEntryProperty::parse(s).ok_or(()) + } +} + +impl Element for JournalEntryValue { + type Property = JournalEntryProperty; + + fn try_parse

(key: &Key<'_, Self::Property>, value: &str) -> Option { + match key { + Key::Property(JournalEntryProperty::Id) => { + Id::from_str(value).ok().map(JournalEntryValue::Id) + } + _ => None, + } + } + + fn to_cow(&self) -> Cow<'static, str> { + match self { + JournalEntryValue::Id(id) => id.to_string().into(), + } + } +} + +/// One condition of an `inbuxa:JournalEntry/query` filter. Several in one +/// filter object must all hold. +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum JournalFilter { + /// From this time on (UTC date). + After(String), + /// Before this time (UTC date). + Before(String), + /// Part of the sender's address. + Sender(String), + /// Part of a recipient's address. + Recipient(String), + /// Part of the sender's or a recipient's address. + Address(String), + /// `outgoing`, `incoming` or `internal`. + Direction(String), + /// Words that must all be in the subject. + Text(String), + MessageId(String), + JournalId(Id), + _T(String), +} + +impl Default for JournalFilter { + fn default() -> Self { + JournalFilter::_T(String::new()) + } +} + +impl<'de> DeserializeArguments<'de> for JournalFilter { + fn deserialize_argument(&mut self, key: &str, map: &mut A) -> Result<(), A::Error> + where + A: serde::de::MapAccess<'de>, + { + hashify::fnc_map!(key.as_bytes(), + b"after" => { + *self = JournalFilter::After(map.next_value()?); + }, + b"before" => { + *self = JournalFilter::Before(map.next_value()?); + }, + b"sender" => { + *self = JournalFilter::Sender(map.next_value()?); + }, + b"recipient" => { + *self = JournalFilter::Recipient(map.next_value()?); + }, + b"address" => { + *self = JournalFilter::Address(map.next_value()?); + }, + b"direction" => { + *self = JournalFilter::Direction(map.next_value()?); + }, + b"text" => { + *self = JournalFilter::Text(map.next_value()?); + }, + b"messageId" => { + *self = JournalFilter::MessageId(map.next_value()?); + }, + b"journalId" => { + *self = JournalFilter::JournalId(map.next_value()?); + }, + _ => { + *self = JournalFilter::_T(key.to_string()); + let _ = map.next_value::()?; + } + ); + + Ok(()) + } +} + +/// Entries sort newest first, by `receivedAt`; nothing else. +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum JournalComparator { + ReceivedAt, + _T(String), +} + +impl Default for JournalComparator { + fn default() -> Self { + JournalComparator::_T(String::new()) + } +} + +impl<'de> DeserializeArguments<'de> for JournalComparator { + fn deserialize_argument(&mut self, key: &str, map: &mut A) -> Result<(), A::Error> + where + A: serde::de::MapAccess<'de>, + { + if key == "property" { + let value = map.next_value::>()?; + *self = if value == "receivedAt" { + JournalComparator::ReceivedAt + } else { + JournalComparator::_T(value.into_owned()) + }; + } else { + let _ = map.next_value::()?; + } + Ok(()) + } +} + +macro_rules! journal_object { + ($object:ty, $filter:ty, $comparator:ty) => { + impl JmapObject for $object { + type Property = JournalEntryProperty; + + type Element = JournalEntryValue; + + type Id = Id; + + type Filter = $filter; + + type Comparator = $comparator; + + type GetArguments = (); + + type SetArguments<'de> = (); + + type QueryArguments = (); + + type CopyArguments = (); + + type ParseArguments = (); + + const ID_PROPERTY: Self::Property = JournalEntryProperty::Id; + } + }; +} + +journal_object!(JournalEntry, JournalFilter, JournalComparator); +journal_object!(JournalExport, (), ()); +journal_object!(JournalVerification, (), ()); + +impl From for JournalEntryValue { + fn from(id: Id) -> Self { + JournalEntryValue::Id(id) + } +} + +impl JmapObjectId for JournalEntryValue { + fn as_id(&self) -> Option { + match self { + JournalEntryValue::Id(id) => Some(*id), + } + } + + fn as_any_id(&self) -> Option { + match self { + JournalEntryValue::Id(id) => Some(AnyId::Id(*id)), + } + } + + fn as_id_ref(&self) -> Option<&str> { + None + } + + fn try_set_id(&mut self, new_id: AnyId) -> bool { + if let AnyId::Id(id) = new_id { + *self = JournalEntryValue::Id(id); + true + } else { + false + } + } +} + +impl JmapObjectId for JournalEntryProperty { + fn as_id(&self) -> Option { + None + } + + fn as_any_id(&self) -> Option { + None + } + + fn as_id_ref(&self) -> Option<&str> { + None + } + + fn try_set_id(&mut self, _: AnyId) -> bool { + false + } +} diff --git a/crates/jmap-proto/src/object/mod.rs b/crates/jmap-proto/src/object/mod.rs index c7108ac..9a965f5 100644 --- a/crates/jmap-proto/src/object/mod.rs +++ b/crates/jmap-proto/src/object/mod.rs @@ -32,6 +32,7 @@ pub mod inbuxa_legal_hold; // inbuxa: legal hold pub mod inbuxa_mail_rule; // inbuxa: DLP and mail flow rules pub mod inbuxa_security_acceptance; // inbuxa: accepted security to-do items pub mod inbuxa_journal; // inbuxa: journaling +pub mod inbuxa_journal_entry; // inbuxa: journaling, search and export pub mod inbuxa_held_message; // inbuxa: mail held for review pub mod inbuxa_hold_export; // inbuxa: legal hold exports pub mod inbuxa_explanation; // inbuxa: "Explain this" with the local model diff --git a/crates/jmap-proto/src/references/eval.rs b/crates/jmap-proto/src/references/eval.rs index c6cd377..a944476 100644 --- a/crates/jmap-proto/src/references/eval.rs +++ b/crates/jmap-proto/src/references/eval.rs @@ -94,6 +94,9 @@ impl Response<'_> { GetResponseMethod::Journal(response) => { response.eval_jptr(path, &mut results) } + GetResponseMethod::JournalEntry(response) => { + response.eval_jptr(path, &mut results) + } GetResponseMethod::HeldMessage(response) => { response.eval_jptr(path, &mut results) } diff --git a/crates/jmap-proto/src/references/resolve.rs b/crates/jmap-proto/src/references/resolve.rs index d8f8686..0032fd4 100644 --- a/crates/jmap-proto/src/references/resolve.rs +++ b/crates/jmap-proto/src/references/resolve.rs @@ -57,6 +57,7 @@ impl Response<'_> { GetRequestMethod::MailRule(request) => request.resolve_references(self)?, GetRequestMethod::SecurityAcceptance(request) => request.resolve_references(self)?, GetRequestMethod::Journal(request) => request.resolve_references(self)?, + GetRequestMethod::JournalEntry(request) => request.resolve_references(self)?, GetRequestMethod::HeldMessage(request) => request.resolve_references(self)?, GetRequestMethod::HoldExport(request) => request.resolve_references(self)?, GetRequestMethod::ProtocolPolicy(request) => request.resolve_references(self)?, @@ -139,6 +140,12 @@ impl Response<'_> { SetRequestMethod::Journal(request) => { request.resolve_references(self, 1, false)? } + SetRequestMethod::JournalExport(request) => { + request.resolve_references(self, 1, false)? + } + SetRequestMethod::JournalVerification(request) => { + request.resolve_references(self, 1, false)? + } SetRequestMethod::HeldMessage(request) => { request.resolve_references(self, 1, false)? } diff --git a/crates/jmap-proto/src/request/method.rs b/crates/jmap-proto/src/request/method.rs index d89d7c0..d9ef681 100644 --- a/crates/jmap-proto/src/request/method.rs +++ b/crates/jmap-proto/src/request/method.rs @@ -73,6 +73,9 @@ pub enum MethodObject { HeldMessage, // inbuxa: journaling Journal, + JournalEntry, + JournalExport, + JournalVerification, TenantProtocolPolicy, } @@ -115,7 +118,10 @@ impl MethodObject { | MethodObject::MailRule | MethodObject::SecurityAcceptance | MethodObject::HeldMessage - | MethodObject::Journal => Capability::Inbuxa, + | MethodObject::Journal + | MethodObject::JournalEntry + | MethodObject::JournalExport + | MethodObject::JournalVerification => Capability::Inbuxa, MethodObject::ProtocolPolicy => Capability::Inbuxa, MethodObject::TenantProtocolPolicy => Capability::Inbuxa, } @@ -317,6 +323,12 @@ impl MethodName { (MethodFunction::Set, MethodObject::SecurityAcceptance) => "inbuxa:SecurityAcceptance/set", (MethodFunction::Get, MethodObject::Journal) => "inbuxa:Journal/get", (MethodFunction::Set, MethodObject::Journal) => "inbuxa:Journal/set", + (MethodFunction::Get, MethodObject::JournalEntry) => "inbuxa:JournalEntry/get", + (MethodFunction::Query, MethodObject::JournalEntry) => "inbuxa:JournalEntry/query", + (MethodFunction::Set, MethodObject::JournalExport) => "inbuxa:JournalExport/set", + (MethodFunction::Set, MethodObject::JournalVerification) => { + "inbuxa:JournalVerification/set" + } (MethodFunction::Get, MethodObject::HeldMessage) => "inbuxa:HeldMessage/get", (MethodFunction::Set, MethodObject::HeldMessage) => "inbuxa:HeldMessage/set", (MethodFunction::Get, MethodObject::HoldExport) => "inbuxa:HoldExport/get", @@ -479,6 +491,10 @@ impl MethodName { "inbuxa:SecurityAcceptance/set" => (MethodObject::SecurityAcceptance, MethodFunction::Set), "inbuxa:Journal/get" => (MethodObject::Journal, MethodFunction::Get), "inbuxa:Journal/set" => (MethodObject::Journal, MethodFunction::Set), + "inbuxa:JournalEntry/get" => (MethodObject::JournalEntry, MethodFunction::Get), + "inbuxa:JournalEntry/query" => (MethodObject::JournalEntry, MethodFunction::Query), + "inbuxa:JournalExport/set" => (MethodObject::JournalExport, MethodFunction::Set), + "inbuxa:JournalVerification/set" => (MethodObject::JournalVerification, MethodFunction::Set), "inbuxa:HeldMessage/get" => (MethodObject::HeldMessage, MethodFunction::Get), "inbuxa:HeldMessage/set" => (MethodObject::HeldMessage, MethodFunction::Set), "inbuxa:HoldExport/get" => (MethodObject::HoldExport, MethodFunction::Get), @@ -555,6 +571,9 @@ impl Display for MethodObject { MethodObject::MailRule => "inbuxa:MailRule", MethodObject::SecurityAcceptance => "inbuxa:SecurityAcceptance", MethodObject::Journal => "inbuxa:Journal", + MethodObject::JournalEntry => "inbuxa:JournalEntry", + MethodObject::JournalExport => "inbuxa:JournalExport", + MethodObject::JournalVerification => "inbuxa:JournalVerification", MethodObject::HeldMessage => "inbuxa:HeldMessage", MethodObject::HoldExport => "inbuxa:HoldExport", MethodObject::ProtocolPolicy => "inbuxa:ProtocolPolicy", diff --git a/crates/jmap-proto/src/request/mod.rs b/crates/jmap-proto/src/request/mod.rs index 48397a6..48c0de7 100644 --- a/crates/jmap-proto/src/request/mod.rs +++ b/crates/jmap-proto/src/request/mod.rs @@ -127,6 +127,7 @@ pub enum GetRequestMethod { MailRule(Box>), SecurityAcceptance(Box>), Journal(Box>), + JournalEntry(Box>), HeldMessage(Box>), HoldExport(Box>), ProtocolPolicy(Box>), @@ -169,6 +170,8 @@ pub enum SetRequestMethod<'x> { Box>, ), Journal(Box>), + JournalExport(Box>), + JournalVerification(Box>), HeldMessage(Box>), HoldExport(Box>), ProtocolPolicy(Box>), @@ -203,6 +206,7 @@ pub enum QueryRequestMethod { ShareNotification(Box>), Registry(Box>), AuditEvent(Box>), + JournalEntry(Box>), } #[derive(Debug)] diff --git a/crates/jmap-proto/src/request/parser.rs b/crates/jmap-proto/src/request/parser.rs index 0d09910..518e8f3 100644 --- a/crates/jmap-proto/src/request/parser.rs +++ b/crates/jmap-proto/src/request/parser.rs @@ -669,6 +669,34 @@ impl<'de> Visitor<'de> for CallVisitor { } }, // inbuxa: journaling + (MethodFunction::Get, MethodObject::JournalEntry) => match seq.next_element() { + Ok(Some(value)) => RequestMethod::Get(GetRequestMethod::JournalEntry(value)), + Err(err) => RequestMethod::invalid(err), + Ok(None) => { + return Err(de::Error::invalid_length(1, &self)); + } + }, + (MethodFunction::Query, MethodObject::JournalEntry) => match seq.next_element() { + Ok(Some(value)) => RequestMethod::Query(QueryRequestMethod::JournalEntry(value)), + Err(err) => RequestMethod::invalid(err), + Ok(None) => { + return Err(de::Error::invalid_length(1, &self)); + } + }, + (MethodFunction::Set, MethodObject::JournalExport) => match seq.next_element() { + Ok(Some(value)) => RequestMethod::Set(SetRequestMethod::JournalExport(value)), + Err(err) => RequestMethod::invalid(err), + Ok(None) => { + return Err(de::Error::invalid_length(1, &self)); + } + }, + (MethodFunction::Set, MethodObject::JournalVerification) => match seq.next_element() { + Ok(Some(value)) => RequestMethod::Set(SetRequestMethod::JournalVerification(value)), + Err(err) => RequestMethod::invalid(err), + Ok(None) => { + return Err(de::Error::invalid_length(1, &self)); + } + }, (MethodFunction::Get, MethodObject::Journal) => match seq.next_element() { Ok(Some(value)) => RequestMethod::Get(GetRequestMethod::Journal(value)), Err(err) => RequestMethod::invalid(err), diff --git a/crates/jmap-proto/src/response/mod.rs b/crates/jmap-proto/src/response/mod.rs index 03f84ac..500ab4e 100644 --- a/crates/jmap-proto/src/response/mod.rs +++ b/crates/jmap-proto/src/response/mod.rs @@ -114,6 +114,7 @@ pub enum GetResponseMethod { MailRule(GetResponse), SecurityAcceptance(GetResponse), Journal(GetResponse), + JournalEntry(GetResponse), HeldMessage(GetResponse), HoldExport(GetResponse), ProtocolPolicy(GetResponse), @@ -156,6 +157,8 @@ pub enum SetResponseMethod { Box>, ), Journal(Box>), + JournalExport(Box>), + JournalVerification(Box>), HeldMessage(Box>), HoldExport(Box>), Explanation(Box>), @@ -865,6 +868,24 @@ impl<'x> From> for Respon } // inbuxa: journaling +impl<'x> From> for ResponseMethod<'x> { + fn from(value: GetResponse) -> Self { + ResponseMethod::Get(GetResponseMethod::JournalEntry(value)) + } +} + +impl<'x> From> for ResponseMethod<'x> { + fn from(value: SetResponse) -> Self { + ResponseMethod::Set(SetResponseMethod::JournalExport(Box::new(value))) + } +} + +impl<'x> From> for ResponseMethod<'x> { + fn from(value: SetResponse) -> Self { + ResponseMethod::Set(SetResponseMethod::JournalVerification(Box::new(value))) + } +} + impl<'x> From> for ResponseMethod<'x> { fn from(value: GetResponse) -> Self { ResponseMethod::Get(GetResponseMethod::Journal(value)) diff --git a/crates/jmap/src/api/auth.rs b/crates/jmap/src/api/auth.rs index 710b6c9..4884c05 100644 --- a/crates/jmap/src/api/auth.rs +++ b/crates/jmap/src/api/auth.rs @@ -118,6 +118,7 @@ impl JmapAuthorization for AccessToken { } // inbuxa: journaling (JR-18) GetRequestMethod::Journal(_) => Permission::SysJournalGet, + GetRequestMethod::JournalEntry(_) => Permission::SysJournalSearch, GetRequestMethod::HoldExport(_) => Permission::SysLegalHoldExport, // inbuxa: accepted security items are read by whoever may // see the server's security settings @@ -301,6 +302,20 @@ impl JmapAuthorization for AccessToken { Permission::SysJournalUpdate, Permission::SysJournalUpdate, ), + SetRequestMethod::JournalExport(s) => validate_set( + s, + self, + Permission::SysJournalExport, + Permission::SysJournalExport, + Permission::SysJournalExport, + ), + SetRequestMethod::JournalVerification(s) => validate_set( + s, + self, + Permission::SysJournalGet, + Permission::SysJournalGet, + Permission::SysJournalGet, + ), // inbuxa: accepting a security to-do item, or removing // an acceptance; nothing is ever edited SetRequestMethod::SecurityAcceptance(s) => { @@ -481,6 +496,9 @@ impl JmapAuthorization for AccessToken { | MethodObject::SecurityAcceptance | MethodObject::HeldMessage | MethodObject::Journal + | MethodObject::JournalEntry + | MethodObject::JournalExport + | MethodObject::JournalVerification | MethodObject::ProtocolPolicy | MethodObject::TenantProtocolPolicy => Permission::JmapEmailChanges, // inbuxa: x:MaskedEmail/changes reads what /get reads @@ -539,6 +557,8 @@ impl JmapAuthorization for AccessToken { QueryRequestMethod::ShareNotification(_) => Permission::JmapShareNotificationQuery, // inbuxa: the audit log (AU-9) QueryRequestMethod::AuditEvent(_) => Permission::SysAuditGet, + // inbuxa: journaling (JR-15) + QueryRequestMethod::JournalEntry(_) => Permission::SysJournalSearch, QueryRequestMethod::Registry(_) => { let MethodObject::Registry(object_type) = object else { unreachable!() diff --git a/crates/jmap/src/api/request.rs b/crates/jmap/src/api/request.rs index 0bd8ffe..bff61e1 100644 --- a/crates/jmap/src/api/request.rs +++ b/crates/jmap/src/api/request.rs @@ -291,6 +291,12 @@ impl RequestHandler for Server { SetResponseMethod::Journal(set_response) => { set_response.update_created_ids(&mut response); } + SetResponseMethod::JournalExport(set_response) => { + set_response.update_created_ids(&mut response); + } + SetResponseMethod::JournalVerification(set_response) => { + set_response.update_created_ids(&mut response); + } SetResponseMethod::HeldMessage(set_response) => { set_response.update_created_ids(&mut response); } @@ -530,6 +536,12 @@ impl RequestHandler for Server { resolve_account_id(&mut req.account_id, method_name.obj, access_token)?; crate::inbuxa::journal::get(self, access_token, *req).await?.into() } + GetRequestMethod::JournalEntry(mut req) => { + resolve_account_id(&mut req.account_id, method_name.obj, access_token)?; + crate::inbuxa::journal_entry::get(self, access_token, session, *req) + .await? + .into() + } // inbuxa: the audit log (AU-9) GetRequestMethod::AuditEvent(mut req) => { resolve_account_id(&mut req.account_id, method_name.obj, access_token)?; @@ -726,6 +738,13 @@ impl RequestHandler for Server { .await? .into() } + // inbuxa: journaling (JR-15) + QueryRequestMethod::JournalEntry(mut req) => { + resolve_account_id(&mut req.account_id, method_name.obj, access_token)?; + crate::inbuxa::journal_entry::query(self, access_token, session, *req) + .await? + .into() + } QueryRequestMethod::Registry(mut req) => { resolve_account_id(&mut req.account_id, method_name.obj, access_token)?; assert_registry_account(self, method_name.obj, access_token, req.account_id) @@ -1054,6 +1073,19 @@ impl RequestHandler for Server { .await? .into() } + // inbuxa: journaling (JR-6, JR-16) + SetRequestMethod::JournalExport(mut req) => { + resolve_account_id(&mut req.account_id, method_name.obj, access_token)?; + crate::inbuxa::journal_entry::export_set(self, access_token, session, *req) + .await? + .into() + } + SetRequestMethod::JournalVerification(mut req) => { + resolve_account_id(&mut req.account_id, method_name.obj, access_token)?; + crate::inbuxa::journal_entry::verification_set(self, access_token, session, *req) + .await? + .into() + } // inbuxa: inbuxa:Explanation/set ("Explain this") SetRequestMethod::Explanation(mut req) => { resolve_account_id(&mut req.account_id, method_name.obj, access_token)?; diff --git a/crates/jmap/src/changes/get.rs b/crates/jmap/src/changes/get.rs index 3e50441..f18b7de 100644 --- a/crates/jmap/src/changes/get.rs +++ b/crates/jmap/src/changes/get.rs @@ -433,6 +433,9 @@ impl IntermediateChangesResponse { | MethodObject::MailRule | MethodObject::SecurityAcceptance | MethodObject::Journal + | MethodObject::JournalEntry + | MethodObject::JournalExport + | MethodObject::JournalVerification | MethodObject::HeldMessage | MethodObject::ProtocolPolicy | MethodObject::TenantProtocolPolicy diff --git a/crates/jmap/src/inbuxa/journal_entry.rs b/crates/jmap/src/inbuxa/journal_entry.rs new file mode 100644 index 0000000..61227ad --- /dev/null +++ b/crates/jmap/src/inbuxa/journal_entry.rs @@ -0,0 +1,791 @@ +/* + * SPDX-FileCopyrightText: 2026 Coffey Labs + * + * SPDX-License-Identifier: AGPL-3.0-only + */ + +//! The journal over JMAP (journaling spec, JR-6, JR-15 to JR-17): +//! +//! - `inbuxa:JournalEntry/query` and `/get`: searching and reading what was +//! journaled (`sysJournalSearch`). `report` is the whole journal report, +//! only when asked for. +//! - `inbuxa:JournalExport/set`: a ZIP of the reports a filter matches, in +//! the hold export's shape (`sysJournalExport`). +//! - `inbuxa:JournalVerification/set`: rechecks every chain and every report +//! (`sysJournalGet`). +//! +//! Every search, read and export is written to the audit log first; if it +//! can't be, nothing is returned (JR-17). All of it is the server's: nobody +//! in a tenant reaches it. + +use common::{Server, auth::AccessToken}; +use http_proto::HttpSessionData; +use inbuxa_features::{ + audit::{Action, Outcome, Record, Target}, + journal::{ + Direction, + entries::{self, ChainReport, Entry, EntryId, Filter, MAX_QUERY_LIMIT}, + }, +}; +use jmap_proto::{ + error::set::SetError, + method::{ + get::{GetRequest, GetResponse}, + query::{Filter as QueryFilter, QueryRequest, QueryResponse}, + set::{SetRequest, SetResponse}, + }, + object::inbuxa_journal_entry::{ + JournalEntry, JournalEntryProperty as P, JournalEntryValue, JournalExport, JournalFilter, + JournalVerification, + }, + request::IntoValid, + types::{date::UTCDate, state::State}, +}; +use jmap_tools::{Key, Map, Value}; +use sha2::{Digest, Sha256}; +use std::{ + borrow::Cow, + io::{Cursor, Write}, + str::FromStr, +}; +use types::id::Id; +use zip::{CompressionMethod, ZipWriter, write::SimpleFileOptions}; + +type JValue = Value<'static, P, JournalEntryValue>; + +/// Properties a get returns unless asked otherwise: all but the report. +const LISTED: &[P] = &[ + P::Id, + P::ReceivedAt, + P::Direction, + P::Sender, + P::Authenticated, + P::Recipients, + P::Subject, + P::MessageId, + P::JournalIds, + P::Held, + P::Size, + P::Sha256, + P::ExpiresAt, +]; + +/// Most reports one export holds, and most bytes. +const MAX_EXPORT_ENTRIES: usize = 10_000; +const MAX_EXPORT_BYTES: u64 = 1024 * 1024 * 1024; +/// Most of one report `get` returns as text. +const MAX_REPORT_TEXT: usize = 10 * 1024 * 1024; + +fn server_level(access_token: &AccessToken) -> trc::Result<()> { + if access_token.tenant_id().is_some() { + Err(trc::JmapEvent::Forbidden + .into_err() + .details("The journal is the server's.")) + } else { + Ok(()) + } +} + +fn date(seconds: u64) -> JValue { + Value::Str(UTCDate::from_timestamp(seconds as i64).to_string().into()) +} + +fn text(value: &str) -> JValue { + Value::Str(value.to_string().into()) +} + +fn json_to_value(json: serde_json::Value) -> JValue { + match json { + serde_json::Value::Null => Value::Null, + serde_json::Value::Bool(b) => Value::Bool(b), + serde_json::Value::Number(n) => match n.as_u64() { + Some(n) => Value::Number(n.into()), + None => Value::Number(n.as_i64().unwrap_or_default().into()), + }, + serde_json::Value::String(s) => Value::Str(Cow::Owned(s)), + serde_json::Value::Array(items) => { + Value::Array(items.into_iter().map(json_to_value).collect()) + } + serde_json::Value::Object(map) => { + let mut out = Map::with_capacity(map.len()); + for (key, value) in map { + out.insert_unchecked(Key::Owned(key), json_to_value(value)); + } + Value::Object(out) + } + } +} + +fn entry_value(id: EntryId, entry: &Entry, report: Option<&str>, properties: &[P]) -> JValue { + let mut out = Map::with_capacity(properties.len()); + for property in properties { + let value = match property { + P::Id => Value::Element(JournalEntryValue::Id(Id::new(id.to_u64()))), + P::ReceivedAt => date(entry.at), + P::Direction => text(entry.direction.as_str()), + P::Sender => text(&entry.sender), + P::Authenticated => Value::Bool(entry.authenticated), + P::Recipients => Value::Array(entry.recipients.iter().map(|r| text(r)).collect()), + P::Subject => text(&entry.subject), + P::MessageId => text(&entry.message_id), + P::JournalIds => Value::Array( + entry + .journals + .iter() + .map(|j| text(&Id::from(*j).to_string())) + .collect(), + ), + P::Held => Value::Bool(entry.held), + P::Size => Value::Number(entry.size.into()), + P::Sha256 => text(&entry.sha256), + P::ExpiresAt => date(entry.expires_at), + P::Report => report.map_or(Value::Null, text), + _ => Value::Null, + }; + out.insert_unchecked(Key::Property(property.clone()), value); + } + Value::Object(out) +} + +/// Writes a record before anything is returned; an error means nothing +/// may be (JR-17). +async fn record( + server: &Server, + access_token: &AccessToken, + session: &HttpSessionData, + action: Action, + target_id: Option, + target_name: Option, + details: String, + reason: Option, +) -> trc::Result<()> { + server + .audit_append(&Record { + at: store::write::now() * 1000, + actor: server.audit_actor(access_token).await, + via: access_token.origin().cloned(), + remote_ip: Some(session.remote_ip), + action, + target: Target { + kind: "inbuxa:JournalEntry".into(), + id: target_id, + name: target_name, + ..Default::default() + }, + changes: vec![], + details: Some(details), + reason, + outcome: Outcome::success(), + }) + .await + .map(|_| ()) + .map_err(|err| { + err.details("The audit log couldn't be written, so the journal wasn't read.") + }) +} + +async fn report_bytes(server: &Server, entry: &Entry) -> trc::Result>> { + match entry.blob_hash() { + Some(hash) => { + server + .blob_store() + .get_blob(hash.as_slice(), 0..usize::MAX) + .await + } + None => Ok(None), + } +} + +/// `inbuxa:JournalEntry/get`: the entries named. Listing them is recorded +/// once; each report read is recorded on its own. +pub async fn get( + server: &Server, + access_token: &AccessToken, + session: &HttpSessionData, + mut request: GetRequest, +) -> trc::Result> { + server_level(access_token)?; + let properties = request.unwrap_properties(LISTED); + let (ids, not_found) = request.unwrap_ids(server.core.jmap.get_max_objects)?; + let mut response = GetResponse { + account_id: request.account_id.into(), + state: None, + list: Vec::new(), + not_found, + }; + let Some(ids) = ids else { + return Err(trc::JmapEvent::RequestTooLarge + .into_err() + .details("Name the entries to get; use inbuxa:JournalEntry/query to find them.")); + }; + let mut found = Vec::with_capacity(ids.len()); + for id in ids { + let entry_id = EntryId::from_u64(id.id()); + match entries::get(server.store(), entry_id).await? { + Some(entry) => found.push((entry_id, entry)), + None => response.push_not_found(id), + } + } + if found.is_empty() { + return Ok(response); + } + let with_report = properties.contains(&P::Report); + if !with_report { + record( + server, + access_token, + session, + Action::BlobAccess, + None, + None, + format!("Listed {} journal entries", found.len()), + None, + ) + .await?; + } + for (entry_id, entry) in found { + let report = if with_report { + record( + server, + access_token, + session, + Action::BlobAccess, + Some(Id::new(entry_id.to_u64()).to_string()), + Some(entry.subject.clone()), + format!("Read a journaled message from {}", entry.sender), + None, + ) + .await?; + report_bytes(server, &entry).await?.map(|bytes| { + let end = bytes.len().min(MAX_REPORT_TEXT); + String::from_utf8_lossy(&bytes[..end]).into_owned() + }) + } else { + None + }; + response.list.push(entry_value( + entry_id, + &entry, + report.as_deref(), + &properties, + )); + } + Ok(response) +} + +fn seconds(value: &str) -> Result { + UTCDate::from_str(value) + .map(|date| date.timestamp().max(0) as u64) + .map_err(|_| format!("{value} isn't a UTC date.")) +} + +fn direction(value: &str) -> Result { + match value { + "outgoing" => Ok(Direction::Outgoing), + "incoming" => Ok(Direction::Incoming), + "internal" => Ok(Direction::Internal), + "any" => Ok(Direction::Any), + other => Err(format!("{other} isn't a direction.")), + } +} + +/// The conditions of a query filter, all of which must hold. `Or` and +/// `Not` aren't supported. +fn build_filter(conditions: Vec>) -> trc::Result { + let unsupported = |why: String| trc::JmapEvent::UnsupportedFilter.into_err().details(why); + let mut filter = Filter::default(); + for condition in conditions { + match condition { + QueryFilter::Property(condition) => match condition { + JournalFilter::After(date) => { + filter.after = Some(seconds(&date).map_err(unsupported)?) + } + JournalFilter::Before(date) => { + filter.before = Some(seconds(&date).map_err(unsupported)?) + } + JournalFilter::Sender(s) => filter.sender = Some(s), + JournalFilter::Recipient(r) => filter.recipient = Some(r), + JournalFilter::Address(a) => filter.address = Some(a), + JournalFilter::Direction(d) => { + filter.direction = Some(direction(&d).map_err(unsupported)?) + } + JournalFilter::Text(t) => filter.text = Some(t), + JournalFilter::MessageId(m) => filter.message_id = Some(m), + JournalFilter::JournalId(id) => filter.journal_id = Some(id.document_id()), + JournalFilter::_T(other) => { + return Err(unsupported(format!("Unknown filter property {other}."))); + } + }, + QueryFilter::And | QueryFilter::Close => {} + QueryFilter::Or | QueryFilter::Not => { + return Err(unsupported( + "Journal searches take conditions that must all hold; OR and NOT aren't \ + supported." + .into(), + )); + } + } + } + Ok(filter) +} + +fn filter_text(filter: &Filter) -> String { + serde_json::to_string(filter).unwrap_or_default() +} + +/// `inbuxa:JournalEntry/query`: newest first. The search is recorded, with +/// its terms, before anything is returned. +pub async fn query( + server: &Server, + access_token: &AccessToken, + session: &HttpSessionData, + request: QueryRequest, +) -> trc::Result { + server_level(access_token)?; + let filter = build_filter(request.filter)?; + let position = request.position.unwrap_or(0); + if position < 0 || request.anchor.is_some() { + return Err(trc::JmapEvent::UnsupportedFilter + .into_err() + .details("Journal searches page by a position from the start.")); + } + let limit = request + .limit + .unwrap_or(MAX_QUERY_LIMIT) + .min(MAX_QUERY_LIMIT); + let count_all = request.calculate_total.unwrap_or(false); + record( + server, + access_token, + session, + Action::BlobAccess, + None, + None, + format!("Searched the journal: {}", filter_text(&filter)), + None, + ) + .await?; + let (ids, total) = + entries::query(server.store(), &filter, position as usize, limit, count_all).await?; + Ok(QueryResponse { + account_id: request.account_id, + query_state: State::Initial, + can_calculate_changes: false, + position, + ids: ids.into_iter().map(|id| Id::new(id.to_u64())).collect(), + total: count_all.then_some(total), + limit: Some(limit), + }) +} + +/// An export's filter, as sent: the query's conditions in one object. +fn export_filter(value: Option>) -> Result { + let json: serde_json::Value = value + .map(Into::into) + .unwrap_or(serde_json::Value::Object(Default::default())); + let serde_json::Value::Object(map) = json else { + return Err("The filter is an object of conditions.".into()); + }; + let mut filter = Filter::default(); + for (key, value) in map { + let text = || { + value + .as_str() + .map(str::to_string) + .ok_or_else(|| format!("{key} is text.")) + }; + match key.as_str() { + "after" => filter.after = Some(seconds(&text()?)?), + "before" => filter.before = Some(seconds(&text()?)?), + "sender" => filter.sender = Some(text()?), + "recipient" => filter.recipient = Some(text()?), + "address" => filter.address = Some(text()?), + "direction" => filter.direction = Some(direction(&text()?)?), + "text" => filter.text = Some(text()?), + "messageId" => filter.message_id = Some(text()?), + "journalId" => { + filter.journal_id = Some( + Id::from_str(&text()?) + .map_err(|_| "journalId is a journal's id.".to_string())? + .document_id(), + ) + } + other => return Err(format!("Unknown filter property {other}.")), + } + } + Ok(filter) +} + +fn csv(field: &str) -> String { + if field.contains([',', '"', '\n', '\r']) { + format!("\"{}\"", field.replace('"', "\"\"")) + } else { + field.to_string() + } +} + +fn hex(bytes: &[u8]) -> String { + bytes.iter().map(|b| format!("{b:02x}")).collect() +} + +/// A ZIP of reports in the hold export's shape: each report as `.eml`, +/// `manifest.csv` with the envelope and a SHA-256 per file, the entries +/// whose report couldn't be read in `exceptions.csv`, and +/// `manifest.sha256` over both. Returns its bytes and how many reports went +/// in. +pub(crate) fn build_zip( + items: &[(EntryId, Entry, Option>)], +) -> trc::Result<(Vec, usize)> { + let fail = |err: zip::result::ZipError| { + trc::StoreEvent::UnexpectedError + .into_err() + .details("Failed to write the export") + .reason(err) + }; + let options = SimpleFileOptions::default().compression_method(CompressionMethod::Deflated); + let mut zip = ZipWriter::new(Cursor::new(Vec::new())); + let mut manifest = String::from( + "path,receivedAt,direction,sender,recipients,subject,messageId,queueId,size,sha256\n", + ); + let mut exceptions = String::from("entry,receivedAt,sender,subject,reason\n"); + let mut written = 0u64; + let mut count = 0; + for (id, entry, bytes) in items { + let received = UTCDate::from_timestamp(entry.at as i64).to_string(); + let Some(bytes) = bytes else { + exceptions.push_str(&format!( + "{},{},{},{},{}\n", + Id::new(id.to_u64()), + received, + csv(&entry.sender), + csv(&entry.subject), + "The report couldn't be read." + )); + continue; + }; + written += bytes.len() as u64; + if written > MAX_EXPORT_BYTES { + return Err(trc::StoreEvent::UnexpectedError.into_err().details( + "The reports are larger than one export can hold (1 GB). Narrow the search.", + )); + } + let path = format!( + "reports/{}-{:x}.eml", + received.replace(':', ""), + entry.queue_id + ); + zip.start_file(path.as_str(), options).map_err(fail)?; + zip.write_all(bytes).map_err(|e| fail(e.into()))?; + manifest.push_str(&format!( + "{},{},{},{},{},{},{},{:x},{},{}\n", + csv(&path), + received, + entry.direction.as_str(), + csv(&entry.sender), + csv(&entry.recipients.join(" ")), + csv(&entry.subject), + csv(&entry.message_id), + entry.queue_id, + bytes.len(), + hex(&Sha256::digest(bytes)) + )); + count += 1; + } + let manifest_hash = hex(&Sha256::digest(manifest.as_bytes())); + let exceptions_hash = hex(&Sha256::digest(exceptions.as_bytes())); + zip.start_file("manifest.csv", options).map_err(fail)?; + zip.write_all(manifest.as_bytes()) + .map_err(|e| fail(e.into()))?; + zip.start_file("exceptions.csv", options).map_err(fail)?; + zip.write_all(exceptions.as_bytes()) + .map_err(|e| fail(e.into()))?; + zip.start_file("manifest.sha256", options).map_err(fail)?; + zip.write_all( + format!("{manifest_hash} manifest.csv\n{exceptions_hash} exceptions.csv\n").as_bytes(), + ) + .map_err(|e| fail(e.into()))?; + Ok((zip.finish().map_err(fail)?.into_inner(), count)) +} + +/// `inbuxa:JournalExport/set`: create `{filter, reason}`; the created +/// object names the ZIP's blob (the caller's), its size, how many reports it +/// holds and its SHA-256. A reason is required; the export is recorded +/// before it's built. +pub async fn export_set( + server: &Server, + access_token: &AccessToken, + session: &HttpSessionData, + mut request: SetRequest<'_, JournalExport>, +) -> trc::Result> { + server_level(access_token)?; + let mut response = SetResponse::from_request(&request, server.core.jmap.set_max_objects)?; + for (id, _) in request.unwrap_update().into_valid() { + response.not_updated.append( + id, + SetError::forbidden().with_description("Exports can't be changed."), + ); + } + for id in request.unwrap_destroy().into_valid() { + response.not_destroyed.append( + id, + SetError::forbidden().with_description("Exports aren't kept to destroy."), + ); + } + + for (client_id, value) in request.unwrap_create() { + let mut filter_value = None; + let mut reason = None; + let mut invalid = None; + for (key, value) in value.into_expanded_object() { + match (&key, value) { + (Key::Property(P::Filter), value) => filter_value = Some(value.into_owned()), + (Key::Property(P::Reason), Value::Str(r)) => { + reason = Some(r.trim().chars().take(500).collect::()) + } + _ => { + invalid = Some(SetError::invalid_properties().with_property(key.into_owned())); + break; + } + } + } + if let Some(error) = invalid { + response.not_created.append(client_id, error); + continue; + } + let Some(reason) = reason.filter(|r| !r.is_empty()) else { + response.not_created.append( + client_id, + SetError::invalid_properties() + .with_property(P::Reason) + .with_description( + "Say why: a reason is required and is kept in the audit log.", + ), + ); + continue; + }; + let filter = match export_filter(filter_value) { + Ok(filter) => filter, + Err(why) => { + response.not_created.append( + client_id, + SetError::invalid_properties() + .with_property(P::Filter) + .with_description(why), + ); + continue; + } + }; + + let (ids, total) = + entries::query(server.store(), &filter, 0, MAX_EXPORT_ENTRIES, true).await?; + if total > MAX_EXPORT_ENTRIES { + response.not_created.append( + client_id, + SetError::invalid_properties() + .with_property(P::Filter) + .with_description(format!( + "{total} entries match; one export holds {MAX_EXPORT_ENTRIES}. Narrow the search." + )), + ); + continue; + } + + // Recorded first: no export leaves without its record + record( + server, + access_token, + session, + Action::Export, + None, + None, + format!( + "Exported {} journal entries: {}", + ids.len(), + filter_text(&filter) + ), + Some(reason), + ) + .await?; + + let mut items = Vec::with_capacity(ids.len()); + for id in ids { + if let Some(entry) = entries::get(server.store(), id).await? { + let bytes = report_bytes(server, &entry).await?; + items.push((id, entry, bytes)); + } + } + let (bytes, count) = build_zip(&items)?; + let blob = server + .put_jmap_blob(access_token.account_id(), &bytes) + .await?; + + let mut created = Map::with_capacity(5); + created.insert_unchecked( + Key::Property(P::Id), + Value::Element(JournalEntryValue::Id(Id::new(store::write::now()))), + ); + created.insert_unchecked( + Key::Property(P::BlobId), + Value::Str(blob.to_string().into()), + ); + created.insert_unchecked( + Key::Property(P::Size), + Value::Number((bytes.len() as u64).into()), + ); + created.insert_unchecked( + Key::Property(P::Count), + Value::Number((count as u64).into()), + ); + created.insert_unchecked( + Key::Property(P::Sha256), + Value::Str(hex(&Sha256::digest(&bytes)).into()), + ); + response.created.insert(client_id, Value::Object(created)); + } + Ok(response) +} + +fn summary(chains: &[ChainReport]) -> String { + if chains.is_empty() { + return "The journal is empty.".into(); + } + chains + .iter() + .map(|chain| match (&chain.broken_at, &chain.reason) { + (Some(at), Some(reason)) => format!("node {}: broken at {at}: {reason}", chain.node), + _ => format!( + "node {}: {} entries and {} purged verified ({} to {})", + chain.node, chain.entries, chain.purged, chain.first_seq, chain.last_seq + ), + }) + .collect::>() + .join("; ") +} + +/// `inbuxa:JournalVerification/set`: create `{}` to recheck every node's +/// chain and every report against its entry (JR-6). Recorded, with what it +/// found. +pub async fn verification_set( + server: &Server, + access_token: &AccessToken, + session: &HttpSessionData, + mut request: SetRequest<'_, JournalVerification>, +) -> trc::Result> { + server_level(access_token)?; + let mut response = SetResponse::from_request(&request, server.core.jmap.set_max_objects)?; + for (id, _) in request.unwrap_update().into_valid() { + response.not_updated.append(id, SetError::forbidden()); + } + for id in request.unwrap_destroy().into_valid() { + response.not_destroyed.append(id, SetError::forbidden()); + } + for (client_id, _) in request.unwrap_create() { + let chains = entries::verify(server.store(), Some(server.blob_store())).await?; + let verified = chains.iter().all(|chain| chain.broken_at.is_none()); + let entry = server + .audit_append(&Record { + at: store::write::now() * 1000, + actor: server.audit_actor(access_token).await, + via: access_token.origin().cloned(), + remote_ip: Some(session.remote_ip), + action: Action::Verify, + target: Target { + kind: "inbuxa:JournalEntry".into(), + ..Default::default() + }, + changes: vec![], + details: Some(summary(&chains)), + reason: None, + outcome: if verified { + Outcome::success() + } else { + Outcome::refused("chainBroken", None) + }, + }) + .await + .ok(); + + let mut created = Map::with_capacity(3); + created.insert_unchecked( + Key::Property(P::Id), + Value::Element(JournalEntryValue::Id(Id::new( + entry.map_or(0, |entry| entry.to_u64()), + ))), + ); + created.insert_unchecked(Key::Property(P::Verified), Value::Bool(verified)); + created.insert_unchecked( + Key::Property(P::Chains), + json_to_value(serde_json::to_value(&chains).unwrap_or_default()), + ); + response.created.insert(client_id, Value::Object(created)); + } + Ok(response) +} + +#[cfg(test)] +mod tests { + use super::*; + + fn entry(queue_id: u64) -> Entry { + Entry { + queue_id, + at: 1_790_000_000, + direction: Direction::Outgoing, + sender: "alice@example.com".into(), + authenticated: true, + recipients: vec!["pay@bank.example".into()], + subject: "Q3, final".into(), + message_id: "".into(), + accounts: vec![], + tenants: vec![], + journals: vec![1], + held: false, + blob: String::new(), + size: 0, + sha256: String::new(), + expires_at: 0, + } + } + + #[test] + fn exports_list_every_report_and_what_was_missing() { + let items = vec![ + ( + EntryId { node: 1, seq: 1 }, + entry(0x1a), + Some(b"report one".to_vec()), + ), + (EntryId { node: 1, seq: 2 }, entry(0x1b), None), + ]; + let (bytes, count) = build_zip(&items).unwrap(); + assert_eq!(count, 1); + let mut zip = zip::ZipArchive::new(Cursor::new(bytes)).unwrap(); + let mut read = |name: &str| { + let mut out = String::new(); + std::io::Read::read_to_string(&mut zip.by_name(name).unwrap(), &mut out).unwrap(); + out + }; + let manifest = read("manifest.csv"); + assert!(manifest.contains("\"Q3, final\""), "{manifest}"); + assert!(manifest.contains(&hex(&Sha256::digest(b"report one")))); + assert!(read("exceptions.csv").contains("couldn't be read")); + let sums = read("manifest.sha256"); + assert!(sums.contains(&hex(&Sha256::digest(manifest.as_bytes())))); + } + + #[test] + fn export_filters_parse() { + let filter: Value<'_, P, JournalEntryValue> = json_to_value(serde_json::json!({ + "sender": "alice", "direction": "outgoing", "journalId": "b", + "after": "2026-09-01T00:00:00Z" + })); + let filter = export_filter(Some(filter)).unwrap(); + assert_eq!(filter.sender.as_deref(), Some("alice")); + assert_eq!(filter.direction, Some(Direction::Outgoing)); + assert_eq!(filter.journal_id, Some(1)); + assert!(filter.after.is_some()); + let bad: Value<'_, P, JournalEntryValue> = + json_to_value(serde_json::json!({"colour": "red"})); + assert!(export_filter(Some(bad)).is_err()); + } +} diff --git a/crates/jmap/src/inbuxa/mod.rs b/crates/jmap/src/inbuxa/mod.rs index 92be663..d058281 100644 --- a/crates/jmap/src/inbuxa/mod.rs +++ b/crates/jmap/src/inbuxa/mod.rs @@ -13,6 +13,7 @@ pub mod legal_hold; pub mod mail_rule; pub mod security_acceptance; pub mod journal; +pub mod journal_entry; pub mod held_message; pub mod dlp_settings; pub mod hold_export; diff --git a/docs/spec/features/journaling.md b/docs/spec/features/journaling.md index 4d844e1..b8d7c86 100644 --- a/docs/spec/features/journaling.md +++ b/docs/spec/features/journaling.md @@ -313,6 +313,27 @@ Phase 3 (`feature/journal-archive`): in one SMTP session; before, a second message in the same session kept the first one's route. +Phase 4 (`feature/journal-search`): + +- `inbuxa:JournalEntry/query` (after, before, sender, recipient, address, + direction, subject words, Message-ID, journal; newest first, pages of up + to 500) and `/get` (`report`, the whole journal report up to 10 MB of + text, only when asked for), with `sysJournalSearch`. +- Recording (JR-17) happens before anything is returned, and nothing is + returned if it can't be written: a search with its terms, a listing + once per call, each report read on its own (as `blobAccess`, the + action reads of someone's mail already use), each export with its + reason (`export`), each check (`verify`). No new audit actions, so an + older version reads every record. +- `inbuxa:JournalExport/set` builds the ZIP in the request, like the audit + log's export, rather than as a task: at most 10,000 reports and 1 GB, + and a search that matches more is refused with the count, to narrow. + The ZIP has the reports as `.eml`, `manifest.csv` with the envelope and + a SHA-256 per file, `exceptions.csv` for reports that couldn't be read, + and `manifest.sha256`. +- `inbuxa:JournalVerification/set` rechecks the chains and every report + against its entry, with `sysJournalGet`. + ## Known gaps - A message a person saves to Sent over IMAP, or sends through another diff --git a/resources/privacy/catalog.toml b/resources/privacy/catalog.toml index 5f46cd7..0c0c0a2 100644 --- a/resources/privacy/catalog.toml +++ b/resources/privacy/catalog.toml @@ -136,6 +136,36 @@ retention = "object-life" note = ["content"] acceptedBy = ["identifier"] +[object."inbuxa:JournalEntry"] +file = "inbuxa_journal_entry.rs" +default = "none" +whose = ["holder", "correspondent"] +where = ["data-store", "blob-store"] +scope = "server" +retention = { setting = "inbuxa:Journal.retentionDays" } +[object."inbuxa:JournalEntry".properties] +sender = ["identifier", "contact"] +recipients = ["identifier", "contact"] +subject = ["content"] +messageId = ["identifier"] +report = ["content", "identifier", "contact", "metadata"] + +[object."inbuxa:JournalExport"] +file = "inbuxa_journal_entry.rs" +default = "none" +whose = ["holder", "correspondent"] +where = ["blob-store"] +scope = "server" +retention = { setting = "x:Jmap.uploadTtl" } +[object."inbuxa:JournalExport".properties] +blobId = ["identifier", "contact", "content"] +filter = ["identifier", "contact"] +reason = ["content"] + +[object."inbuxa:JournalVerification"] +file = "inbuxa_journal_entry.rs" +default = "none" + [object."inbuxa:LegalHold"] file = "inbuxa_legal_hold.rs" default = "none" diff --git a/tests/src/system/journal.rs b/tests/src/system/journal.rs index 4c9ac4f..5a67745 100644 --- a/tests/src/system/journal.rs +++ b/tests/src/system/journal.rs @@ -18,7 +18,11 @@ use inbuxa_features::journal::{ entries::{self, Entry, EntryId}, report, }; -use registry::schema::structs::{Expression, MtaStageAuth}; +use registry::schema::{ + prelude::{ObjectType, Property}, + structs::{CustomRoles, Expression, MtaStageAuth, Role, UserRoles}, +}; +use registry::types::map::Map; use serde_json::{Value, json}; use std::str::FromStr; use store::{Deserialize, write::BatchBuilder}; @@ -682,6 +686,223 @@ pub async fn archive(test: &mut TestServer) { inbuxa_features::journal::invalidate(); } +/// Phase 4: searching, reading and exporting over JMAP, by a Compliance +/// Officer, each recorded; administrators set journals up but don't read +/// them; the chain check. +pub async fn search(test: &mut TestServer) { + println!("Running journal search tests..."); + let admin = test.account("admin@example.com"); + let sender = admin + .create_user_account( + "search-sender@example.com", + "search-sender-secret-7301", + "Search sender", + &[], + vec![], + ) + .await; + let (_, response) = call( + &sender, + "Identity/set", + json!({"create": {"i": {"name": "Sender", "email": "search-sender@example.com"}}}), + ) + .await; + let identity = response["created"]["i"]["id"].as_str().unwrap().to_string(); + let (_, response) = call( + &sender, + "Mailbox/set", + json!({"create": {"m": {"name": "Search drafts"}}}), + ) + .await; + let mailbox = response["created"]["m"]["id"].as_str().unwrap().to_string(); + let (_, response) = call( + &admin, + "inbuxa:Journal/set", + json!({"create": {"s": {"name": "Search sender", "enabled": true, "direction": "any", + "scope": {"accounts": [sender.id_string()]}, "retentionDays": 30}}}), + ) + .await; + let journal_id = response["created"]["s"]["id"] + .as_str() + .unwrap_or_else(|| panic!("{response}")) + .to_string(); + inbuxa_features::journal::invalidate(); + for subject in ["Budget draft", "Budget final", "Lunch"] { + let response = send( + &sender, + &identity, + &mailbox, + &["someone@elsewhere.org"], + &["someone@elsewhere.org"], + subject, + ) + .await; + assert!(response["created"].get("s").is_some(), "{response}"); + } + + // Administrators set journals up but don't read them + let (name, response) = call( + &admin, + "inbuxa:JournalEntry/query", + json!({"filter": {"sender": "search-sender"}}), + ) + .await; + assert_eq!(name, "error", "{response}"); + + // A Compliance Officer does + let mut officer_role = None; + for id in admin + .registry_query_ids( + ObjectType::Role, + Vec::<(&str, &str)>::new(), + Vec::<&str>::new(), + ) + .await + { + let role = admin.registry_get::(id).await; + if role.description == "Compliance Officer" && role.member_tenant_id.is_none() { + officer_role = Some(id); + } + } + let officer = admin + .create_user_account( + "journal-officer@example.com", + "journal-officer-secret-7302", + "Officer", + &[], + vec![], + ) + .await; + admin + .registry_update_object( + ObjectType::Account, + officer.id(), + json!({Property::Roles: UserRoles::Custom(CustomRoles { + role_ids: Map::new(vec![officer_role.expect("the officer role")]), + })}), + ) + .await; + let (_, response) = call( + &officer, + "inbuxa:JournalEntry/query", + json!({"filter": {"sender": "search-sender", "text": "budget"}, "calculateTotal": true}), + ) + .await; + assert_eq!(response["total"], 2, "{response}"); + let ids = response["ids"].clone(); + let (_, response) = call(&officer, "inbuxa:JournalEntry/get", json!({"ids": ids})).await; + let list = response["list"].as_array().unwrap(); + assert_eq!(list.len(), 2, "{response}"); + assert_eq!(list[0]["subject"], "Budget final", "newest first"); + assert_eq!(list[0]["direction"], "outgoing"); + assert_eq!(list[0]["journalIds"], json!([journal_id])); + assert!(list[0]["report"].is_null(), "only when asked for"); + let first = list[0]["id"].as_str().unwrap().to_string(); + let (_, response) = call( + &officer, + "inbuxa:JournalEntry/get", + json!({"ids": [first.clone()], "properties": ["subject", "report"]}), + ) + .await; + let report = response["list"][0]["report"].as_str().unwrap_or_default(); + assert!( + report.contains("Subject: Journal report: Budget final"), + "{response}" + ); + assert!(report.contains("Sender: search-sender@example.com\r\n")); + let (_, response) = call( + &officer, + "inbuxa:JournalEntry/query", + json!({"filter": {"journalId": journal_id, "direction": "incoming"}}), + ) + .await; + assert_eq!(response["ids"], json!([]), "{response}"); + let (name, _) = call( + &officer, + "inbuxa:JournalEntry/query", + json!({"filter": {"colour": "red"}}), + ) + .await; + assert_eq!(name, "error"); + + // Exports need a reason, and hold every report the filter matches + let (_, response) = call( + &officer, + "inbuxa:JournalExport/set", + json!({"create": {"x": {"filter": {"sender": "search-sender"}}}}), + ) + .await; + assert_eq!( + response["notCreated"]["x"]["properties"], + json!(["reason"]), + "{response}" + ); + let (_, response) = call( + &officer, + "inbuxa:JournalExport/set", + json!({"create": {"x": {"filter": {"sender": "search-sender"}, "reason": "Case 12"}}}), + ) + .await; + let export = &response["created"]["x"]; + assert_eq!(export["count"], 3, "{response}"); + assert!(export["blobId"].as_str().is_some()); + assert_eq!(export["sha256"].as_str().map(str::len), Some(64)); + + // The chain check, which the officer may run too + let (_, response) = call( + &officer, + "inbuxa:JournalVerification/set", + json!({"create": {"v": {}}}), + ) + .await; + assert_eq!(response["created"]["v"]["verified"], true, "{response}"); + + // The officer changes no journals + let (name, _) = call( + &officer, + "inbuxa:Journal/set", + json!({"destroy": [journal_id.clone()]}), + ) + .await; + assert_eq!(name, "error"); + + // Every search, listing, read, export and check is recorded + let (_, response) = call( + &admin, + "inbuxa:AuditEvent/query", + json!({"filter": {"targetKind": "inbuxa:JournalEntry", "actorId": officer.id_string()}}), + ) + .await; + let ids = response["ids"].clone(); + let (_, response) = call(&admin, "inbuxa:AuditEvent/get", json!({"ids": ids})).await; + let details: Vec = response["list"] + .as_array() + .unwrap() + .iter() + .map(|e| format!("{} {}", e["action"], e["details"])) + .collect(); + for expected in [ + "Searched the journal", + "Listed 2 journal entries", + "Read a journaled message from search-sender@example.com", + "Exported 3 journal entries", + "verify", + ] { + assert!( + details.iter().any(|d| d.contains(expected)), + "{expected}: {details:?}" + ); + } + + call( + &admin, + "inbuxa:Journal/set", + json!({"destroy": [journal_id]}), + ) + .await; + inbuxa_features::journal::invalidate(); +} + struct Raw(Vec); impl Deserialize for Raw { @@ -703,6 +924,7 @@ pub async fn journal_tests() { test.insert_account(admin); self::test(&mut test).await; self::archive(&mut test).await; + self::search(&mut test).await; if test.is_reset() { test.temp_dir.delete(); } -- 2.54.0