Compare commits

..
13 Commits
Author SHA1 Message Date
jcoffey-dev ffcfde0b5a Merge pull request 'Release 2026.9.29' (#121) from release/2026.9.29-pr into main
publish / version (push) Successful in 11s
ci / fork-checks (push) Successful in 1m6s
ci / build (push) Successful in 32m20s
publish / publish-amd64 (push) Successful in 39m12s
publish / release (push) Successful in 5s
publish / publish-arm64 (push) Successful in 45m10s
publish / binaries (push) Successful in 41s
publish / announce (push) Successful in 22s
2026-09-29 05:43:53 +00:00
jcoffey-dev f1f05db790 Release 2026.9.29
ci / fork-checks (pull_request) Successful in 46s
ci / build (pull_request) Successful in 4m19s
2026-09-28 22:38:59 -07:00
jcoffey-dev e1076a04b2 Merge pull request 'Journaling spec: built, and the console as built' (#120) from spec/journaling-built into main
ci / fork-checks (push) Successful in 32s
ci / build (push) Canceled after 19m51s
2026-09-29 05:23:58 +00:00
jcoffey-dev 8fc8d94bbc Merge pull request 'Journaling: a Journal link in Management › Compliance' (#119) from feature/journal-menu into main
ci / fork-checks (push) Canceled after 22s
ci / build (push) Canceled after 22s
2026-09-29 05:23:36 +00:00
jcoffey-dev 0c600a63fa Journaling spec: built, and the console as built
ci / fork-checks (pull_request) Successful in 19s
ci / build (pull_request) Successful in 8m0s
2026-09-28 22:09:07 -07:00
jcoffey-dev f78925b316 Journaling: a Journal link in Management › Compliance
ci / fork-checks (pull_request) Successful in 56s
ci / build (pull_request) Successful in 17m1s
The console's journal page (CustomComponent/Journal), after Data Loss
Prevention; the console shows it to those who may see journals.
2026-09-28 22:06:12 -07:00
jcoffey-dev 6ee7ba1b7e Merge pull request 'Logs: a total only when it's known, not the query cap' (#117) from fix/log-query-total into main
ci / fork-checks (push) Successful in 17s
ci / build (push) Canceled after 23m14s
2026-09-29 05:00:18 +00:00
jcoffey-dev a992caf810 Merge pull request 'Journaling: search, read and export over JMAP, and the chain check' (#118) from feature/journal-search into main
ci / fork-checks (push) Successful in 17s
ci / build (push) Canceled after 6m46s
2026-09-29 04:53:28 +00:00
jcoffey-dev daa486f7e7 Journaling: search, read and export over JMAP, and the chain check
ci / fork-checks (pull_request) Successful in 16s
ci / build (pull_request) Successful in 8m0s
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.
2026-09-28 21:45:09 -07:00
jcoffey-dev 9c29fb2bea Logs: a total only when it's known, not the query cap
ci / fork-checks (pull_request) Successful in 55s
ci / build (pull_request) Successful in 17m3s
2026-09-28 21:42:53 -07:00
jcoffey-dev abd5811420 Merge pull request 'Journaling: outside archives, and Journal it in mail flow rules' (#116) from feature/journal-archive into main
ci / fork-checks (push) Successful in 50s
ci / build (push) Canceled after 19m30s
2026-09-29 04:33:59 +00:00
jcoffey-dev 80051539d5 Merge pull request 'Security to-do list: accepted items, kept on the server' (#114) from feature/security-acceptances into main
ci / fork-checks (push) Canceled after 11s
ci / build (push) Canceled after 10s
2026-09-29 04:33:46 +00:00
jcoffey-dev 64550ebbd0 Journaling: outside archives, and Journal it in mail flow rules
ci / fork-checks (pull_request) Successful in 1m19s
ci / build (pull_request) Successful in 5m44s
Phase 3 of the journaling spec.

- A journal's destination: builtIn (true for journals stored before) and
  archiveAddress, at least one. Reports to an archive are queued from the
  empty sender, one per address, flagged so they're never journaled.
- A pending record per report. When the queue lets go of one without
  delivering it (refused, expired, deleted), it becomes its own entry in
  the built-in journal under the sending journals' retention, the
  journal's archiveFailures (count, last time, reason) goes up, and the
  audit log records it; if that can't be written it stays queued.
- Journal it: a rule action naming a journal, on mail flow rules and
  beside a DLP rule's block, warn or hold. A journal whose scope chooses
  nobody takes only what rules send it.
- The report lists recipients a rule added or redirected to under
  "Added by rule", by rule name.
- A rule's route is cleared between messages in one SMTP session, with the
  new journal marks; a second message used to keep the first one's route.

tests/src/system/journal.rs: destination validation, a rule-only journal
fed by a rule that also adds a recipient, an unreachable archive's report
kept in the built-in journal with the failure counted, a report delivered
to an archive here and not journaled itself.
2026-09-28 21:27:44 -07:00
34 changed files with 2643 additions and 78 deletions
+120
View File
@@ -0,0 +1,120 @@
/*
* SPDX-FileCopyrightText: 2026 Coffey Labs
*
* SPDX-License-Identifier: AGPL-3.0-only
*/
//! Reports on their way to an outside archive (JR-7). Keys, after `J`:
//!
//! - `o` + the report's queue id: what goes into the built-in journal if
//! the archive never takes the report, as JSON. Cleared once it's
//! delivered or kept.
//! - `w` + journal id (u32): how often that journal's archive didn't take a
//! report, and the last time and reason, for the console's warning.
use super::{FEATURE, Json, entries::Entry};
use serde::{Deserialize as SerdeDeserialize, Serialize as SerdeSerialize};
use store::{
SUBSPACE_INBUXA, Serialize, Store, ValueKey,
write::{AnyClass, BatchBuilder, ValueClass},
};
use trc::AddContext;
const KIND_PENDING: u8 = b'o';
const KIND_FAILURES: u8 = b'w';
/// A report queued to an archive.
#[derive(Debug, Clone, PartialEq, Eq, SerdeSerialize, SerdeDeserialize)]
#[serde(rename_all = "camelCase")]
pub struct Pending {
pub address: String,
/// The entry, should the archive not take it: its own, with the
/// sending journals' retention, whatever else the built-in journal has.
pub entry: Entry,
}
/// How a journal's archive has been taking its reports.
#[derive(Debug, Clone, Default, PartialEq, Eq, SerdeSerialize, SerdeDeserialize)]
#[serde(rename_all = "camelCase")]
pub struct Failures {
pub count: u64,
/// Seconds.
pub last_at: u64,
pub last_reason: String,
}
fn class(kind: u8, id: &[u8]) -> ValueClass {
let mut key = Vec::with_capacity(2 + id.len());
key.push(FEATURE);
key.push(kind);
key.extend_from_slice(id);
ValueClass::Any(AnyClass {
subspace: SUBSPACE_INBUXA,
key,
})
}
pub async fn set_pending(data: &Store, queue_id: u64, pending: &Pending) -> trc::Result<()> {
let mut batch = BatchBuilder::new();
batch.set(
class(KIND_PENDING, &queue_id.to_be_bytes()),
Json(pending).serialize()?,
);
data.write(batch.build_all())
.await
.caused_by(trc::location!())?;
Ok(())
}
pub async fn pending(data: &Store, queue_id: u64) -> trc::Result<Option<Pending>> {
Ok(data
.get_value::<Json<Pending>>(ValueKey::from(class(KIND_PENDING, &queue_id.to_be_bytes())))
.await
.caused_by(trc::location!())?
.map(|Json(pending)| pending))
}
pub async fn clear_pending(data: &Store, queue_id: u64) -> trc::Result<()> {
let mut batch = BatchBuilder::new();
batch.clear(class(KIND_PENDING, &queue_id.to_be_bytes()));
data.write(batch.build_all())
.await
.caused_by(trc::location!())?;
Ok(())
}
pub async fn failures(data: &Store, journal_id: u32) -> trc::Result<Failures> {
Ok(data
.get_value::<Json<Failures>>(ValueKey::from(class(
KIND_FAILURES,
&journal_id.to_be_bytes(),
)))
.await
.caused_by(trc::location!())?
.map(|Json(failures)| failures)
.unwrap_or_default())
}
/// Counts one report an archive didn't take, for each of `journals`.
pub async fn record_failure(
data: &Store,
journals: &[u32],
at: u64,
reason: &str,
) -> trc::Result<()> {
for journal_id in journals {
let mut failures = failures(data, *journal_id).await?;
failures.count += 1;
failures.last_at = at;
failures.last_reason = reason.chars().take(500).collect();
let mut batch = BatchBuilder::new();
batch.set(
class(KIND_FAILURES, &journal_id.to_be_bytes()),
Json(&failures).serialize()?,
);
data.write(batch.build_all())
.await
.caused_by(trc::location!())?;
}
Ok(())
}
+180
View File
@@ -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<u64>,
/// Before this time, in seconds.
#[serde(skip_serializing_if = "Option::is_none")]
pub before: Option<u64>,
/// Part of the sender's address, ignoring case.
#[serde(skip_serializing_if = "Option::is_none")]
pub sender: Option<String>,
/// Part of any recipient's address, ignoring case.
#[serde(skip_serializing_if = "Option::is_none")]
pub recipient: Option<String>,
/// Part of the sender's or any recipient's address.
#[serde(skip_serializing_if = "Option::is_none")]
pub address: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub direction: Option<Direction>,
/// Words that must all appear in the subject, ignoring case.
#[serde(skip_serializing_if = "Option::is_none")]
pub text: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub message_id: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub journal_id: Option<u32>,
}
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<EntryId>, 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: "[email protected]".into(),
authenticated: true,
recipients: vec!["[email protected]".into()],
subject: "Q3 figures, final".into(),
message_id: "<[email protected]>".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("[email protected]".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];
+83 -2
View File
@@ -16,6 +16,7 @@
//! with `J`; journals are `j` + id (u32), as JSON. There are few, so they're
//! read whole.
pub mod archive;
pub mod entries;
pub mod report;
@@ -128,6 +129,12 @@ pub struct Journal {
/// How long an entry this journal writes is kept. An entry keeps the
/// retention it was written with (JR-12).
pub retention_days: u32,
/// Whether entries go into the built-in journal (JR-5).
#[serde(default = "yes")]
pub built_in: bool,
/// An outside archive's journal address, sent each report (JR-7).
#[serde(default, skip_serializing_if = "Option::is_none")]
pub archive_address: Option<String>,
#[serde(default)]
pub created_by: String,
#[serde(default)]
@@ -164,13 +171,28 @@ impl Journal {
format!("Keep entries between {MIN_RETENTION_DAYS} and {MAX_RETENTION_DAYS} days."),
);
}
// Neither is a journal only rules send mail to (JR-10)
let chosen = self.scope.lists().iter().any(|list| !list.is_empty());
if self.scope.everyone == chosen {
if self.scope.everyone && chosen {
return invalid(
"scope",
"Journal everyone, or choose accounts, groups, domains or tenants; not both.",
);
}
if !self.built_in && self.archive_address.is_none() {
return invalid(
"builtIn",
"Keep entries in the built-in journal, send them to an archive, or both.",
);
}
if let Some(address) = &self.archive_address
&& !is_address(address)
{
return invalid(
"archiveAddress",
format!("\"{address}\" isn't an email address."),
);
}
if self.scope.lists().iter().any(|list| list.len() > MAX_LIST) {
return invalid("scope", format!("Choose at most {MAX_LIST} of each."));
}
@@ -179,6 +201,11 @@ impl Journal {
/// Whether this journal takes a message going `direction` with these
/// people here on either side.
/// Whether only rules send this journal mail (JR-10).
pub fn rules_only(&self) -> bool {
!self.scope.everyone && self.scope.lists().iter().all(|list| list.is_empty())
}
pub fn takes(&self, direction: Direction, members: &[Member]) -> bool {
self.enabled
&& self.direction.includes(direction)
@@ -186,6 +213,22 @@ impl Journal {
}
}
fn yes() -> bool {
true
}
/// An address an archive can be sent to: one `@`, something either side,
/// nothing that would break an envelope.
fn is_address(address: &str) -> bool {
address.len() <= 320
&& address.split_once('@').is_some_and(|(local, domain)| {
!local.is_empty() && domain.contains('.') && !domain.contains('@')
})
&& !address
.chars()
.any(|c| c.is_whitespace() || c.is_control() || matches!(c, '<' | '>' | ',' | ';'))
}
/// A value stored as JSON.
pub(crate) struct Json<T>(pub T);
@@ -347,6 +390,8 @@ mod tests {
direction: Direction::Any,
scope,
retention_days: 365,
built_in: true,
archive_address: None,
created_by: String::new(),
created_at: 0,
updated_at: 0,
@@ -372,7 +417,11 @@ mod tests {
.validate()
.is_ok()
);
assert!(journal(Scope::default()).validate().is_err());
// Nobody chosen: only rules send it mail
let rules_only = journal(Scope::default());
assert!(rules_only.validate().is_ok());
assert!(rules_only.rules_only());
assert!(!rules_only.takes(Direction::Any, &[member(3, vec![7])]));
let both = Scope {
everyone: true,
groups: vec![4],
@@ -381,6 +430,38 @@ mod tests {
assert_eq!(journal(both).validate().unwrap_err().property, "scope");
}
#[test]
fn destinations() {
let mut j = journal(Scope {
everyone: true,
..Default::default()
});
j.built_in = false;
assert_eq!(j.validate().unwrap_err().property, "builtIn");
j.archive_address = Some("[email protected]".into());
assert!(j.validate().is_ok());
for bad in [
"archive",
"a@b",
"a [email protected]",
"<[email protected]>",
"a@[email protected]",
] {
j.archive_address = Some(bad.into());
assert_eq!(
j.validate().unwrap_err().property,
"archiveAddress",
"{bad}"
);
}
// Stored before destinations existed: the built-in journal
let old: Journal = serde_json::from_str(
r#"{"name":"Old","direction":"any","scope":{"everyone":true},"retentionDays":30}"#,
)
.unwrap();
assert!(old.built_in && old.archive_address.is_none());
}
#[test]
fn retention_has_bounds() {
let mut j = journal(Scope {
+27 -1
View File
@@ -20,6 +20,8 @@ use sha2::{Digest, Sha256};
pub struct Recipient {
pub address: String,
pub orcpt: Option<String>,
/// The mail flow rule that added or redirected to it.
pub added_by: Option<String>,
}
/// What the queue knows about a message.
@@ -46,6 +48,8 @@ pub struct Fields {
pub bcc: Vec<String>,
/// A list's address, and its members among the recipients.
pub expanded: Vec<(String, Vec<String>)>,
/// A rule's name, and the recipients it added.
pub added: Vec<(String, Vec<String>)>,
}
/// One line's worth of a value: no line breaks, no control characters.
@@ -106,7 +110,12 @@ pub fn fields(envelope: &Envelope<'_>, original: &[u8]) -> Fields {
.as_deref()
.map(orcpt_address)
.filter(|via| !via.is_empty() && *via != address);
if header_to.contains(&address) {
if let Some(rule) = &rcpt.added_by {
match fields.added.iter_mut().find(|(name, _)| name == rule) {
Some((_, added)) => added.push(line(&rcpt.address)),
None => fields.added.push((line(rule), vec![line(&rcpt.address)])),
}
} else if header_to.contains(&address) {
fields.to.push(line(&rcpt.address));
} else if header_cc.contains(&address) {
fields.cc.push(line(&rcpt.address));
@@ -159,6 +168,9 @@ pub fn text(envelope: &Envelope<'_>, fields: &Fields) -> String {
for (list, members) in &fields.expanded {
field("Expanded", &format!("{list} -> {}", members.join(", ")));
}
for (rule, added) in &fields.added {
field("Added by rule", &format!("{rule} -> {}", added.join(", ")));
}
if envelope.held {
field("Held for review", "yes");
}
@@ -265,6 +277,7 @@ The figures.\r\n";
Recipient {
address: address.into(),
orcpt: orcpt.map(Into::into),
added_by: None,
}
}
@@ -337,6 +350,19 @@ The figures.\r\n";
assert_eq!(original(&report), Some(unterminated));
}
#[test]
fn rule_added_recipients_say_so() {
let mut copied = rcpt("[email protected]", None);
copied.added_by = Some("Copy finance".into());
let recipients = [rcpt("[email protected]", None), copied];
let env = envelope(&recipients);
let fields = fields(&env, ORIGINAL);
assert!(fields.bcc.is_empty(), "{fields:?}");
assert!(
text(&env, &fields).contains("Added by rule: Copy finance -> [email protected]\r\n")
);
}
#[test]
fn values_stay_on_one_line() {
let recipients = [rcpt("[email protected]", None)];
+73 -2
View File
@@ -95,6 +95,33 @@ pub(crate) mod jmap_ids {
}
}
/// One id in the same form.
pub(crate) mod jmap_id {
use serde::{Deserialize, Deserializer, Serializer, de::Error};
use std::str::FromStr;
use types::id::Id;
pub fn serialize<S: Serializer>(id: &u32, serializer: S) -> Result<S::Ok, S::Error> {
serializer.serialize_str(&Id::from(*id).to_string())
}
#[derive(Deserialize)]
#[serde(untagged)]
enum Either {
Text(String),
Number(u32),
}
pub fn deserialize<'de, D: Deserializer<'de>>(deserializer: D) -> Result<u32, D::Error> {
match Either::deserialize(deserializer)? {
Either::Number(n) => Ok(n),
Either::Text(text) => Id::from_str(&text)
.map(|id| id.document_id())
.map_err(|_| D::Error::custom(format!("\"{text}\" isn't an id"))),
}
}
}
/// A detector and the least it must find.
#[derive(Debug, Clone, PartialEq, Eq, SerdeSerialize, SerdeDeserialize)]
#[serde(rename_all = "camelCase")]
@@ -222,6 +249,11 @@ pub enum Action {
Route {
queue: String,
},
/// Journaling spec, JR-10: a copy into this journal, whatever its scope.
Journal {
#[serde(with = "jmap_id")]
journal: u32,
},
// DLP actions
Block {
notice: String,
@@ -311,10 +343,16 @@ impl Rule {
if self.direction != Direction::Outgoing {
return Err(invalid("direction", "DLP rules check outgoing mail only."));
}
if dlp_actions != 1 || self.actions.len() != 1 {
// One of block, warn or hold; journaling may go with it
if dlp_actions != 1
|| self
.actions
.iter()
.any(|a| !a.is_dlp() && !matches!(a, Action::Journal { .. }))
{
return Err(invalid(
"actions",
"A DLP rule has exactly one action: block, warn or hold.",
"A DLP rule has exactly one action: block, warn or hold, and may also journal the message.",
));
}
}
@@ -478,6 +516,7 @@ fn validate_action(action: &Action) -> Result<(), String> {
}
Action::Refuse { text: t } => text(t, "refusal text"),
Action::Route { queue } => text(queue, "queue"),
Action::Journal { .. } => Ok(()),
Action::Block { notice } | Action::Warn { notice } | Action::Hold { notice, .. } => {
text(notice, "notice")
}
@@ -619,6 +658,38 @@ mod tests {
}
}
#[test]
fn journal_action_goes_with_either_kind() {
let hold = Action::Hold {
notice: "Held.".into(),
notify_sender: false,
};
let journal = Action::Journal { journal: 3 };
assert!(
rule(Kind::Dlp, vec![hold.clone(), journal.clone()])
.validate()
.is_ok()
);
assert!(rule(Kind::Dlp, vec![journal.clone()]).validate().is_err());
assert!(
rule(
Kind::Dlp,
vec![hold, Action::PrefixSubject { text: "x".into() }]
)
.validate()
.is_err()
);
assert!(
rule(Kind::Transport, vec![journal.clone()])
.validate()
.is_ok()
);
let json = serde_json::to_value(&journal).unwrap();
assert_eq!(json, serde_json::json!({"type": "journal", "journal": "d"}));
let back: Action = serde_json::from_value(json).unwrap();
assert_eq!(back, journal);
}
#[test]
fn wire_format() {
let json = r#"{"name":"Cards","kind":"dlp","direction":"outgoing",
@@ -31,6 +31,12 @@ pub enum JournalProperty {
Scope,
/// How long an entry is kept; each keeps what it was written with.
RetentionDays,
/// Whether entries go into the built-in journal.
BuiltIn,
/// An outside archive's journal address.
ArchiveAddress,
/// Reports the archive didn't take: how many, when and why last.
ArchiveFailures,
CreatedBy,
CreatedAt,
UpdatedAt,
@@ -59,6 +65,9 @@ impl Property for JournalProperty {
JournalProperty::Direction => "direction",
JournalProperty::Scope => "scope",
JournalProperty::RetentionDays => "retentionDays",
JournalProperty::BuiltIn => "builtIn",
JournalProperty::ArchiveAddress => "archiveAddress",
JournalProperty::ArchiveFailures => "archiveFailures",
JournalProperty::CreatedBy => "createdBy",
JournalProperty::CreatedAt => "createdAt",
JournalProperty::UpdatedAt => "updatedAt",
@@ -77,6 +86,9 @@ impl JournalProperty {
b"direction" => JournalProperty::Direction,
b"scope" => JournalProperty::Scope,
b"retentionDays" => JournalProperty::RetentionDays,
b"builtIn" => JournalProperty::BuiltIn,
b"archiveAddress" => JournalProperty::ArchiveAddress,
b"archiveFailures" => JournalProperty::ArchiveFailures,
b"createdBy" => JournalProperty::CreatedBy,
b"createdAt" => JournalProperty::CreatedAt,
b"updatedAt" => JournalProperty::UpdatedAt,
@@ -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<Self> {
// 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<Self> {
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<Self, Self::Err> {
JournalEntryProperty::parse(s).ok_or(())
}
}
impl Element for JournalEntryValue {
type Property = JournalEntryProperty;
fn try_parse<P>(key: &Key<'_, Self::Property>, value: &str) -> Option<Self> {
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<A>(&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::<serde::de::IgnoredAny>()?;
}
);
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<A>(&mut self, key: &str, map: &mut A) -> Result<(), A::Error>
where
A: serde::de::MapAccess<'de>,
{
if key == "property" {
let value = map.next_value::<Cow<str>>()?;
*self = if value == "receivedAt" {
JournalComparator::ReceivedAt
} else {
JournalComparator::_T(value.into_owned())
};
} else {
let _ = map.next_value::<serde::de::IgnoredAny>()?;
}
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<Id> for JournalEntryValue {
fn from(id: Id) -> Self {
JournalEntryValue::Id(id)
}
}
impl JmapObjectId for JournalEntryValue {
fn as_id(&self) -> Option<Id> {
match self {
JournalEntryValue::Id(id) => Some(*id),
}
}
fn as_any_id(&self) -> Option<AnyId> {
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<Id> {
None
}
fn as_any_id(&self) -> Option<AnyId> {
None
}
fn as_id_ref(&self) -> Option<&str> {
None
}
fn try_set_id(&mut self, _: AnyId) -> bool {
false
}
}
+1
View File
@@ -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
+3
View File
@@ -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)
}
@@ -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)?
}
+20 -1
View File
@@ -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",
+4
View File
@@ -127,6 +127,7 @@ pub enum GetRequestMethod {
MailRule(Box<GetRequest<crate::object::inbuxa_mail_rule::MailRule>>),
SecurityAcceptance(Box<GetRequest<crate::object::inbuxa_security_acceptance::SecurityAcceptance>>),
Journal(Box<GetRequest<crate::object::inbuxa_journal::Journal>>),
JournalEntry(Box<GetRequest<crate::object::inbuxa_journal_entry::JournalEntry>>),
HeldMessage(Box<GetRequest<crate::object::inbuxa_held_message::HeldMessage>>),
HoldExport(Box<GetRequest<crate::object::inbuxa_hold_export::HoldExport>>),
ProtocolPolicy(Box<GetRequest<crate::object::inbuxa_protocol_policy::ProtocolPolicy>>),
@@ -169,6 +170,8 @@ pub enum SetRequestMethod<'x> {
Box<SetRequest<'x, crate::object::inbuxa_security_acceptance::SecurityAcceptance>>,
),
Journal(Box<SetRequest<'x, crate::object::inbuxa_journal::Journal>>),
JournalExport(Box<SetRequest<'x, crate::object::inbuxa_journal_entry::JournalExport>>),
JournalVerification(Box<SetRequest<'x, crate::object::inbuxa_journal_entry::JournalVerification>>),
HeldMessage(Box<SetRequest<'x, crate::object::inbuxa_held_message::HeldMessage>>),
HoldExport(Box<SetRequest<'x, crate::object::inbuxa_hold_export::HoldExport>>),
ProtocolPolicy(Box<SetRequest<'x, crate::object::inbuxa_protocol_policy::ProtocolPolicy>>),
@@ -203,6 +206,7 @@ pub enum QueryRequestMethod {
ShareNotification(Box<QueryRequest<ShareNotification>>),
Registry(Box<QueryRequest<Registry>>),
AuditEvent(Box<QueryRequest<crate::object::inbuxa_audit::AuditEvent>>),
JournalEntry(Box<QueryRequest<crate::object::inbuxa_journal_entry::JournalEntry>>),
}
#[derive(Debug)]
+28
View File
@@ -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),
+21
View File
@@ -114,6 +114,7 @@ pub enum GetResponseMethod {
MailRule(GetResponse<crate::object::inbuxa_mail_rule::MailRule>),
SecurityAcceptance(GetResponse<crate::object::inbuxa_security_acceptance::SecurityAcceptance>),
Journal(GetResponse<crate::object::inbuxa_journal::Journal>),
JournalEntry(GetResponse<crate::object::inbuxa_journal_entry::JournalEntry>),
HeldMessage(GetResponse<crate::object::inbuxa_held_message::HeldMessage>),
HoldExport(GetResponse<crate::object::inbuxa_hold_export::HoldExport>),
ProtocolPolicy(GetResponse<crate::object::inbuxa_protocol_policy::ProtocolPolicy>),
@@ -156,6 +157,8 @@ pub enum SetResponseMethod {
Box<SetResponse<crate::object::inbuxa_security_acceptance::SecurityAcceptance>>,
),
Journal(Box<SetResponse<crate::object::inbuxa_journal::Journal>>),
JournalExport(Box<SetResponse<crate::object::inbuxa_journal_entry::JournalExport>>),
JournalVerification(Box<SetResponse<crate::object::inbuxa_journal_entry::JournalVerification>>),
HeldMessage(Box<SetResponse<crate::object::inbuxa_held_message::HeldMessage>>),
HoldExport(Box<SetResponse<crate::object::inbuxa_hold_export::HoldExport>>),
Explanation(Box<SetResponse<crate::object::inbuxa_explanation::Explanation>>),
@@ -865,6 +868,24 @@ impl<'x> From<SetResponse<crate::object::inbuxa_mail_rule::MailRule>> for Respon
}
// inbuxa: journaling
impl<'x> From<GetResponse<crate::object::inbuxa_journal_entry::JournalEntry>> for ResponseMethod<'x> {
fn from(value: GetResponse<crate::object::inbuxa_journal_entry::JournalEntry>) -> Self {
ResponseMethod::Get(GetResponseMethod::JournalEntry(value))
}
}
impl<'x> From<SetResponse<crate::object::inbuxa_journal_entry::JournalExport>> for ResponseMethod<'x> {
fn from(value: SetResponse<crate::object::inbuxa_journal_entry::JournalExport>) -> Self {
ResponseMethod::Set(SetResponseMethod::JournalExport(Box::new(value)))
}
}
impl<'x> From<SetResponse<crate::object::inbuxa_journal_entry::JournalVerification>> for ResponseMethod<'x> {
fn from(value: SetResponse<crate::object::inbuxa_journal_entry::JournalVerification>) -> Self {
ResponseMethod::Set(SetResponseMethod::JournalVerification(Box::new(value)))
}
}
impl<'x> From<GetResponse<crate::object::inbuxa_journal::Journal>> for ResponseMethod<'x> {
fn from(value: GetResponse<crate::object::inbuxa_journal::Journal>) -> Self {
ResponseMethod::Get(GetResponseMethod::Journal(value))
+20
View File
@@ -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!()
+32
View File
@@ -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)?;
+3
View File
@@ -433,6 +433,9 @@ impl IntermediateChangesResponse {
| MethodObject::MailRule
| MethodObject::SecurityAcceptance
| MethodObject::Journal
| MethodObject::JournalEntry
| MethodObject::JournalExport
| MethodObject::JournalVerification
| MethodObject::HeldMessage
| MethodObject::ProtocolPolicy
| MethodObject::TenantProtocolPolicy
+56 -11
View File
@@ -11,7 +11,10 @@
//! Changing or removing a journal never touches what it has taken.
use common::{Server, auth::AccessToken};
use inbuxa_features::journal::{self, Journal as Stored};
use inbuxa_features::journal::{
self, Journal as Stored,
archive::{self, Failures},
};
use jmap_proto::{
error::set::SetError,
method::{
@@ -37,13 +40,22 @@ const ALL: &[P] = &[
P::Direction,
P::Scope,
P::RetentionDays,
P::BuiltIn,
P::ArchiveAddress,
P::ArchiveFailures,
P::CreatedBy,
P::CreatedAt,
P::UpdatedAt,
];
/// Properties the server sets; a client that sends them is refused.
const SERVER_SET: &[P] = &[P::Id, P::CreatedBy, P::CreatedAt, P::UpdatedAt];
const SERVER_SET: &[P] = &[
P::Id,
P::ArchiveFailures,
P::CreatedBy,
P::CreatedAt,
P::UpdatedAt,
];
fn server_level(access_token: &AccessToken) -> trc::Result<()> {
if access_token.tenant_id().is_some() {
@@ -86,7 +98,7 @@ fn date(seconds: u64) -> JValue {
Value::Str(UTCDate::from_timestamp(seconds as i64).to_string().into())
}
fn to_value(journal: &Stored, properties: &[P]) -> JValue {
fn to_value(journal: &Stored, failures: &Failures, properties: &[P]) -> JValue {
let json = serde_json::to_value(journal).unwrap_or_default();
let mut out = Map::with_capacity(properties.len());
for property in properties {
@@ -94,6 +106,32 @@ fn to_value(journal: &Stored, properties: &[P]) -> JValue {
P::Id => Value::Element(JournalValue::Id(Id::from(journal.id))),
P::CreatedAt => date(journal.created_at),
P::UpdatedAt => date(journal.updated_at),
P::ArchiveAddress => journal
.archive_address
.as_ref()
.map_or(Value::Null, |a| Value::Str(a.clone().into())),
// JR-7: what the console warns about
P::ArchiveFailures => {
let mut out = Map::with_capacity(3);
out.insert_unchecked(Key::Borrowed("count"), Value::Number(failures.count.into()));
out.insert_unchecked(
Key::Borrowed("lastAt"),
if failures.count > 0 {
date(failures.last_at)
} else {
Value::Null
},
);
out.insert_unchecked(
Key::Borrowed("lastReason"),
if failures.count > 0 {
Value::Str(failures.last_reason.clone().into())
} else {
Value::Null
},
);
Value::Object(out)
}
other => json
.get(other.to_cow().as_ref())
.cloned()
@@ -162,21 +200,28 @@ pub async fn get(
not_found,
};
let journals = journal::all(server.store()).await?;
match ids {
None => {
response.list = journals
.iter()
.map(|journal| to_value(journal, &properties))
.collect()
}
let wanted: Vec<&Stored> = match ids {
None => journals.iter().collect(),
Some(ids) => {
let mut wanted = Vec::with_capacity(ids.len());
for id in ids {
match journal_id(id).and_then(|id| journals.iter().find(|j| j.id == id)) {
Some(journal) => response.list.push(to_value(journal, &properties)),
Some(journal) => wanted.push(journal),
None => response.push_not_found(id),
}
}
wanted
}
};
for journal in wanted {
let failures = if properties.contains(&P::ArchiveFailures) {
archive::failures(server.store(), journal.id).await?
} else {
Failures::default()
};
response
.list
.push(to_value(journal, &failures, &properties));
}
Ok(response)
}
+791
View File
@@ -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<String>,
target_name: Option<String>,
details: String,
reason: Option<String>,
) -> 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<Option<Vec<u8>>> {
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<JournalEntry>,
) -> trc::Result<GetResponse<JournalEntry>> {
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<u64, String> {
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<Direction, String> {
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<QueryFilter<JournalFilter>>) -> trc::Result<Filter> {
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<JournalEntry>,
) -> trc::Result<QueryResponse> {
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<Value<'_, P, JournalEntryValue>>) -> Result<Filter, String> {
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<Vec<u8>>)],
) -> trc::Result<(Vec<u8>, 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<SetResponse<JournalExport>> {
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::<String>())
}
_ => {
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::<Vec<_>>()
.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<SetResponse<JournalVerification>> {
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: "[email protected]".into(),
authenticated: true,
recipients: vec!["[email protected]".into()],
subject: "Q3, final".into(),
message_id: "<[email protected]>".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());
}
}
+1
View File
@@ -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;
+9
View File
@@ -146,6 +146,15 @@ pub(crate) async fn log_query(
})?;
response.anchor_found = true;
// inbuxa: the total is only known when the first page reached the end
// of the logs; counting them all would mean reading every file on every
// page. Upstream answered the query cap (5000) as the total, so a
// two-line log read "of 5000".
response.response.total = (req.request.calculate_total.unwrap_or(false)
&& anchor == 0
&& response.response.ids.len() < limit)
.then_some(response.response.ids.len());
Ok(response)
}
+8
View File
@@ -102,6 +102,10 @@ pub struct SessionData {
pub dlp_refusal: Option<DlpRefusal>,
// inbuxa: a mail flow rule's route for this message
pub mailflow_queue: Option<String>,
// inbuxa: journaling (JR-3, JR-10): journals rules sent this message
// to, and recipients rules added, by rule name
pub journal_marks: Vec<u32>,
pub journal_added: Vec<(String, String)>,
}
/// inbuxa: a DATA refusal by DLP rules: blocked, or a warning the sender
@@ -189,6 +193,8 @@ impl SessionData {
dlp_override: None,
dlp_refusal: None,
mailflow_queue: None,
journal_marks: Vec::new(),
journal_added: Vec::new(),
}
}
}
@@ -315,6 +321,8 @@ impl SessionData {
dlp_override: None,
dlp_refusal: None,
mailflow_queue: None,
journal_marks: Vec::new(),
journal_added: Vec::new(),
}
}
}
+15 -3
View File
@@ -763,19 +763,27 @@ impl<T: SessionStream> Session<T> {
}
for change in envelope {
match change {
super::mailflow::EnvelopeChange::AddRecipient(address) => {
super::mailflow::EnvelopeChange::AddRecipient(address, rule) => {
if !self
.data
.rcpt_to
.iter()
.any(|r| r.address_lcase.eq_ignore_ascii_case(&address))
{
self.data.journal_added.push((address.to_lowercase(), rule));
self.data.rcpt_to.push(SessionAddress::new(address));
}
}
super::mailflow::EnvelopeChange::Redirect(addresses) => {
super::mailflow::EnvelopeChange::Redirect(addresses, rule) => {
self.data.journal_added = addresses
.iter()
.map(|a| (a.to_lowercase(), rule.clone()))
.collect();
self.data.rcpt_to = addresses.into_iter().map(SessionAddress::new).collect();
}
super::mailflow::EnvelopeChange::Journal(journal) => {
self.data.journal_marks.push(journal);
}
super::mailflow::EnvelopeChange::Route(queue) => {
self.data.mailflow_queue = Some(queue);
}
@@ -882,7 +890,11 @@ impl<T: SessionStream> Session<T> {
.with_dkim_signers(dkim_signers)
.with_original_raw_message(original_message)
.with_original_authenticated_message(auth_message)
.with_metadata(metadata),
.with_metadata(metadata)
.with_journal(
std::mem::take(&mut self.data.journal_marks),
std::mem::take(&mut self.data.journal_added),
),
)
.await
{
+28 -5
View File
@@ -63,9 +63,12 @@ pub struct HeldDraft {
/// What a transport rule changes about where a message goes.
pub enum EnvelopeChange {
AddRecipient(String),
Redirect(Vec<String>),
/// An address, and the rule that added it.
AddRecipient(String, String),
Redirect(Vec<String>, String),
Route(String),
/// Journaling spec, JR-10: a journal the message goes to.
Journal(u32),
}
/// `[override: reason]` at the start of a subject: the reason, and the
@@ -402,11 +405,17 @@ impl<T: SessionStream> Session<T> {
rewrite::prefix_subject(now, text)
}
RuleAction::AddRecipient { address } => {
changes.push(EnvelopeChange::AddRecipient(address.clone()));
changes.push(EnvelopeChange::AddRecipient(
address.clone(),
matched.name.clone(),
));
None
}
RuleAction::Redirect { addresses } => {
changes.push(EnvelopeChange::Redirect(addresses.clone()));
changes.push(EnvelopeChange::Redirect(
addresses.clone(),
matched.name.clone(),
));
None
}
RuleAction::Route { queue } => {
@@ -423,7 +432,8 @@ impl<T: SessionStream> Session<T> {
}
RuleAction::Block { .. }
| RuleAction::Warn { .. }
| RuleAction::Hold { .. } => None,
| RuleAction::Hold { .. }
| RuleAction::Journal { .. } => None,
};
if next.is_some() {
current = next;
@@ -450,6 +460,19 @@ impl<T: SessionStream> Session<T> {
.await;
}
}
// JR-10: journals any matched rule sends the message to,
// DLP rules included
for matched in &outcome.matched {
for action in &matched.actions {
if let RuleAction::Journal { journal } = action
&& !changes
.iter()
.any(|c| matches!(c, EnvelopeChange::Journal(j) if j == journal))
{
changes.push(EnvelopeChange::Journal(*journal));
}
}
}
match hold {
Some(draft) => Checked::Hold {
draft,
+4
View File
@@ -494,6 +494,10 @@ impl<T: AsyncWrite + AsyncRead + Unpin> Session<T> {
self.data.delivery_by = 0;
self.data.future_release = 0;
self.data.rcpt_oks = 0;
// inbuxa: what mail flow rules decided was for the last message only
self.data.mailflow_queue = None;
self.data.journal_marks.clear();
self.data.journal_added.clear();
}
pub fn reset_tls(&mut self) {
+165 -44
View File
@@ -4,16 +4,22 @@
* SPDX-License-Identifier: AGPL-3.0-only
*/
//! inbuxa: journaling (journaling spec, JR-1 to JR-5, JR-11): the copy
//! taken as a message is queued, after DLP and transport rules, so it has
//! the envelope the message actually leaves or arrives with.
//! inbuxa: journaling (journaling spec, JR-1 to JR-11): the copy taken as a
//! message is queued, after DLP and transport rules, so it has the envelope
//! the message actually leaves or arrives with; reports to outside archives,
//! and what happens when an archive doesn't take one.
use crate::queue::{FROM_AUTHENTICATED, FROM_AUTOGENERATED, FROM_DSN, FROM_REPORT, Message};
use crate::queue::{
FROM_AUTHENTICATED, FROM_AUTOGENERATED, FROM_DSN, FROM_REPORT, Message, MessageSource, Status,
spool::{QueueParams, SmtpSpool},
};
use common::Server;
use inbuxa_features::{
audit::{Action, Actor, Outcome, Record, Target},
hold::Member,
journal::{
self, Direction,
archive::{self, Pending},
entries::{self, Entry},
report::{self, Envelope, Recipient},
},
@@ -22,6 +28,15 @@ use inbuxa_features::{
use store::write::{BatchBuilder, BlobLink, BlobOp, now};
use types::blob_hash::BlobHash;
/// What mail flow rules decided about a message at DATA (JR-3, JR-10).
#[derive(Debug, Clone, Default)]
pub struct Hints {
/// Journals a rule sent it to.
pub marks: Vec<u32>,
/// Recipients a rule added (lowercase), and the rule's name.
pub added: Vec<(String, String)>,
}
/// Marks a journal report the server queued itself, so it's never
/// journaled (JR-2). Free in the message flags (the MAIL parameters use
/// the low bits, the sources bits 32 to 37).
@@ -35,6 +50,7 @@ pub async fn capture(
queue_id: u64,
message: &Message,
raw: &[u8],
hints: &Hints,
) -> trc::Result<()> {
if message.flags & (FROM_JOURNAL | FROM_REPORT) != 0 {
return Ok(());
@@ -77,9 +93,10 @@ pub async fn capture(
}
}
let direction = Direction::of(sender_local, any_remote, any_local);
// A journal takes it through its scope, or because a rule sent it there
let taken: Vec<&journal::Journal> = journals
.iter()
.filter(|j| j.takes(direction, &members))
.filter(|j| j.takes(direction, &members) || hints.marks.contains(&j.id))
.collect();
if taken.is_empty() {
return Ok(());
@@ -98,6 +115,11 @@ pub async fn capture(
.map(|rcpt| Recipient {
address: rcpt.address.to_string(),
orcpt: rcpt.orcpt.as_deref().map(Into::into),
added_by: hints
.added
.iter()
.find(|(address, _)| address.eq_ignore_ascii_case(&rcpt.address))
.map(|(_, rule)| rule.clone()),
})
.collect();
let envelope = Envelope {
@@ -112,48 +134,147 @@ pub async fn capture(
let host = server.core.network.server_name.as_str();
let (bytes, fields) = report::build(&envelope, raw, &format!("postmaster@{host}"), host);
// The report's blob, reserved until the entry links it
let hash = BlobHash::generate(&bytes);
let mut batch = BatchBuilder::new();
batch.set(
BlobOp::Link {
hash: hash.clone(),
to: BlobLink::Temporary { until: at + 120 },
},
vec![],
);
server.store().write(batch.build_all()).await?;
server
.blob_store()
.put_blob(hash.as_slice(), &bytes, server.core.email.compression)
.await?;
let retention_days = taken
.iter()
.map(|j| j.retention_days)
.max()
.unwrap_or_default();
let mut tenants: Vec<u32> = members.iter().filter_map(|m| m.tenant).collect();
tenants.sort_unstable();
tenants.dedup();
let entry = Entry {
queue_id,
at,
direction,
sender: message.return_path.to_string(),
authenticated: envelope.authenticated,
recipients: recipients.iter().map(|r| r.address.clone()).collect(),
subject: fields.subject,
message_id: fields.message_id,
accounts: members.iter().map(|m| m.account).collect(),
tenants,
journals: taken.iter().map(|j| j.id).collect(),
held,
blob: entries::hex(hash.as_slice()),
size: bytes.len() as u64,
sha256: entries::sha256(&bytes),
expires_at: at + u64::from(retention_days) * 86_400,
let hash = BlobHash::generate(&bytes);
let entry_for = |journals: &[&journal::Journal]| {
let retention_days = journals
.iter()
.map(|j| j.retention_days)
.max()
.unwrap_or_default();
Entry {
queue_id,
at,
direction,
sender: message.return_path.to_string(),
authenticated: envelope.authenticated,
recipients: recipients.iter().map(|r| r.address.clone()).collect(),
subject: fields.subject.clone(),
message_id: fields.message_id.clone(),
accounts: members.iter().map(|m| m.account).collect(),
tenants: tenants.clone(),
journals: journals.iter().map(|j| j.id).collect(),
held,
blob: entries::hex(hash.as_slice()),
size: bytes.len() as u64,
sha256: entries::sha256(&bytes),
expires_at: at + u64::from(retention_days) * 86_400,
}
};
entries::append(server.store(), server.core.network.node_id, &entry).await?;
// The built-in journal: one entry, however many journals keep it there
let built_in: Vec<&journal::Journal> = taken.iter().copied().filter(|j| j.built_in).collect();
if !built_in.is_empty() {
// The report's blob, reserved until the entry links it
let mut batch = BatchBuilder::new();
batch.set(
BlobOp::Link {
hash: hash.clone(),
to: BlobLink::Temporary { until: at + 120 },
},
vec![],
);
server.store().write(batch.build_all()).await?;
server
.blob_store()
.put_blob(hash.as_slice(), &bytes, server.core.email.compression)
.await?;
entries::append(
server.store(),
server.core.network.node_id,
&entry_for(&built_in),
)
.await?;
}
// Outside archives: one report per address (JR-4, JR-7)
let mut addresses: Vec<(String, Vec<&journal::Journal>)> = Vec::new();
for journal in &taken {
if let Some(address) = &journal.archive_address {
let address = address.to_lowercase();
match addresses.iter_mut().find(|(a, _)| *a == address) {
Some((_, journals)) => journals.push(journal),
None => addresses.push((address, vec![journal])),
}
}
}
for (address, journals) in addresses {
// From nobody: an archive's refusal comes back to no one, and the
// queue's own record of it is what counts (settle, below)
let mut report = server.new_message("", MessageSource::Autogenerated, 0);
report.message.flags |= FROM_JOURNAL;
report.add_expanded_recipient(&address, server).await;
let pending = Pending {
address,
entry: entry_for(&journals),
};
archive::set_pending(server.store(), report.queue_id, &pending).await?;
let report_id = report.queue_id;
// Boxed: queueing the report comes back through this function
let queued = Box::pin(report.queue(QueueParams::new(&bytes, 0, server))).await;
if !queued {
archive::clear_pending(server.store(), report_id).await?;
return Err(trc::StoreEvent::UnexpectedError
.into_err()
.details("Failed to queue a journal report"));
}
}
Ok(())
}
/// JR-7: a journal report is leaving the queue. Delivered, its pending
/// record goes; not delivered (refused, expired, or deleted from the
/// queue), it goes into the built-in journal instead, and the journals
/// that sent it count a failure. An error means nothing was settled, and
/// the report must stay queued.
pub async fn settle(server: &Server, queue_id: u64, message: &Message) -> trc::Result<()> {
let store = server.store();
let Some(pending) = archive::pending(store, queue_id).await? else {
return Ok(());
};
let delivered = !message.recipients.is_empty()
&& message
.recipients
.iter()
.all(|rcpt| matches!(rcpt.status, Status::Completed(_)));
if !delivered {
let reason = if message
.recipients
.iter()
.any(|rcpt| matches!(rcpt.status, Status::PermanentFailure(_)))
{
"the archive refused it"
} else {
"it wasn't delivered before leaving the queue"
};
entries::append(store, server.core.network.node_id, &pending.entry).await?;
let at = now();
archive::record_failure(store, &pending.entry.journals, at, reason).await?;
server
.audit_note(Record {
at: at * 1000,
actor: Actor::system("Journal"),
via: None,
remote_ip: None,
action: Action::Create,
target: Target {
kind: "inbuxa:JournalEntry".into(),
id: Some(format!("{:x}", pending.entry.queue_id)),
name: None,
account_id: None,
tenant_id: None,
},
changes: vec![],
details: Some(format!(
"A journal report to {} wasn't delivered ({reason}); kept in the built-in journal",
pending.address
)),
reason: None,
outcome: Outcome::success(),
})
.await;
}
archive::clear_pending(store, queue_id).await
}
+39
View File
@@ -370,6 +370,8 @@ pub(crate) struct QueueParams<'x, 'y> {
pub session_id: u64,
pub server: &'y Server,
pub train_spam: Option<(bool, String)>,
// inbuxa: journaling, JR-3, JR-10
pub journal: crate::queue::journal::Hints,
}
impl MessageWrapper {
@@ -390,6 +392,7 @@ impl MessageWrapper {
server,
train_spam,
metadata,
journal,
..
} = params;
let event = self.message.queued_event();
@@ -462,6 +465,7 @@ impl MessageWrapper {
self.queue_id,
&self.message,
message.as_ref(),
&journal,
)
.await
{
@@ -813,6 +817,20 @@ impl MessageWrapper {
}
pub async fn remove(self, server: &Server, prev_event: Option<u64>) -> bool {
// inbuxa: journaling, JR-7: a journal report the archive never took
// goes into the built-in journal before it leaves the queue
if self.message.flags & crate::queue::journal::FROM_JOURNAL != 0
&& let Err(err) =
crate::queue::journal::settle(server, self.queue_id, &self.message).await
{
trc::error!(
err.details("Failed to settle a journal report; it stays queued.")
.span_id(self.span_id)
.caused_by(trc::location!())
);
return false;
}
let mut batch = BatchBuilder::new();
if let Some(prev_event) = prev_event {
@@ -987,6 +1005,19 @@ impl MessageWrapper {
server: &Server,
prev_events: AHashMap<QueueName, u64>,
) -> bool {
// inbuxa: journaling, JR-7, as in `remove`
if self.message.flags & crate::queue::journal::FROM_JOURNAL != 0
&& let Err(err) =
crate::queue::journal::settle(server, self.queue_id, &self.message).await
{
trc::error!(
err.details("Failed to settle a journal report; it stays queued.")
.span_id(self.span_id)
.caused_by(trc::location!())
);
return false;
}
let mut batch = BatchBuilder::new();
for (queue_name, due) in prev_events {
@@ -1142,9 +1173,17 @@ impl<'x, 'y> QueueParams<'x, 'y> {
original_raw_message: None,
original_authenticated_message: None,
metadata: Vec::new(),
journal: Default::default(),
}
}
/// inbuxa: journals mail flow rules sent the message to, and the
/// recipients they added, by rule name.
pub fn with_journal(mut self, marks: Vec<u32>, added: Vec<(String, String)>) -> Self {
self.journal = crate::queue::journal::Hints { marks, added };
self
}
pub fn with_train_spam(mut self, train_spam: Option<(bool, String)>) -> Self {
self.train_spam = train_spam;
self
+1 -1
View File
@@ -81,7 +81,7 @@ fn legacy_setting(name: &str, is_set: impl Fn(&str) -> bool) -> Option<String> {
#[macro_export]
macro_rules! brand_version {
() => {
"2026.9.28.5"
"2026.9.29"
};
}
+57 -3
View File
@@ -1,6 +1,7 @@
# Feature spec: journaling
Status: **approved 2026-09-28**, with the answers under [Settled](#settled).
Status: **approved 2026-09-28**, with the answers under [Settled](#settled);
**built 2026-09-29** (phases 2–5, see [As built](#as-built)), not yet released.
Not a rebuild of an upstream feature, so it has no line in SPEC.md §4's table.
Rule IDs: **JR-**.
@@ -290,8 +291,61 @@ fills in what it left open:
- **`inbuxa:JournalEntry`** (get, query) and **Check the journal** over
JMAP come in phase 4 with search, so every read is audited from the first
version that allows one. Phase 2 has `inbuxa:Journal` only.
- **Outside archives** (a journal's destination) come in phase 3; every
journal writes to the built-in journal until then.
Phase 3 (`feature/journal-archive`):
- **Destinations** are two properties of a journal: `builtIn` (true for
journals stored before phase 3) and `archiveAddress`. At least one.
- **Journals only rules use**: a journal whose scope chooses nobody takes
only what a **Journal it** action sends it. The action goes on mail flow
rules, and on a DLP rule beside its block, warn or hold (a blocked
message isn't queued, so it isn't journaled).
- **Reports to an archive** are queued from the empty sender, so a refusal
comes back to no one; a pending record per report says what to keep.
When the queue lets go of a report without delivering it (refused,
expired, or deleted from the queue), the report becomes its own entry in
the built-in journal under the sending journals' retention, even when
another journal already kept the message there, the journal's
`archiveFailures` (count, last time, reason) goes up, and the audit log
records it. If that can't be written, the report stays queued.
- **Added by rule** lists recipients a transport rule added or redirected
to, by rule name, instead of counting them as Bcc.
- A rule's route (and now its journal marks) is cleared between messages
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`.
Phase 5 (inbuxa-admin #62, server #119 for the menu, docs inbuxa.org #32):
- One console page, **Management › Compliance › Journal**, with two tabs
instead of the two pages §3 named: **Search** (for those who may search)
and **Journals** (the editor, on/off, delete, archive warnings, and Check
the journal).
- The warnings §3 put on the Overview (undelivered archive reports, a
node-local blob store) aren't there: archive failures show on each
journal, and there's no blob-store warning yet.
- The rule editor's **Journal it**, on mail flow rules and as an optional
second action on DLP rules.
## Known gaps
Binary file not shown.
+30
View File
@@ -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"
Binary file not shown.
+1 -1
View File
@@ -1 +1 @@
0Ay1tm9k_D94V7sLFfK7pIxwt3QdfYK_VBNs-rLGjhs
D8e0s1e4Umau4gRh5MEW24KtsawKGPNvS-6LWrIFbtQ
+464 -4
View File
@@ -18,8 +18,13 @@ 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};
const USING: &[&str] = &[
@@ -27,6 +32,7 @@ const USING: &[&str] = &[
"urn:ietf:params:jmap:mail",
"urn:ietf:params:jmap:submission",
"urn:inbuxa:jmap",
"urn:inbuxa:jmap:registry",
];
async fn call(account: &Account, method: &str, mut arguments: Value) -> (String, Value) {
@@ -171,8 +177,11 @@ pub async fn test(test: &mut TestServer) {
"both": {"name": "Both", "enabled": true, "direction": "any",
"scope": {"everyone": true, "accounts": [sender.id_string()]},
"retentionDays": 365},
"none": {"name": "None", "enabled": true, "direction": "any",
"scope": {}, "retentionDays": 365},
"none": {"name": "Nowhere", "enabled": true, "direction": "any",
"scope": {"everyone": true}, "retentionDays": 365, "builtIn": false},
"badaddr": {"name": "Bad archive", "enabled": true, "direction": "any",
"scope": {"everyone": true}, "retentionDays": 365,
"archiveAddress": "not an address"},
"server": {"name": "Mine", "enabled": true, "direction": "any",
"scope": {"everyone": true}, "retentionDays": 365,
"createdBy": "me"},
@@ -183,7 +192,7 @@ pub async fn test(test: &mut TestServer) {
}}),
)
.await;
for refused in ["short", "both", "none", "server"] {
for refused in ["short", "both", "none", "badaddr", "server"] {
assert_eq!(
response["notCreated"][refused]["type"], "invalidProperties",
"{refused}: {response}"
@@ -197,6 +206,14 @@ pub async fn test(test: &mut TestServer) {
response["notCreated"]["both"]["properties"],
json!(["scope"])
);
assert_eq!(
response["notCreated"]["none"]["properties"],
json!(["builtIn"])
);
assert_eq!(
response["notCreated"]["badaddr"]["properties"],
json!(["archiveAddress"])
);
let everything = response["created"]["all"]["id"]
.as_str()
.unwrap_or_else(|| panic!("{response}"))
@@ -445,6 +462,447 @@ pub async fn test(test: &mut TestServer) {
);
}
/// Phase 3: journals only rules send mail to, recipients a rule added,
/// reports sent to an outside archive, and what happens when the archive
/// doesn't take one.
pub async fn archive(test: &mut TestServer) {
println!("Running journal archive tests...");
let admin = test.account("[email protected]");
let sender = admin
.create_user_account(
"[email protected]",
"archive-sender-secret-7201",
"Archive sender",
&[],
vec![],
)
.await;
let vault = admin
.create_user_account(
"[email protected]",
"journal-vault-secret-7202",
"Journal vault",
&[],
vec![],
)
.await;
let (_, response) = call(
&sender,
"Identity/set",
json!({"create": {"i": {"name": "Sender", "email": "[email protected]"}}}),
)
.await;
let identity = response["created"]["i"]["id"].as_str().unwrap().to_string();
let (_, response) = call(
&sender,
"Mailbox/set",
json!({"create": {"m": {"name": "Archive drafts"}}}),
)
.await;
let mailbox = response["created"]["m"]["id"].as_str().unwrap().to_string();
let (_, response) = call(
&admin,
"inbuxa:Journal/set",
json!({"create": {
"rules": {"name": "Only what rules send", "enabled": true, "direction": "any",
"scope": {}, "retentionDays": 30},
"local": {"name": "To the vault", "enabled": true, "direction": "outgoing",
"scope": {"accounts": [sender.id_string()]}, "retentionDays": 30,
"builtIn": false, "archiveAddress": "[email protected]"},
"remote": {"name": "To an outside archive", "enabled": true, "direction": "internal",
"scope": {"accounts": [sender.id_string()]}, "retentionDays": 30,
"builtIn": false, "archiveAddress": "[email protected]"}
}}),
)
.await;
let id = |name: &str| {
response["created"][name]["id"]
.as_str()
.unwrap_or_else(|| panic!("{name}: {response}"))
.to_string()
};
let (rules_only, local, remote) = (id("rules"), id("local"), id("remote"));
let number = |id: &str| types::id::Id::from_str(id).unwrap().document_id();
let (_, response) = call(
&admin,
"inbuxa:MailRule/set",
json!({"create": {"r": {
"name": "Copy and journal", "kind": "transport", "direction": "outgoing",
"conditions": [{"type": "words", "words": ["journal-me"]}],
"actions": [
{"type": "addRecipient", "address": "[email protected]"},
{"type": "journal", "journal": rules_only.clone()}
]
}}}),
)
.await;
let rule = response["created"]["r"]["id"]
.as_str()
.unwrap_or_else(|| panic!("{response}"))
.to_string();
inbuxa_features::journal::invalidate();
// A rule sends it to a journal whose scope takes nobody, and says who
// it added
let response = send(
&sender,
&identity,
&mailbox,
&["[email protected]"],
&["[email protected]"],
"Marked journal-me",
)
.await;
assert!(response["created"].get("s").is_some(), "{response}");
let entry = all_entries(test)
.await
.into_iter()
.map(|(_, e)| e)
.find(|e| e.subject == "Marked journal-me" && e.journals.contains(&number(&rules_only)))
.expect("journaled by the rule");
assert_eq!(entry.journals, vec![number(&rules_only)], "{entry:?}");
let text = String::from_utf8_lossy(&report_of(test, &entry).await).into_owned();
assert!(
text.contains("Added by rule: Copy and journal -> [email protected]\r\n"),
"{text}"
);
assert!(!text.contains("Bcc:"), "{text}");
// The same message went to the outside archive, which can't be reached
// from here: once it leaves the queue (given up on, or deleted), it's
// kept in the built-in journal
let fallback = |entries: &[(EntryId, Entry)]| {
entries
.iter()
.any(|(_, e)| e.subject == "Marked journal-me" && e.journals == vec![number(&remote)])
};
let mut deleted = false;
for _ in 0..100 {
if fallback(&all_entries(test).await) {
break;
}
let (_, response) = call(&admin, "x:QueuedMessage/get", json!({"ids": null})).await;
if let Some(queued) = response["list"]
.as_array()
.unwrap()
.iter()
.find(|m| m.to_string().contains("[email protected]"))
{
let queued_id = queued["id"].as_str().unwrap().to_string();
let (_, response) = call(
&admin,
"x:QueuedMessage/set",
json!({"destroy": [queued_id.clone()]}),
)
.await;
deleted = response["destroyed"] == json!([queued_id]);
}
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
}
let kept: Vec<Entry> = all_entries(test)
.await
.into_iter()
.map(|(_, e)| e)
.filter(|e| e.subject == "Marked journal-me")
.collect();
assert_eq!(kept.len(), 2, "{kept:?}");
assert!(kept.iter().any(|e| e.journals == vec![number(&remote)]));
let (_, response) = call(
&admin,
"inbuxa:Journal/get",
json!({"ids": [remote.clone()]}),
)
.await;
let failures = &response["list"][0]["archiveFailures"];
assert_eq!(failures["count"], 1, "{response}");
assert_eq!(
failures["lastReason"],
if deleted {
"it wasn't delivered before leaving the queue"
} else {
"the archive refused it"
},
"{response}"
);
let (_, response) = call(
&admin,
"inbuxa:Journal/get",
json!({"ids": [local.clone()]}),
)
.await;
assert_eq!(response["list"][0]["archiveFailures"]["count"], 0);
// Delivered to an archive here: the report arrives, and nothing goes
// into the built-in journal for that journal
let response = send(
&sender,
&identity,
&mailbox,
&["[email protected]"],
&["[email protected]"],
"To the vault",
)
.await;
assert!(response["created"].get("s").is_some(), "{response}");
let mut arrived = Vec::new();
for _ in 0..100 {
let (_, response) = call(
&vault,
"Email/query",
json!({"filter": {"subject": "Journal report: To the vault"}}),
)
.await;
arrived = response["ids"].as_array().cloned().unwrap_or_default();
if !arrived.is_empty() {
break;
}
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
}
assert_eq!(arrived.len(), 1, "the report arrived");
assert!(entry_for(test, "To the vault").await.is_none());
assert!(
all_entries(test)
.await
.iter()
.all(|(_, e)| !e.subject.starts_with("Journal report")),
"reports aren't journaled"
);
let (_, response) = call(
&admin,
"inbuxa:Journal/get",
json!({"ids": [local.clone()]}),
)
.await;
assert_eq!(response["list"][0]["archiveFailures"]["count"], 0);
call(&admin, "inbuxa:MailRule/set", json!({"destroy": [rule]})).await;
call(
&admin,
"inbuxa:Journal/set",
json!({"destroy": [rules_only, local, remote]}),
)
.await;
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("[email protected]");
let sender = admin
.create_user_account(
"[email protected]",
"search-sender-secret-7301",
"Search sender",
&[],
vec![],
)
.await;
let (_, response) = call(
&sender,
"Identity/set",
json!({"create": {"i": {"name": "Sender", "email": "[email protected]"}}}),
)
.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,
&["[email protected]"],
&["[email protected]"],
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::<Role>(id).await;
if role.description == "Compliance Officer" && role.member_tenant_id.is_none() {
officer_role = Some(id);
}
}
let officer = admin
.create_user_account(
"[email protected]",
"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: [email protected]\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<String> = 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 [email protected]",
"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<u8>);
impl Deserialize for Raw {
@@ -465,6 +923,8 @@ pub async fn journal_tests() {
let admin = test.create_admin_account("[email protected]").await;
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();
}