use std::time::Duration;
use serde_json::Value;
#[derive(Debug)]
pub enum DeliveryMode {
Now,
TurnEnd,
Quiet {
quiet_for: Duration,
max_delay: Duration,
},
}
pub(crate) const CONFIRM_CAP: Duration = Duration::from_secs(10);
pub(crate) const SENDER_LABEL_MAX_CHARS: usize = 64;
pub(crate) fn validate_sender_label(label: &str) -> Result<(), String> {
if label.contains(']') || label.contains('\n') || label.contains('\r') {
return Err("sender label must not contain ']' or newlines".into());
}
if label.chars().count() > SENDER_LABEL_MAX_CHARS {
return Err(format!(
"sender label must be at most {SENDER_LABEL_MAX_CHARS} chars"
));
}
Ok(())
}
pub(crate) async fn confirm_route(
target: uuid::Uuid,
cap: Duration,
) -> (&'static str, String, &'static str) {
use crate::brain::agent::service::session_routes::turn_probe;
let mid_turn = |t| turn_probe(t).is_some_and(|probe| probe());
if mid_turn(target) {
return (
"queued_pending_drain",
"Confirmed queued: the target is mid-turn; the message injects at its next \
tool-loop boundary."
.into(),
"mid_turn",
);
}
let deadline = tokio::time::Instant::now() + cap;
while tokio::time::Instant::now() < deadline {
tokio::time::sleep(Duration::from_millis(500)).await;
if mid_turn(target) {
return (
"woke",
"Confirmed end-to-end: the target was idle and has started a turn on the \
message."
.into(),
"wake_confirmed",
);
}
}
(
"delivered",
format!(
"Routed to session {target}, but no wake was observed within {}s — the \
target may be parked (channel not claimed since boot) or slow to pick the \
message up. Re-check via session_search before resending.",
cap.as_secs()
),
"unconfirmed",
)
}
pub(crate) fn resolve_mode(
mode: Option<&str>,
interrupt: Option<bool>,
delivery: Option<&Value>,
) -> Result<DeliveryMode, String> {
fn secs(parent: Option<&Value>, key: &str, default: u64) -> Result<Duration, String> {
match parent.and_then(|d| d.get(key)) {
None => Ok(Duration::from_secs(default)),
Some(v) => {
let n = v
.as_u64()
.ok_or_else(|| format!("delivery.{key} must be a non-negative integer"))?;
Ok(Duration::from_secs(n))
}
}
}
let resolved = match mode {
None => None,
Some(known @ ("now" | "turn-end")) => Some(known),
Some("quiet") => {
if interrupt == Some(true) {
return Err(
"delivery.mode 'quiet' and interrupt=true disagree — quiet defers, \
interrupt derails"
.into(),
);
}
let quiet_for = secs(delivery, "quiet_for_secs", 60)?;
let max_delay = secs(delivery, "max_delay_secs", 1800)?;
return Ok(DeliveryMode::Quiet {
quiet_for,
max_delay,
});
}
Some(other) => {
return Err(format!(
"delivery.mode '{other}' is not available yet — use 'now', 'turn-end' or 'quiet'"
));
}
};
match (resolved, interrupt) {
(Some("turn-end"), None | Some(true)) | (None, Some(true)) => Ok(DeliveryMode::TurnEnd),
(Some("now"), None | Some(false)) | (None, None | Some(false)) => Ok(DeliveryMode::Now),
_ => Err("delivery.mode and interrupt disagree — pass one, not both".into()),
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn sender_label_rejects_framing_breakers_and_overlong() {
assert!(validate_sender_label("ok label").is_ok());
let brk = validate_sender_label("bad]label").unwrap_err();
assert!(brk.contains("must not contain"), "got: {brk}");
let nl = validate_sender_label("bad\nlabel").unwrap_err();
assert!(nl.contains("must not contain"), "got: {nl}");
let long =
validate_sender_label("x".repeat(SENDER_LABEL_MAX_CHARS + 1).as_str()).unwrap_err();
assert!(long.contains("at most"), "got: {long}");
assert!(validate_sender_label("x".repeat(SENDER_LABEL_MAX_CHARS).as_str()).is_ok());
}
}