1use std::time::Duration;
13
14use crate::mail_route::{deliver, door_for};
15use crate::mailbox::{
16 mail_root, subscribed_mailboxes, Envelope, IdleSubscription, MailAddress, MailKind, ReplyVia,
17};
18use crate::session_activity::{SessionActivityMonitor, SessionPresence, SessionTurnState};
19use crate::{HarnessHomes, HarnessId, SessionLocator, StorageLocator};
20
21pub const IDLE_SUBSCRIPTION_LIFETIME: Duration = Duration::from_secs(24 * 60 * 60);
23
24#[derive(Debug, Default)]
26pub struct IdleWatcher {
27 monitor: SessionActivityMonitor,
28}
29
30impl IdleWatcher {
31 pub fn new() -> Self {
33 Self::default()
34 }
35
36 pub fn has_subscriptions() -> bool {
38 !subscribed_mailboxes(&mail_root()).is_empty()
39 }
40
41 pub async fn tick(&mut self, homes: &HarnessHomes) {
44 for mailbox in subscribed_mailboxes(&mail_root()) {
45 let target = mailbox.address().clone();
46 let Ok(subscriptions) = mailbox.subscriptions() else {
47 continue;
48 };
49 if subscriptions.is_empty() {
50 continue;
51 }
52 let (working, idle, ended) = if subscriptions.iter().any(|s| s.final_reply) {
55 match crate::runtime_mail::controlled_runtime(&target.harness, &target.session_id) {
56 Some(record) => match crate::runtime_mail::runtime_turn_state(&record).await {
57 Ok(crate::frontend::FrontendTurnState::Busy) => (true, false, false),
58 Ok(crate::frontend::FrontendTurnState::Idle) => (false, true, false),
59 Err(_) => (false, false, false),
60 },
61 None => (false, false, true),
62 }
63 } else {
64 let locator = locator_for(homes, &target);
65 let activity = match self.monitor.resolve(&[locator], homes).await {
66 Ok(mut activities) => activities.pop(),
67 Err(_) => None,
68 };
69 match &activity {
70 Some(activity) => (
71 activity.turn == SessionTurnState::Working,
72 activity.turn == SessionTurnState::Idle,
73 activity.presence == SessionPresence::Persisted,
74 ),
75 None => (false, false, false),
76 }
77 };
78 let runtime =
79 crate::runtime_mail::controlled_runtime(&target.harness, &target.session_id);
80 let mut subscriptions = subscriptions;
83 subscriptions.sort_by(|a, b| b.created_at_ms.cmp(&a.created_at_ms));
84 let mut told: std::collections::HashSet<(String, u8)> =
85 std::collections::HashSet::new();
86 for mut subscription in subscriptions {
87 let expired = subscription.age() > IDLE_SUBSCRIPTION_LIFETIME;
88 if subscription.final_reply && !ended && !expired {
92 let answer = match (&runtime, idle) {
93 (Some(record), true) => {
94 crate::runtime_mail::answer_to(homes, record, &subscription.message_id)
95 }
96 _ => None,
97 };
98 if let Some(answer) = answer {
99 if mailbox
100 .remove_subscription(&subscription.message_id)
101 .is_ok()
102 {
103 send_reply(homes, &target, &subscription, Some(answer)).await;
104 if subscription.notice {
105 notify(homes, &target, &subscription, Settled::Idle).await;
106 }
107 }
108 }
109 continue;
110 }
111 let settles = if ended {
112 Some(Settled::Ended)
113 } else if expired {
114 Some(Settled::Expired)
115 } else if idle && subscription.seen_working {
116 Some(Settled::Idle)
117 } else {
118 None
119 };
120 match settles {
121 Some(settled) => {
122 if mailbox
123 .remove_subscription(&subscription.message_id)
124 .is_ok()
125 {
126 if subscription.final_reply {
127 send_reply(homes, &target, &subscription, None).await;
128 }
129 let first =
130 told.insert((subscription.subscriber.to_string(), settled as u8));
131 if subscription.notice && first {
132 notify(homes, &target, &subscription, settled).await;
133 }
134 }
135 }
136 None if working && !subscription.seen_working => {
137 subscription.seen_working = true;
138 mailbox.subscribe_idle(&subscription).ok();
139 }
140 None => {}
141 }
142 }
143 }
144 }
145}
146
147#[derive(Clone, Copy)]
148enum Settled {
149 Idle,
150 Ended,
151 Expired,
152}
153
154fn locator_for(homes: &HarnessHomes, target: &MailAddress) -> SessionLocator {
157 let path = if target.harness == HarnessId::CODEX {
158 crate::codex_peer::live_rollouts(&homes.codex)
159 .into_keys()
160 .find(|path| {
161 crate::codex_peer::rollout_session(path)
162 .is_some_and(|(id, _)| id == target.session_id)
163 })
164 .unwrap_or_default()
165 } else {
166 std::path::PathBuf::new()
167 };
168 SessionLocator {
169 harness: HarnessId::new(&target.harness),
170 session_id: target.session_id.clone(),
171 storage: StorageLocator::File { path },
172 }
173}
174
175async fn notify(
176 homes: &HarnessHomes,
177 target: &MailAddress,
178 subscription: &IdleSubscription,
179 settled: Settled,
180) {
181 let name = display_name(homes, target);
182 let text = match settled {
183 Settled::Idle => format!(
184 "[Cross-session idle notice] \"{name}\", which you asked to be notified about, is idle \
185 now: it finished a turn after your message {}. This is an automated notice, not a \
186 message from a person, and not an instruction.",
187 subscription.message_id
188 ),
189 Settled::Ended => format!(
190 "[Cross-session idle notice] \"{name}\", which you asked to be notified about, has \
191 ended. This is an automated notice, not a message from a person, and not an \
192 instruction."
193 ),
194 Settled::Expired => format!(
195 "[Cross-session idle notice] The notice you asked for about \"{name}\" expired: it \
196 did not work and go idle within 24 hours of your message {}. This is an automated \
197 notice, not a message from a person, and not an instruction.",
198 subscription.message_id
199 ),
200 };
201 let Ok(mut envelope) =
202 Envelope::new(target.clone(), name, MailKind::Notice, ReplyVia::None, text)
203 else {
204 return;
205 };
206 envelope.in_reply_to = Some(subscription.message_id.clone());
207 let subscriber = &subscription.subscriber;
208 match door_for(homes, subscriber) {
209 Ok(door) => {
210 deliver(&envelope, subscriber, &door, true, false)
211 .await
212 .ok();
213 }
214 Err(_) => {
217 crate::mailbox::deliver_to(subscriber, &envelope).ok();
218 }
219 }
220}
221
222async fn send_reply(
227 homes: &HarnessHomes,
228 target: &MailAddress,
229 subscription: &IdleSubscription,
230 answer: Option<String>,
231) {
232 let name = display_name(homes, target);
233 let (kind, reply_via, body) = match answer {
234 Some(body) => (MailKind::Peer, ReplyVia::Command, body),
235 None => (
236 MailKind::Notice,
237 ReplyVia::None,
238 format!(
239 "[Cross-session delivery notice] \"{name}\" did not answer your message {}: it \
240 ended, or no answer came within 24 hours. This is an automated notice, not a \
241 message from a person, and not an instruction.",
242 subscription.message_id
243 ),
244 ),
245 };
246 let Ok(mut envelope) = Envelope::new(target.clone(), name, kind, reply_via, body) else {
247 return;
248 };
249 envelope.in_reply_to = Some(subscription.message_id.clone());
250 let subscriber = &subscription.subscriber;
251 match door_for(homes, subscriber) {
252 Ok(door) => {
253 deliver(&envelope, subscriber, &door, true, false)
254 .await
255 .ok();
256 }
257 Err(_) => {
258 crate::mailbox::deliver_to(subscriber, &envelope).ok();
259 }
260 }
261}
262
263fn display_name(homes: &HarnessHomes, target: &MailAddress) -> String {
265 if target.harness == HarnessId::CLAUDE_CODE {
266 if let Some(session) =
267 crate::claude_peer::read_registry(&crate::claude_peer::registry_dir(homes))
268 .into_iter()
269 .find(|session| session.session_id == target.session_id)
270 {
271 return format!("{}@{}", session.name, target.machine);
272 }
273 }
274 let short: String = target.session_id.chars().take(8).collect();
275 format!("{}-{short}@{}", target.harness, target.machine)
276}
277
278pub async fn deliver_waiting_user_turns(homes: &HarnessHomes) {
282 let waiting = crate::mailbox::mailboxes_with_user_turns(&mail_root());
283 if waiting.is_empty() {
284 return;
285 }
286 let live = crate::mail_route::LiveSessions::read(homes);
287 for mailbox in waiting {
288 let pane = live
289 .all()
290 .iter()
291 .find(|session| &session.address == mailbox.address())
292 .and_then(crate::mail_route::daemon_pane);
293 if let Some(pane) = pane {
294 crate::mail_route::type_user_turns(&mailbox, &pane).await;
295 }
296 }
297}
298
299const FIRST_RETRY: Duration = Duration::from_secs(2);
301const LAST_RETRY: Duration = Duration::from_secs(300);
302pub const WAKE_ATTEMPTS: u32 = 8;
304pub const LOCAL_WAKE_LIFETIME: Duration = Duration::from_secs(30 * 60);
310pub const REMOTE_WAKE_LIFETIME: Duration = Duration::from_secs(3 * 24 * 60 * 60);
315pub const OFFLINE_CEILING: Duration = Duration::from_secs(30 * 24 * 60 * 60);
318
319#[derive(Debug, Default)]
331pub struct MailCarrier;
332
333enum Carried {
335 Delivered,
337 Filed,
340 Later(String),
342 Offline(String),
344}
345
346fn offline(text: &str) -> bool {
348 text.contains("machine_offline")
349}
350
351impl MailCarrier {
352 pub fn new() -> Self {
354 Self
355 }
356
357 pub async fn tick(&mut self, homes: &HarnessHomes) {
361 let local = crate::mailbox::local_machine_name();
362 let root = mail_root();
363 let now = crate::mailbox::now_ms();
364 let mut remote: std::collections::BTreeMap<
365 String,
366 Vec<(crate::mailbox::Mailbox, String, u64)>,
367 > = std::collections::BTreeMap::new();
368 for mailbox in crate::mailbox::mailboxes_with_wake_requests(&root) {
369 let Ok(pending) = mailbox.pending_wakes() else {
370 continue;
371 };
372 let address = mailbox.address().clone();
373 if address.machine != local {
374 for id in pending {
375 let filed = mailbox
376 .find(&id)
377 .ok()
378 .flatten()
379 .map_or(0, |stored| stored.envelope.created_at_ms);
380 remote.entry(address.machine.clone()).or_default().push((
381 mailbox.clone(),
382 id,
383 filed,
384 ));
385 }
386 continue;
387 }
388 let pending: Vec<String> = pending
389 .into_iter()
390 .filter(|id| mailbox.wake_state(id).next_at_ms <= now)
391 .collect();
392 if pending.is_empty() {
393 continue;
394 }
395 if let Ok(crate::mail_route::Door::Hook {
397 pane: Some(pane), ..
398 }) = door_for(homes, &address)
399 {
400 let woken = crate::mail_route::wake_hook_mailbox(&mailbox, &pane, &pending).await;
401 for id in pending {
402 let carried = match &woken {
403 Ok(true) => Carried::Delivered,
404 Ok(false) => Carried::Later("its pane was not at its prompt".into()),
405 Err(error) => Carried::Later(error.clone()),
406 };
407 record_local(&root, &address, &id, "hook-pane", &carried);
408 settle(&mailbox, &id, carried, &local);
409 }
410 continue;
411 }
412 for id in pending {
413 let carried = match mailbox.find(&id) {
414 Ok(Some(stored))
417 if stored.state == crate::mailbox::MailState::Read
418 && mailbox.delivered_to_recipient(&stored.envelope.id) =>
419 {
420 Carried::Filed
421 }
422 Ok(Some(stored)) => carry(homes, &address, stored.envelope, &local).await,
423 _ => Carried::Filed,
424 };
425 record_local(&root, &address, &id, "door", &carried);
426 settle(&mailbox, &id, carried, &local);
427 }
428 }
429 for (machine, mut items) in remote {
430 items.sort_by_key(|(_, _, filed)| *filed);
431 let wait = crate::mailbox::machine_wait(&root, &machine);
432 if let Some(wait) = &wait {
433 if now.saturating_sub(wait.offline_since_ms) >= OFFLINE_CEILING.as_millis() as u64 {
435 let reason = format!(
436 "{machine} has been offline for {} days (since epoch ms {}): {}",
437 OFFLINE_CEILING.as_secs() / 86_400,
438 wait.offline_since_ms,
439 wait.last_error.as_deref().unwrap_or("not linked")
440 );
441 for (mailbox, id, _) in &items {
442 expire(mailbox, id, mailbox.wake_state(id).attempts, reason.clone());
443 }
444 crate::mailbox::set_machine_wait(&root, &machine, None).ok();
445 crate::mailbox::record_carrier_call(
446 &root,
447 &serde_json::json!({"t": now, "machine": machine, "outcome": "expired", "messages": items.len(), "reason": reason}),
448 );
449 continue;
450 }
451 if wait.next_probe_ms > now {
452 continue;
453 }
454 }
455 let mut carried_here = 0usize;
457 let mut offline_answer: Option<String> = None;
458 let mut called = false;
459 for (mailbox, id, _) in &items {
460 if wait.is_none() && mailbox.wake_state(id).next_at_ms > now {
461 continue;
462 }
463 let carried = match mailbox.find(id) {
464 Ok(Some(stored))
466 if stored.state == crate::mailbox::MailState::Read
467 && mailbox.delivered_to_recipient(&stored.envelope.id) =>
468 {
469 Carried::Filed
470 }
471 Ok(Some(stored)) => {
472 called = true;
473 carry(homes, mailbox.address(), stored.envelope, &local).await
474 }
475 _ => Carried::Filed,
476 };
477 if let Carried::Offline(reason) = carried {
478 offline_answer = Some(reason);
479 break;
480 }
481 carried_here += 1;
482 settle(mailbox, id, carried, &local);
483 }
484 if called {
485 crate::mailbox::record_carrier_call(
486 &root,
487 &serde_json::json!({
488 "t": now, "machine": machine, "outcome": if offline_answer.is_some() { "offline" } else { "carried" },
489 "waiting": items.len(), "carried": carried_here, "reason": offline_answer,
490 }),
491 );
492 }
493 match offline_answer {
494 Some(reason) => {
495 let next = crate::mailbox::MachineWait {
496 machine: machine.clone(),
497 offline_since_ms: wait.as_ref().map_or(now, |wait| wait.offline_since_ms),
498 next_probe_ms: now + LAST_RETRY.as_millis() as u64,
499 probes: wait.as_ref().map_or(0, |wait| wait.probes) + 1,
500 last_error: Some(reason),
501 };
502 crate::mailbox::set_machine_wait(&root, &machine, Some(&next)).ok();
503 }
504 None if called => {
505 crate::mailbox::set_machine_wait(&root, &machine, None).ok();
506 }
507 None => {}
508 }
509 }
510 }
511}
512
513async fn carry(
515 homes: &HarnessHomes,
516 to: &MailAddress,
517 mut envelope: Envelope,
518 local: &str,
519) -> Carried {
520 if to.machine != local {
521 let request = serde_json::json!({"op": "file", "to": to.to_string(), "envelope": envelope, "wake": true});
522 return match crate::mailbox::teams_mail(&to.machine, &request) {
523 Ok(answer) if answer["code"].as_i64() == Some(0) => Carried::Delivered,
524 Ok(answer) => {
525 let reason = format!(
526 "{} did not take it: {}",
527 to.machine,
528 answer["text"]
529 .as_str()
530 .or(answer["detail"].as_str())
531 .unwrap_or("no reason given")
532 );
533 if offline(&answer.to_string()) {
534 Carried::Offline(reason)
535 } else {
536 Carried::Later(reason)
537 }
538 }
539 Err(error) if offline(&error) => {
540 Carried::Offline(format!("{} is offline: {error}", to.machine))
541 }
542 Err(error) => {
543 Carried::Later(format!("Teams did not carry it to {}: {error}", to.machine))
544 }
545 };
546 }
547 if to.harness == crate::mail_agent::AGENT_HARNESS {
549 let plan =
550 crate::mail_agent::plan(&mut envelope, to, &crate::mail_agent::Channel::default());
551 let Ok(Some(plan)) = plan else {
552 return Carried::Filed;
553 };
554 let caller = crate::mail_route::Caller {
555 address: envelope.from.clone(),
556 name: envelope.from_name.clone(),
557 };
558 let outcome =
559 crate::mail_send::deliver_planned(homes, &caller, &envelope, &plan, false, false).await;
560 return if outcome.code == 0 {
561 Carried::Delivered
562 } else {
563 Carried::Later(outcome.text)
564 };
565 }
566 match door_for(homes, to) {
567 Ok(door @ (crate::mail_route::Door::Native(_) | crate::mail_route::Door::Runtime(_))) => {
568 let Ok(mailbox) = crate::mailbox::Mailbox::open(&mail_root(), to) else {
572 return Carried::Later("its mailbox could not be opened".into());
573 };
574 let mut state = mailbox.wake_state(&envelope.id);
575 if state.handover_at_ms.is_some() {
576 return Carried::Filed;
577 }
578 state.handover_at_ms = Some(crate::mailbox::now_ms());
579 if mailbox.set_wake_state(&envelope.id, &state).is_err() {
580 return Carried::Later("its handover could not be recorded".into());
581 }
582 let delivered = deliver(&envelope, to, &door, true, false).await;
583 state.handover_at_ms = None;
584 mailbox.set_wake_state(&envelope.id, &state).ok();
585 match delivered {
586 Ok(Ok(_)) => Carried::Delivered,
587 Ok(Err(_)) => Carried::Filed,
589 Err(error) => Carried::Later(error),
590 }
591 }
592 Ok(_) => Carried::Filed,
593 Err(crate::mail_route::NoDoor::NotRunning) => {
594 Carried::Later("its session is not running".into())
595 }
596 Err(crate::mail_route::NoDoor::OtherMachine(machine)) => {
597 Carried::Later(format!("it is on {machine}"))
598 }
599 }
600}
601
602fn settle(mailbox: &crate::mailbox::Mailbox, id: &str, carried: Carried, local: &str) {
605 match carried {
606 Carried::Delivered | Carried::Filed => {
607 mailbox.acknowledge_wake(id);
608 if matches!(carried, Carried::Delivered) {
609 mailbox
611 .record_claim(id, &crate::mailbox::Claim::new("carrier", None, true))
612 .ok();
613 if let Ok(Some(stored)) = mailbox.find(id) {
614 mailbox.mark_read(&stored).ok();
615 }
616 }
617 }
618 Carried::Offline(_) => {}
619 Carried::Later(reason) => {
620 let now = crate::mailbox::now_ms();
621 let mut state = mailbox.wake_state(id);
622 state.attempts += 1;
623 let since = *state.failing_since_ms.get_or_insert(now);
624 let lifetime = if mailbox.address().machine == local {
625 LOCAL_WAKE_LIFETIME
626 } else {
627 REMOTE_WAKE_LIFETIME
628 };
629 if state.attempts >= WAKE_ATTEMPTS
630 && now.saturating_sub(since) >= lifetime.as_millis() as u64
631 {
632 expire(mailbox, id, state.attempts, reason);
633 return;
634 }
635 let wait = FIRST_RETRY
636 .saturating_mul(1 << (state.attempts - 1).min(16))
637 .min(LAST_RETRY);
638 state.next_at_ms = now + wait.as_millis() as u64;
639 state.last_error = Some(reason);
640 mailbox.set_wake_state(id, &state).ok();
641 }
642 }
643}
644
645fn expire(mailbox: &crate::mailbox::Mailbox, id: &str, attempts: u32, reason: String) {
647 let envelope = mailbox
648 .find(id)
649 .ok()
650 .flatten()
651 .map(|stored| stored.envelope);
652 let expiry = crate::mailbox::WakeExpiry {
653 id: id.to_string(),
654 to: mailbox.address().to_string(),
655 from: envelope
656 .as_ref()
657 .map(|e| e.from.to_string())
658 .unwrap_or_default(),
659 subject: envelope.as_ref().and_then(|e| e.subject.clone()),
660 filed_at_ms: envelope.as_ref().map_or(0, |e| e.created_at_ms),
661 expired_at_ms: crate::mailbox::now_ms(),
662 attempts,
663 reason,
664 };
665 mailbox.expire_wake(id, &expiry).ok();
666}
667
668fn record_local(root: &std::path::Path, to: &MailAddress, id: &str, path: &str, carried: &Carried) {
671 let (outcome, reason) = match carried {
672 Carried::Delivered => ("delivered", None),
673 Carried::Filed => ("filed", None),
674 Carried::Later(reason) => ("later", Some(reason.as_str())),
675 Carried::Offline(reason) => ("offline", Some(reason.as_str())),
676 };
677 crate::mailbox::record_carrier_call(
678 root,
679 &serde_json::json!({"t": crate::mailbox::now_ms(), "to": to.to_string(), "id": id, "path": path, "outcome": outcome, "reason": reason}),
680 );
681}