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}