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}