Skip to main content

supercode_harness/
mail_watch.rs

1//! Idle notices for every harness, from session activity.
2//!
3//! Knowing that a session finished does not need control of it: supercode's
4//! normalized activity ([`crate::session_activity`]) reports busy/idle for a
5//! native Claude session, a running Codex, and any runtime supercode hosts.
6//! A sender that asked for `--notify-when-idle` gets one notice when the
7//! receiver is seen working after the message and then idle, one when the
8//! receiver ends, or one saying the subscription expired after
9//! [`IDLE_SUBSCRIPTION_LIFETIME`]. The notice is delivered through the one
10//! router ([`crate::mail_route`]), so each subscriber gets it by its own door.
11
12use 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
21/// How long a subscription waits for the receiver to work and go idle.
22pub const IDLE_SUBSCRIPTION_LIFETIME: Duration = Duration::from_secs(24 * 60 * 60);
23
24/// Polls the activity of every session someone subscribed to.
25#[derive(Debug, Default)]
26pub struct IdleWatcher {
27    monitor: SessionActivityMonitor,
28}
29
30impl IdleWatcher {
31    /// A watcher with no state yet.
32    pub fn new() -> Self {
33        Self::default()
34    }
35
36    /// Whether any subscription is waiting.
37    pub fn has_subscriptions() -> bool {
38        !subscribed_mailboxes(&mail_root()).is_empty()
39    }
40
41    /// One pass: settle every subscription whose receiver finished a turn,
42    /// ended, or outlived the subscription.
43    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            // A runtime supercode hosts reports its own turn; for the rest,
53            // session activity reads the harness's lifecycle.
54            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            // One turn settles every message a sender subscribed with; the sender hears it once,
81            // about its latest message, and its earlier subscriptions settle with it silently.
82            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                // A reply is the runtime's answer to this message: it settles
89                // once the runtime is idle and has answered, however fast
90                // the turn was.
91                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
154/// The locator activity resolution needs: a Codex session by its rollout
155/// file; Claude sessions and hosted runtimes by id.
156fn 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        // A subscriber on another machine, or one no longer running, still
215        // finds the notice in its mailbox.
216        Err(_) => {
217            crate::mailbox::deliver_to(subscriber, &envelope).ok();
218        }
219    }
220}
221
222/// The receiver's answer to a message that went in `reply-via=final-message`,
223/// sent back to its sender as the reply. With no answer (the receiver ended,
224/// or the subscription expired first), a notice says so, so the sender is
225/// never left waiting in silence.
226async 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
263/// The name a notice calls its session by.
264fn 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
278/// Type the user's waiting turns into their sessions' panes, each once its
279/// composer is empty. A turn waits when the person has a draft; the mailbox
280/// is where it waits.
281pub 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
299/// How long a wake waits after a failed attempt, at first and at most (doubling between).
300const FIRST_RETRY: Duration = Duration::from_secs(2);
301const LAST_RETRY: Duration = Duration::from_secs(300);
302/// Counted failures a wake has at least before it may expire.
303pub const WAKE_ATTEMPTS: u32 = 8;
304/// How long a wake to a session on this machine keeps trying, from its first counted failure. Its basis is a session
305/// restart: the board relaunches or resumes a lost card session within its launch window (300 s, `launchWindowSeconds`),
306/// the machine daemon restarts its watch loop within 5 s, and a person resuming a closed session does so within
307/// minutes; thirty minutes outlasts several of those. Past it the session is not coming back soon, and the message
308/// waits unread for whenever it does (a session reads its unread mail first).
309pub const LOCAL_WAKE_LIFETIME: Duration = Duration::from_secs(30 * 60);
310/// How long a wake to another machine keeps trying, from its first counted failure (a refusal by that machine or by
311/// Teams, or Teams unreachable). Three days outlasts a weekend with the server down. A refusal saying the machine is
312/// offline (a reboot, a closed lid) is never counted: that mail waits for the machine, however long, and is carried
313/// within [`LAST_RETRY`] of its link coming back.
314pub const REMOTE_WAKE_LIFETIME: Duration = Duration::from_secs(3 * 24 * 60 * 60);
315
316/// The mailbox's own delivery (D140): every message filed with a wake (by a filer that does not wait, such as a board's
317/// dispatcher, or by a send to a hooked session) is handed to its session by the session's own door, by the daemon that
318/// runs this watch loop. Mail for another machine is carried there through Teams (`supercode teams mail`), the only
319/// cross-machine leg, and filed with the same wake, so that machine's watch loop delivers it. A failed attempt waits
320/// [`FIRST_RETRY`] doubling to [`LAST_RETRY`], kept in the wake itself, so a restarted watcher waits as the last one
321/// did. A wake expires only once it has failed [`WAKE_ATTEMPTS`] counted times and its first counted failure is older
322/// than its bound ([`LOCAL_WAKE_LIFETIME`] here, [`REMOTE_WAKE_LIFETIME`] for another machine): the message stays filed
323/// and unread, and the expiry is recorded in its mailbox (`supercode message expired`), which the board's diagnostics
324/// read. Mail for a machine that is offline is not failing: it is offered again every [`LAST_RETRY`] until the machine
325/// links, and never expires for that.
326#[derive(Debug, Default)]
327pub struct MailCarrier;
328
329/// What one attempt at a wake came to.
330enum Carried {
331    /// Delivered: the wake is done and the message is read.
332    Delivered,
333    /// Filed where its reader looks for it (an operator's mailbox, a hook that reads it at the next tool call, a
334    /// session with no door): the wake is done and the message stays unread for its reader.
335    Filed,
336    /// Not now, and why (its session is not running, the attempt failed): a counted failure, tried again after a wait.
337    Later(String),
338    /// Its machine is offline (not linked to Teams): not a failure; offered again at [`LAST_RETRY`].
339    Offline(String),
340}
341
342/// Whether a cross-machine answer says the receiving machine is offline (not linked), rather than refusing the mail.
343fn offline(text: &str) -> bool {
344    text.contains("machine_offline")
345}
346
347impl MailCarrier {
348    /// A carrier.
349    pub fn new() -> Self {
350        Self
351    }
352
353    /// One pass over every mailbox holding a wake that is due.
354    pub async fn tick(&mut self, homes: &HarnessHomes) {
355        let local = crate::mailbox::local_machine_name();
356        let now = crate::mailbox::now_ms();
357        let mut offline_now: std::collections::HashMap<String, String> =
358            std::collections::HashMap::new();
359        for mailbox in crate::mailbox::mailboxes_with_wake_requests(&mail_root()) {
360            let Ok(pending) = mailbox.pending_wakes() else {
361                continue;
362            };
363            let address = mailbox.address().clone();
364            let pending: Vec<String> = pending
365                .into_iter()
366                .filter(|id| mailbox.wake_state(id).next_at_ms <= now)
367                .collect();
368            if pending.is_empty() {
369                continue;
370            }
371            // a hooked session in a daemon pane is woken once for all its waiting mail, its hook showing it all
372            if address.machine == local {
373                if let Ok(crate::mail_route::Door::Hook {
374                    pane: Some(pane), ..
375                }) = door_for(homes, &address)
376                {
377                    let woken =
378                        crate::mail_route::wake_hook_mailbox(&mailbox, &pane, &pending).await;
379                    for id in pending {
380                        let carried = match &woken {
381                            Ok(true) => Carried::Delivered,
382                            Ok(false) => Carried::Later("its pane was not at its prompt".into()),
383                            Err(error) => Carried::Later(error.clone()),
384                        };
385                        settle(&mailbox, &id, carried, &local);
386                    }
387                    continue;
388                }
389            }
390            for id in pending {
391                // a machine found offline in this pass is not asked again for its other mail
392                if let Some(reason) = offline_now.get(&address.machine) {
393                    settle(&mailbox, &id, Carried::Offline(reason.clone()), &local);
394                    continue;
395                }
396                let carried = match mailbox.find(&id) {
397                    // its recipient has it already (its own read, or a hand-over by an attempt that ended before it was settled):
398                    // never delivered twice; a read by anyone else does not count
399                    Ok(Some(stored))
400                        if stored.state == crate::mailbox::MailState::Read
401                            && mailbox.delivered_to_recipient(&stored.envelope.id) =>
402                    {
403                        Carried::Filed
404                    }
405                    Ok(Some(stored)) => carry(homes, &address, stored.envelope, &local).await,
406                    _ => Carried::Filed,
407                };
408                if let Carried::Offline(reason) = &carried {
409                    offline_now.insert(address.machine.clone(), reason.clone());
410                }
411                settle(&mailbox, &id, carried, &local);
412            }
413        }
414    }
415}
416
417/// One message toward its reader.
418async fn carry(
419    homes: &HarnessHomes,
420    to: &MailAddress,
421    mut envelope: Envelope,
422    local: &str,
423) -> Carried {
424    if to.machine != local {
425        let request = serde_json::json!({"op": "file", "to": to.to_string(), "envelope": envelope, "wake": true});
426        return match crate::mailbox::teams_mail(&to.machine, &request) {
427            Ok(answer) if answer["code"].as_i64() == Some(0) => Carried::Delivered,
428            Ok(answer) => {
429                let reason = format!(
430                    "{} did not take it: {}",
431                    to.machine,
432                    answer["text"]
433                        .as_str()
434                        .or(answer["detail"].as_str())
435                        .unwrap_or("no reason given")
436                );
437                if offline(&answer.to_string()) {
438                    Carried::Offline(reason)
439                } else {
440                    Carried::Later(reason)
441                }
442            }
443            Err(error) if offline(&error) => {
444                Carried::Offline(format!("{} is offline: {error}", to.machine))
445            }
446            Err(error) => {
447                Carried::Later(format!("Teams did not carry it to {}: {error}", to.machine))
448            }
449        };
450    }
451    // an agent's one mailbox: its plan decides which of its sessions receives it (ADR 0008 decisions 5-9)
452    if to.harness == crate::mail_agent::AGENT_HARNESS {
453        let plan =
454            crate::mail_agent::plan(&mut envelope, to, &crate::mail_agent::Channel::default());
455        let Ok(Some(plan)) = plan else {
456            return Carried::Filed;
457        };
458        let caller = crate::mail_route::Caller {
459            address: envelope.from.clone(),
460            name: envelope.from_name.clone(),
461        };
462        let outcome =
463            crate::mail_send::deliver_planned(homes, &caller, &envelope, &plan, false, false).await;
464        return if outcome.code == 0 {
465            Carried::Delivered
466        } else {
467            Carried::Later(outcome.text)
468        };
469    }
470    match door_for(homes, to) {
471        Ok(door @ (crate::mail_route::Door::Native(_) | crate::mail_route::Door::Runtime(_))) => {
472            // Idempotent by the message's id across a crash: the handover is recorded before it starts, and a wake whose
473            // last handover never settled (its watcher died in it) is not handed over again; it stays unread, which its
474            // session reads first.
475            let Ok(mailbox) = crate::mailbox::Mailbox::open(&mail_root(), to) else {
476                return Carried::Later("its mailbox could not be opened".into());
477            };
478            let mut state = mailbox.wake_state(&envelope.id);
479            if state.handover_at_ms.is_some() {
480                return Carried::Filed;
481            }
482            state.handover_at_ms = Some(crate::mailbox::now_ms());
483            if mailbox.set_wake_state(&envelope.id, &state).is_err() {
484                return Carried::Later("its handover could not be recorded".into());
485            }
486            let delivered = deliver(&envelope, to, &door, true, false).await;
487            state.handover_at_ms = None;
488            mailbox.set_wake_state(&envelope.id, &state).ok();
489            match delivered {
490                Ok(Ok(_)) => Carried::Delivered,
491                // a message its door can never take (too long for a relay) stays filed for the session to read
492                Ok(Err(_)) => Carried::Filed,
493                Err(error) => Carried::Later(error),
494            }
495        }
496        Ok(_) => Carried::Filed,
497        Err(crate::mail_route::NoDoor::NotRunning) => {
498            Carried::Later("its session is not running".into())
499        }
500        Err(crate::mail_route::NoDoor::OtherMachine(machine)) => {
501            Carried::Later(format!("it is on {machine}"))
502        }
503    }
504}
505
506/// Record how a wake went: done, waiting for its next attempt (an offline machine's mail uncounted), or expired at its
507/// bound.
508fn settle(mailbox: &crate::mailbox::Mailbox, id: &str, carried: Carried, local: &str) {
509    match carried {
510        Carried::Delivered | Carried::Filed => {
511            mailbox.acknowledge_wake(id);
512            if matches!(carried, Carried::Delivered) {
513                if let Ok(Some(stored)) = mailbox.find(id) {
514                    mailbox.mark_read(&stored).ok();
515                }
516            }
517        }
518        Carried::Offline(reason) => {
519            let mut state = mailbox.wake_state(id);
520            state.next_at_ms = crate::mailbox::now_ms() + LAST_RETRY.as_millis() as u64;
521            state.last_error = Some(reason);
522            mailbox.set_wake_state(id, &state).ok();
523        }
524        Carried::Later(reason) => {
525            let now = crate::mailbox::now_ms();
526            let mut state = mailbox.wake_state(id);
527            state.attempts += 1;
528            let since = *state.failing_since_ms.get_or_insert(now);
529            let lifetime = if mailbox.address().machine == local {
530                LOCAL_WAKE_LIFETIME
531            } else {
532                REMOTE_WAKE_LIFETIME
533            };
534            if state.attempts >= WAKE_ATTEMPTS
535                && now.saturating_sub(since) >= lifetime.as_millis() as u64
536            {
537                let envelope = mailbox
538                    .find(id)
539                    .ok()
540                    .flatten()
541                    .map(|stored| stored.envelope);
542                let expiry = crate::mailbox::WakeExpiry {
543                    id: id.to_string(),
544                    to: mailbox.address().to_string(),
545                    from: envelope
546                        .as_ref()
547                        .map(|e| e.from.to_string())
548                        .unwrap_or_default(),
549                    subject: envelope.as_ref().and_then(|e| e.subject.clone()),
550                    filed_at_ms: envelope.as_ref().map_or(0, |e| e.created_at_ms),
551                    expired_at_ms: now,
552                    attempts: state.attempts,
553                    reason,
554                };
555                mailbox.expire_wake(id, &expiry).ok();
556                return;
557            }
558            let wait = FIRST_RETRY
559                .saturating_mul(1 << (state.attempts - 1).min(16))
560                .min(LAST_RETRY);
561            state.next_at_ms = now + wait.as_millis() as u64;
562            state.last_error = Some(reason);
563            mailbox.set_wake_state(id, &state).ok();
564        }
565    }
566}