use std::collections::{HashMap, VecDeque};
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
use affinidi_messaging_core::{ConnState, MessagingError, QueueFullGate};
use affinidi_tdk::messaging::ATM;
use affinidi_tdk::messaging::errors::ATMError;
use affinidi_tdk::messaging::profiles::ATMProfile;
use affinidi_tdk::messaging::protocols::message_pickup::{
MessagePickupStatusReply, UnprocessableMessage,
};
use base64::Engine;
use serde::Serialize;
use sha2::{Digest, Sha256};
use tokio::sync::{broadcast, watch};
use tracing::{debug, info, warn};
pub const CHECK_INTERVAL: Duration = Duration::from_secs(30);
pub const STALE_AFTER_SECS: u64 = 30;
pub const UNANSWERED_ALARM: u32 = 3;
pub const MAX_REDELIVERY_BACKOFF: Duration = Duration::from_secs(30 * 60);
const SAME_OLDEST_TOLERANCE_SECS: u64 = 2;
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
#[serde(rename_all = "camelCase")]
#[non_exhaustive]
pub enum InboxCollection {
#[default]
Unknown,
Collecting,
Backlogged,
NotDelivering,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
#[serde(rename_all = "camelCase")]
#[non_exhaustive]
pub struct InboxHealth {
pub state: InboxCollection,
pub waiting: Option<u32>,
pub longest_waited_secs: Option<u64>,
pub unanswered_status_requests: u32,
pub unproductive_redeliveries: u32,
pub last_answered_at: Option<u64>,
pub unprocessable_seen: u64,
pub unprocessable_deleted: u64,
}
#[derive(Clone, Default)]
pub struct InboxWatch {
inner: Arc<Mutex<InboxHealth>>,
}
impl InboxWatch {
pub fn new() -> Self {
Self::default()
}
pub fn snapshot(&self) -> InboxHealth {
self.inner.lock().expect("inbox health mutex").clone()
}
fn update(&self, f: impl FnOnce(&mut InboxHealth)) {
f(&mut self.inner.lock().expect("inbox health mutex"));
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct InboxStatus {
pub waiting: u32,
pub longest_waited_secs: Option<u64>,
pub oldest_received_time: Option<u64>,
}
impl From<&MessagePickupStatusReply> for InboxStatus {
fn from(s: &MessagePickupStatusReply) -> Self {
Self {
waiting: s.message_count,
longest_waited_secs: s.longest_waited_seconds,
oldest_received_time: s.oldest_received_time,
}
}
}
impl InboxStatus {
pub fn needs_catch_up(&self) -> bool {
self.waiting > 0
&& self
.longest_waited_secs
.is_none_or(|waited| waited >= STALE_AFTER_SECS)
}
fn oldest_arrival(&self, now: u64) -> Option<u64> {
self.oldest_received_time.or_else(|| {
self.longest_waited_secs
.map(|waited| now.saturating_sub(waited))
})
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Probe {
Answered(InboxStatus),
Refused,
Unanswered,
NotSent,
}
pub fn classify_status_error(err: &ATMError) -> Probe {
match err {
ATMError::MsgSendError(msg) if msg.contains("No response") => Probe::Unanswered,
ATMError::ProblemReport(..) | ATMError::MediatorError(..) => Probe::Refused,
_ => Probe::NotSent,
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Step {
Nothing,
Redeliver,
NotDelivering,
}
#[derive(Debug, Default)]
pub struct CatchUp {
unanswered: u32,
unproductive: u32,
asked_for_oldest: Option<u64>,
redeliver_not_before: u64,
stuck_reported: bool,
not_delivering_reported: bool,
}
pub fn redelivery_backoff(unproductive: u32) -> Duration {
if unproductive == 0 {
return Duration::ZERO;
}
let factor = 1u32.checked_shl(unproductive - 1).unwrap_or(u32::MAX);
CHECK_INTERVAL
.saturating_mul(factor)
.min(MAX_REDELIVERY_BACKOFF)
}
impl CatchUp {
pub fn step(&mut self, probe: Probe, now: u64) -> Step {
match probe {
Probe::NotSent => Step::Nothing,
Probe::Unanswered => {
self.unanswered = self.unanswered.saturating_add(1);
if self.unanswered >= UNANSWERED_ALARM {
Step::NotDelivering
} else {
Step::Nothing
}
}
Probe::Refused => {
self.unanswered = 0;
self.not_delivering_reported = false;
Step::Nothing
}
Probe::Answered(status) => {
self.unanswered = 0;
self.not_delivering_reported = false;
if !status.needs_catch_up() {
self.unproductive = 0;
self.asked_for_oldest = None;
self.redeliver_not_before = 0;
self.stuck_reported = false;
return Step::Nothing;
}
let oldest = status.oldest_arrival(now);
let same_as_last_ask = match (self.asked_for_oldest, oldest) {
(Some(asked), Some(now_oldest)) => {
asked.abs_diff(now_oldest) <= SAME_OLDEST_TOLERANCE_SECS
}
_ => false,
};
if now < self.redeliver_not_before {
return Step::Nothing;
}
if same_as_last_ask {
self.unproductive = self.unproductive.saturating_add(1);
} else {
self.unproductive = 0;
self.stuck_reported = false;
}
self.asked_for_oldest = oldest;
self.redeliver_not_before =
now.saturating_add(redelivery_backoff(self.unproductive).as_secs());
Step::Redeliver
}
}
}
pub fn unanswered(&self) -> u32 {
self.unanswered
}
pub fn unproductive(&self) -> u32 {
self.unproductive
}
fn take_stuck_report(&mut self) -> bool {
if self.unproductive >= 2 && !self.stuck_reported {
self.stuck_reported = true;
return true;
}
false
}
fn take_not_delivering_report(&mut self) -> bool {
if self.unanswered >= UNANSWERED_ALARM && !self.not_delivering_reported {
self.not_delivering_reported = true;
return true;
}
false
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum WatchExit {
NotDelivering,
}
fn unix_now() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0)
}
pub async fn run_inbox_watch(
atm: Arc<ATM>,
profile: Arc<ATMProfile>,
conn: Option<watch::Receiver<ConnState>>,
health: InboxWatch,
exit_when_not_delivering: bool,
) -> WatchExit {
let mut tick = tokio::time::interval(CHECK_INTERVAL);
tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
tick.tick().await;
let mut catch_up = CatchUp::default();
let alias = profile.inner.alias.clone();
loop {
tick.tick().await;
if conn
.as_ref()
.is_some_and(|c| *c.borrow() != ConnState::Connected)
{
health.update(|h| h.state = InboxCollection::Unknown);
continue;
}
let probe = match atm
.message_pickup()
.send_status_request(&profile, true, Some(Duration::from_secs(10)))
.await
{
Ok(Some(status)) => Probe::Answered(InboxStatus::from(&status)),
Ok(None) => Probe::Refused,
Err(e) => {
let probe = classify_status_error(&e);
debug!(profile = %alias, error = %e, ?probe, "inbox watch: status request failed");
probe
}
};
let was_not_delivering = catch_up.unanswered() >= UNANSWERED_ALARM;
let now = unix_now();
let step = catch_up.step(probe, now);
health.update(|h| {
h.unanswered_status_requests = catch_up.unanswered();
h.unproductive_redeliveries = catch_up.unproductive();
match probe {
Probe::Answered(status) => {
h.waiting = Some(status.waiting);
h.longest_waited_secs = status.longest_waited_secs;
h.last_answered_at = Some(now);
h.state = if status.needs_catch_up() {
InboxCollection::Backlogged
} else {
InboxCollection::Collecting
};
}
Probe::Refused => {
h.last_answered_at = Some(now);
if h.state == InboxCollection::NotDelivering {
h.state = InboxCollection::Unknown;
}
}
Probe::Unanswered if step == Step::NotDelivering => {
h.state = InboxCollection::NotDelivering;
}
_ => {}
}
});
if was_not_delivering && catch_up.unanswered() == 0 {
info!(
profile = %alias,
"the mediator is answering again — inbound messages are reaching this node"
);
}
match step {
Step::Nothing => {}
Step::NotDelivering => {
if catch_up.take_not_delivering_report() {
warn!(
profile = %alias,
unanswered = catch_up.unanswered(),
interval_secs = CHECK_INTERVAL.as_secs(),
"the mediator socket reports connected, but status requests go \
unanswered: nothing sent to this node is reaching it, and its peers \
will be refused `limits.queue.peer` once their queue to it fills"
);
}
if exit_when_not_delivering {
return WatchExit::NotDelivering;
}
}
Step::Redeliver => {
if let Probe::Answered(status) = probe {
if catch_up.take_stuck_report() {
warn!(
profile = %alias,
waiting = status.waiting,
longest_waited_secs = ?status.longest_waited_secs,
redeliveries = catch_up.unproductive(),
next_in_secs = redelivery_backoff(catch_up.unproductive()).as_secs(),
"messages are waiting in the mediator inbox that redelivery does not \
collect — most likely frames this node cannot unpack (see the \
'cannot unpack' warnings, which name their senders); backing off \
redelivery requests"
);
} else {
info!(
profile = %alias,
waiting = status.waiting,
longest_waited_secs = ?status.longest_waited_secs,
"mediator inbox holds messages live delivery did not bring; asking for \
redelivery"
);
}
}
if let Err(e) = atm
.message_pickup()
.toggle_live_delivery(&profile, true)
.await
{
warn!(
profile = %alias,
error = %e,
"inbox watch: could not ask the mediator to redeliver; retrying next tick"
);
}
}
}
}
}
pub const UNPROCESSABLE_DELETE_AFTER: u32 = 2;
const UNPROCESSABLE_CAPACITY: usize = 1024;
const UNPROCESSABLE_TTL: Duration = Duration::from_secs(60 * 60);
const DELETE_ENQUEUE_TIMEOUT: Duration = Duration::from_secs(5);
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct FrameDescription {
pub envelope: &'static str,
pub sender: Option<String>,
pub message_type: Option<String>,
}
fn b64_json(segment: &str) -> Option<serde_json::Value> {
let bytes = base64::engine::general_purpose::URL_SAFE_NO_PAD
.decode(segment.trim_end_matches('='))
.ok()?;
serde_json::from_slice(&bytes).ok()
}
fn str_field(value: &serde_json::Value, key: &str) -> Option<String> {
value.get(key).and_then(|v| v.as_str()).map(str::to_string)
}
pub fn describe_frame(raw: &str) -> FrameDescription {
let Ok(value) = serde_json::from_str::<serde_json::Value>(raw) else {
return FrameDescription {
envelope: "unknown",
..Default::default()
};
};
if value.get("ciphertext").is_some() {
let header = value
.get("protected")
.and_then(|p| p.as_str())
.and_then(b64_json)
.unwrap_or_default();
let alg = str_field(&header, "alg").unwrap_or_default();
let envelope = if alg.starts_with("ECDH-1PU") {
"authcrypt"
} else if alg.starts_with("ECDH-ES") {
"anoncrypt"
} else {
"unknown"
};
let sender = str_field(&header, "skid").or_else(|| {
str_field(&header, "apu").and_then(|apu| {
base64::engine::general_purpose::URL_SAFE_NO_PAD
.decode(apu.trim_end_matches('='))
.ok()
.and_then(|b| String::from_utf8(b).ok())
})
});
return FrameDescription {
envelope,
sender,
message_type: str_field(&header, "typ"),
};
}
if value.get("signatures").is_some() || value.get("signature").is_some() {
let first = value
.get("signatures")
.and_then(|s| s.as_array())
.and_then(|a| a.first())
.cloned()
.unwrap_or_else(|| value.clone());
let sender = first
.get("header")
.and_then(|h| str_field(h, "kid"))
.or_else(|| {
first
.get("protected")
.and_then(|p| p.as_str())
.and_then(b64_json)
.and_then(|h| str_field(&h, "kid"))
});
return FrameDescription {
envelope: "signed",
sender,
message_type: None,
};
}
if value.get("type").is_some() || value.get("body").is_some() {
return FrameDescription {
envelope: "plaintext",
sender: str_field(&value, "from"),
message_type: str_field(&value, "type"),
};
}
FrameDescription {
envelope: "unknown",
..Default::default()
}
}
pub fn frame_id(raw: &str) -> String {
hex::encode(Sha256::digest(raw.as_bytes()))
}
#[derive(Default)]
pub struct Sightings {
counts: HashMap<String, (u32, Instant)>,
order: VecDeque<String>,
}
impl Sightings {
pub fn record(&mut self, id: &str, now: Instant) -> u32 {
self.expire(now);
if let Some((count, _)) = self.counts.get_mut(id) {
*count = count.saturating_add(1);
return *count;
}
self.counts.insert(id.to_string(), (1, now));
self.order.push_back(id.to_string());
while self.order.len() > UNPROCESSABLE_CAPACITY {
if let Some(evicted) = self.order.pop_front() {
self.counts.remove(&evicted);
}
}
1
}
pub fn forget(&mut self, id: &str) {
self.counts.remove(id);
self.order.retain(|x| x != id);
}
fn expire(&mut self, now: Instant) {
while let Some(id) = self.order.front() {
match self.counts.get(id) {
Some((_, first)) if now.duration_since(*first) < UNPROCESSABLE_TTL => break,
_ => {
let id = self.order.pop_front().expect("front was just observed");
self.counts.remove(&id);
}
}
}
}
pub fn len(&self) -> usize {
self.counts.len()
}
pub fn is_empty(&self) -> bool {
self.counts.is_empty()
}
}
pub async fn run_unprocessable_quarantine(
atm: Arc<ATM>,
profile: Arc<ATMProfile>,
mut rx: broadcast::Receiver<UnprocessableMessage>,
health: InboxWatch,
) {
let mut sightings = Sightings::default();
loop {
let frame = match rx.recv().await {
Ok(frame) => frame,
Err(broadcast::error::RecvError::Lagged(skipped)) => {
warn!(
skipped,
"unprocessable-frame reports were dropped before they could be counted; \
those frames stay in the mediator inbox until they are redelivered"
);
continue;
}
Err(broadcast::error::RecvError::Closed) => return,
};
let id = frame
.attachment_id
.clone()
.unwrap_or_else(|| frame_id(&frame.raw));
let described = describe_frame(&frame.raw);
let seen = sightings.record(&id, Instant::now());
health.update(|h| h.unprocessable_seen = h.unprocessable_seen.saturating_add(1));
if seen < UNPROCESSABLE_DELETE_AFTER {
warn!(
frame = %id,
sender = described.sender.as_deref().unwrap_or("<not named>"),
envelope = described.envelope,
message_type = described.message_type.as_deref().unwrap_or("<encrypted>"),
reason = %frame.reason,
"cannot unpack an inbound message; leaving it at the mediator for one more \
delivery in case the failure is transient"
);
continue;
}
match tokio::time::timeout(
DELETE_ENQUEUE_TIMEOUT,
atm.delete_message_background(&profile, &id),
)
.await
{
Ok(Ok(())) => {
sightings.forget(&id);
health.update(|h| {
h.unprocessable_deleted = h.unprocessable_deleted.saturating_add(1)
});
warn!(
frame = %id,
deliveries = seen,
sender = described.sender.as_deref().unwrap_or("<not named>"),
envelope = described.envelope,
message_type = described.message_type.as_deref().unwrap_or("<encrypted>"),
reason = %frame.reason,
"cannot unpack an inbound message — deleting it from the mediator so it \
stops being redelivered and stops counting against its sender's queue"
);
}
Ok(Err(e)) => warn!(
frame = %id,
error = %e,
"could not delete an unprocessable inbound message; it will be redelivered"
),
Err(_) => warn!(
frame = %id,
timeout_secs = DELETE_ENQUEUE_TIMEOUT.as_secs(),
"timed out queueing the delete of an unprocessable inbound message; it will be \
redelivered"
),
}
}
}
pub const UNCOLLECTED_WINDOW: Duration = Duration::from_secs(5 * 60);
pub fn refused_recipient_not_collecting(err: &ATMError) -> bool {
err.http_status()
.and_then(|s| s.queue_full())
.is_some_and(|gate| gate == QueueFullGate::Peer)
}
pub fn messaging_refused_recipient_not_collecting(err: &MessagingError) -> bool {
err.queue_full() == Some(QueueFullGate::Peer)
}
#[derive(Default)]
pub struct UncollectedPeers {
marked: Mutex<HashMap<String, Instant>>,
}
impl UncollectedPeers {
pub fn new() -> Self {
Self::default()
}
pub fn record_refusal(&self, did: &str, now: Instant) -> bool {
let mut marked = self.marked.lock().expect("uncollected peers mutex");
marked.retain(|_, at| now.duration_since(*at) < UNCOLLECTED_WINDOW);
match marked.get(did) {
Some(_) => false,
None => {
marked.insert(did.to_string(), now);
true
}
}
}
pub fn record_accepted(&self, did: &str) -> bool {
self.marked
.lock()
.expect("uncollected peers mutex")
.remove(did)
.is_some()
}
pub fn is_marked(&self, did: &str, now: Instant) -> bool {
self.marked
.lock()
.expect("uncollected peers mutex")
.get(did)
.is_some_and(|at| now.duration_since(*at) < UNCOLLECTED_WINDOW)
}
pub fn observe_send(&self, did: &str, what: &str, refused_not_collecting: bool) -> bool {
if !refused_not_collecting {
return false;
}
if self.record_refusal(did, Instant::now()) {
warn!(
recipient = %did,
send = what,
window_secs = UNCOLLECTED_WINDOW.as_secs(),
"{did} is not collecting its mediator inbox; replies to it are being dropped"
);
} else {
debug!(
recipient = %did,
send = what,
"reply dropped: the recipient is still not collecting its mediator inbox"
);
}
true
}
pub fn observe_delivered(&self, did: &str) {
if self.record_accepted(did) {
info!(recipient = %did, "{did} is collecting its mediator inbox again");
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use affinidi_messaging_core::HttpStatusError;
fn status(waiting: u32, waited: Option<u64>, oldest: Option<u64>) -> Probe {
Probe::Answered(InboxStatus {
waiting,
longest_waited_secs: waited,
oldest_received_time: oldest,
})
}
#[test]
fn vti_50_only_a_message_live_delivery_missed_triggers_catch_up() {
let s = |waiting, waited| InboxStatus {
waiting,
longest_waited_secs: waited,
oldest_received_time: None,
};
assert!(!s(0, None).needs_catch_up(), "empty inbox");
assert!(!s(1, Some(2)).needs_catch_up(), "fresh push in flight");
assert!(s(1, Some(STALE_AFTER_SECS)).needs_catch_up(), "missed push");
assert!(s(1, None).needs_catch_up(), "age unreported: ask");
}
#[test]
fn vti_50_catch_up_runs_well_inside_a_short_message_expiry() {
assert!(CHECK_INTERVAL.as_secs() * 3 < 300);
}
#[test]
fn a_backlog_redelivery_does_not_clear_is_asked_for_less_and_less_often() {
let mut c = CatchUp::default();
let oldest = Some(1_000);
let mut now = 1_100;
let mut asks = Vec::new();
for _ in 0..200 {
if c.step(status(4, None, oldest), now) == Step::Redeliver {
asks.push(now);
}
now += CHECK_INTERVAL.as_secs();
}
let gaps: Vec<u64> = asks.windows(2).map(|w| w[1] - w[0]).collect();
assert!(
gaps.windows(2).all(|w| w[1] >= w[0]),
"gaps never shrink: {gaps:?}"
);
assert!(
*gaps.last().unwrap() >= MAX_REDELIVERY_BACKOFF.as_secs(),
"reaches the cap: {gaps:?}"
);
assert!(asks.len() < 20, "200 checks, {} asks", asks.len());
}
#[test]
fn a_new_backlog_is_asked_for_at_once() {
let mut c = CatchUp::default();
assert_eq!(c.step(status(1, None, Some(100)), 200), Step::Redeliver);
assert_eq!(c.step(status(1, None, Some(100)), 230), Step::Redeliver);
assert_eq!(c.unproductive(), 1);
assert_eq!(c.step(status(1, None, Some(240)), 290), Step::Redeliver);
assert_eq!(c.unproductive(), 0);
}
#[test]
fn a_drained_inbox_resets_the_backoff() {
let mut c = CatchUp::default();
for t in 0..5 {
c.step(status(1, None, Some(10)), 100 + t * 30);
}
assert!(c.unproductive() > 0);
assert_eq!(c.step(status(0, None, None), 400), Step::Nothing);
assert_eq!(c.unproductive(), 0);
assert_eq!(c.step(status(1, None, Some(390)), 430), Step::Redeliver);
}
#[test]
fn age_alone_identifies_the_same_oldest_message() {
let mut c = CatchUp::default();
assert_eq!(c.step(status(1, Some(40), None), 1_000), Step::Redeliver);
assert_eq!(c.step(status(1, Some(71), None), 1_030), Step::Redeliver);
assert_eq!(c.unproductive(), 1);
}
#[test]
fn three_unanswered_requests_while_connected_mean_not_delivering() {
let mut c = CatchUp::default();
assert_eq!(c.step(Probe::Unanswered, 0), Step::Nothing);
assert_eq!(c.step(Probe::Unanswered, 30), Step::Nothing);
assert_eq!(c.step(Probe::Unanswered, 60), Step::NotDelivering);
assert!(c.take_not_delivering_report());
assert_eq!(c.step(Probe::Unanswered, 90), Step::NotDelivering);
assert!(!c.take_not_delivering_report(), "reported once per episode");
assert_eq!(c.step(status(0, None, None), 120), Step::Nothing);
assert_eq!(c.unanswered(), 0);
}
#[test]
fn a_request_that_was_never_sent_is_not_evidence() {
let mut c = CatchUp::default();
for t in 0..10 {
assert_eq!(c.step(Probe::NotSent, t), Step::Nothing);
}
assert_eq!(c.unanswered(), 0);
}
#[test]
fn status_errors_are_classified_by_what_they_prove() {
assert_eq!(
classify_status_error(&ATMError::MsgSendError("No response from API".into())),
Probe::Unanswered
);
assert_eq!(
classify_status_error(&ATMError::TransportError(
"WebSocket message not transmitted: websocket is not connected".into()
)),
Probe::NotSent
);
assert_eq!(
classify_status_error(&ATMError::Disconnected("socket dropped".into())),
Probe::NotSent
);
assert_eq!(
classify_status_error(&ATMError::ProblemReport(
"e.p.x".into(),
"c".into(),
"".into()
)),
Probe::Refused
);
}
#[test]
fn redelivery_backoff_doubles_to_its_cap() {
assert_eq!(redelivery_backoff(0), Duration::ZERO);
assert_eq!(redelivery_backoff(1), CHECK_INTERVAL);
assert_eq!(redelivery_backoff(2), CHECK_INTERVAL * 2);
assert_eq!(redelivery_backoff(40), MAX_REDELIVERY_BACKOFF);
}
fn b64(v: &serde_json::Value) -> String {
base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(v.to_string())
}
#[test]
fn an_authcrypt_frame_names_its_sender_key() {
let protected = b64(&serde_json::json!({
"alg": "ECDH-1PU+A256KW",
"skid": "did:key:z6Mkexample#z6LSexample",
"typ": "application/didcomm-encrypted+json",
}));
let raw = serde_json::json!({"protected": protected, "ciphertext": "x"}).to_string();
let d = describe_frame(&raw);
assert_eq!(d.envelope, "authcrypt");
assert_eq!(d.sender.as_deref(), Some("did:key:z6Mkexample#z6LSexample"));
}
#[test]
fn an_authcrypt_frame_without_skid_names_its_sender_from_apu() {
let apu = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode("did:web:a#k");
let protected = b64(&serde_json::json!({"alg": "ECDH-1PU+A256KW", "apu": apu}));
let raw = serde_json::json!({"protected": protected, "ciphertext": "x"}).to_string();
assert_eq!(describe_frame(&raw).sender.as_deref(), Some("did:web:a#k"));
}
#[test]
fn an_anoncrypt_frame_names_no_sender() {
let protected = b64(&serde_json::json!({"alg": "ECDH-ES+A256KW"}));
let raw = serde_json::json!({"protected": protected, "ciphertext": "x"}).to_string();
let d = describe_frame(&raw);
assert_eq!(d.envelope, "anoncrypt");
assert_eq!(d.sender, None);
}
#[test]
fn a_signed_and_a_plaintext_frame_are_described() {
let jws = serde_json::json!({
"payload": "x",
"signatures": [{"header": {"kid": "did:key:z6Mk#k"}, "signature": "s"}],
})
.to_string();
let d = describe_frame(&jws);
assert_eq!(
(d.envelope, d.sender.as_deref()),
("signed", Some("did:key:z6Mk#k"))
);
let plain = serde_json::json!({
"id": "1", "type": "https://example/x", "from": "did:web:b", "body": {}
})
.to_string();
let d = describe_frame(&plain);
assert_eq!(d.envelope, "plaintext");
assert_eq!(d.sender.as_deref(), Some("did:web:b"));
assert_eq!(d.message_type.as_deref(), Some("https://example/x"));
assert_eq!(describe_frame("-ETSP...").envelope, "unknown");
}
#[test]
fn a_frame_id_is_the_hex_sha256_of_its_bytes() {
assert_eq!(
frame_id("abc"),
"ba7816bf8f01cfea414140de5dae2223b00361a396177a9cb410ff61f20015ad"
);
}
#[test]
fn sightings_count_and_stay_bounded() {
let mut s = Sightings::default();
let t = Instant::now();
assert_eq!(s.record("a", t), 1);
assert_eq!(s.record("a", t), 2);
s.forget("a");
assert_eq!(s.record("a", t), 1);
for i in 0..(UNPROCESSABLE_CAPACITY + 50) {
s.record(&format!("f{i}"), t);
}
assert!(s.len() <= UNPROCESSABLE_CAPACITY);
let later = t + UNPROCESSABLE_TTL + Duration::from_secs(1);
assert_eq!(s.record("f9999", later), 1);
assert_eq!(s.len(), 1);
}
fn peer_refusal() -> ATMError {
HttpStatusError::from_parts(
"send TSP message",
503,
None,
None,
r#"{"httpCode":503,"message":"{\"code\":\"e.p.limits.queue.peer\",\"comment\":\"Too many messages already waiting for this recipient\"}"}"#.to_string(),
)
.into()
}
#[test]
fn a_per_peer_refusal_means_the_recipient_is_not_collecting() {
assert!(refused_recipient_not_collecting(&peer_refusal()));
let sender_gate: ATMError = HttpStatusError::from_parts(
"send TSP message",
503,
None,
None,
r#"{"message":"{\"code\":\"e.p.limits.queue.sender\"}"}"#.to_string(),
)
.into();
assert!(!refused_recipient_not_collecting(&sender_gate));
assert!(!refused_recipient_not_collecting(
&ATMError::TransportError("connection refused".into())
));
let messaging = MessagingError::from(
peer_refusal()
.http_status()
.expect("an http status")
.clone(),
);
assert!(messaging_refused_recipient_not_collecting(&messaging));
}
#[test]
fn a_recipient_is_reported_once_per_window() {
let peers = UncollectedPeers::new();
let t = Instant::now();
assert!(peers.record_refusal("did:key:a", t), "first refusal warns");
assert!(!peers.record_refusal("did:key:a", t + Duration::from_secs(30)));
assert!(peers.record_refusal("did:key:b", t), "per recipient");
assert!(peers.is_marked("did:key:a", t + Duration::from_secs(60)));
assert!(!peers.is_marked("did:key:a", t + UNCOLLECTED_WINDOW));
assert!(peers.record_refusal("did:key:a", t + UNCOLLECTED_WINDOW + Duration::from_secs(1)));
assert!(peers.record_accepted("did:key:a"));
assert!(!peers.is_marked("did:key:a", t + UNCOLLECTED_WINDOW + Duration::from_secs(2)));
}
}