Compare commits

..
16 Commits
Author SHA1 Message Date
jcoffey-dev 7e06a3b1f6 Merge pull request 'Release 2026.9.29.1' (#123) from release/2026.9.29.1-pr into main
ci / fork-checks (push) Successful in 48s
publish / version (push) Successful in 11s
ci / build (push) Successful in 30m5s
publish / publish-amd64 (push) Successful in 32m52s
publish / release (push) Successful in 10s
publish / publish-arm64 (push) Successful in 43m6s
publish / binaries (push) Successful in 36s
publish / announce (push) Successful in 10s
2026-09-29 15:40:42 +00:00
jcoffey-dev e147206e82 Release 2026.9.29.1
ci / fork-checks (pull_request) Successful in 50s
ci / build (pull_request) Successful in 13m48s
2026-09-29 08:25:13 -07:00
jcoffey-dev ca6484c356 Honor the client registration override only in setup and recovery
The recovery administrator signs in before any OAuth client is
registered, so it needs to skip the registration check. Outside
bootstrap and recovery mode, every account now signs in through a
registered client and one of its redirect URIs.
2026-09-29 08:25:13 -07:00
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
35 changed files with 2654 additions and 85 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",
+11 -7
View File
@@ -270,13 +270,17 @@ impl ClientRegistrationHandler for Server {
false
};
// Check if the account is allowed to override client registration
if self
.access_token(account_id)
.await
.caused_by(trc::location!())?
.build()
.has_permission(Permission::OAuthClientOverride)
// Check if the account is allowed to override client registration.
// inbuxa: only while setting up or recovering, when the recovery
// administrator signs in before any client is registered (contract C-5)
let registry = self.registry();
if (registry.is_bootstrap_mode() || registry.is_recovery_mode())
&& self
.access_token(account_id)
.await
.caused_by(trc::location!())?
.build()
.has_permission(Permission::OAuthClientOverride)
{
return Ok(None);
}
@@ -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.1"
};
}
+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();
}