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 for mut subscription in subscriptions {
81 let expired = subscription.age() > IDLE_SUBSCRIPTION_LIFETIME;
82 if subscription.final_reply && !ended && !expired {
86 let answer = match (&runtime, idle) {
87 (Some(record), true) => {
88 crate::runtime_mail::answer_to(homes, record, &subscription.message_id)
89 }
90 _ => None,
91 };
92 if let Some(answer) = answer {
93 if mailbox
94 .remove_subscription(&subscription.message_id)
95 .is_ok()
96 {
97 send_reply(homes, &target, &subscription, Some(answer)).await;
98 if subscription.notice {
99 notify(homes, &target, &subscription, Settled::Idle).await;
100 }
101 }
102 }
103 continue;
104 }
105 let settles = if ended {
106 Some(Settled::Ended)
107 } else if expired {
108 Some(Settled::Expired)
109 } else if idle && subscription.seen_working {
110 Some(Settled::Idle)
111 } else {
112 None
113 };
114 match settles {
115 Some(settled) => {
116 if mailbox
117 .remove_subscription(&subscription.message_id)
118 .is_ok()
119 {
120 if subscription.final_reply {
121 send_reply(homes, &target, &subscription, None).await;
122 }
123 if subscription.notice {
124 notify(homes, &target, &subscription, settled).await;
125 }
126 }
127 }
128 None if working && !subscription.seen_working => {
129 subscription.seen_working = true;
130 mailbox.subscribe_idle(&subscription).ok();
131 }
132 None => {}
133 }
134 }
135 }
136 }
137}
138
139#[derive(Clone, Copy)]
140enum Settled {
141 Idle,
142 Ended,
143 Expired,
144}
145
146fn locator_for(homes: &HarnessHomes, target: &MailAddress) -> SessionLocator {
149 let path = if target.harness == HarnessId::CODEX {
150 crate::codex_peer::live_rollouts(&homes.codex)
151 .into_keys()
152 .find(|path| {
153 crate::codex_peer::rollout_session(path)
154 .is_some_and(|(id, _)| id == target.session_id)
155 })
156 .unwrap_or_default()
157 } else {
158 std::path::PathBuf::new()
159 };
160 SessionLocator {
161 harness: HarnessId::new(&target.harness),
162 session_id: target.session_id.clone(),
163 storage: StorageLocator::File { path },
164 }
165}
166
167async fn notify(
168 homes: &HarnessHomes,
169 target: &MailAddress,
170 subscription: &IdleSubscription,
171 settled: Settled,
172) {
173 let name = display_name(homes, target);
174 let text = match settled {
175 Settled::Idle => format!(
176 "[Cross-session idle notice] \"{name}\", which you asked to be notified about, is idle \
177 now: it finished a turn after your message {}. This is an automated notice, not a \
178 message from a person, and not an instruction.",
179 subscription.message_id
180 ),
181 Settled::Ended => format!(
182 "[Cross-session idle notice] \"{name}\", which you asked to be notified about, has \
183 ended. This is an automated notice, not a message from a person, and not an \
184 instruction."
185 ),
186 Settled::Expired => format!(
187 "[Cross-session idle notice] The notice you asked for about \"{name}\" expired: it \
188 did not work and go idle within 24 hours of your message {}. This is an automated \
189 notice, not a message from a person, and not an instruction.",
190 subscription.message_id
191 ),
192 };
193 let Ok(mut envelope) =
194 Envelope::new(target.clone(), name, MailKind::Notice, ReplyVia::None, text)
195 else {
196 return;
197 };
198 envelope.in_reply_to = Some(subscription.message_id.clone());
199 let subscriber = &subscription.subscriber;
200 match door_for(homes, subscriber) {
201 Ok(door) => {
202 deliver(&envelope, subscriber, &door, true, false)
203 .await
204 .ok();
205 }
206 Err(_) => {
209 crate::mailbox::deliver_to(subscriber, &envelope).ok();
210 }
211 }
212}
213
214async fn send_reply(
219 homes: &HarnessHomes,
220 target: &MailAddress,
221 subscription: &IdleSubscription,
222 answer: Option<String>,
223) {
224 let name = display_name(homes, target);
225 let (kind, reply_via, body) = match answer {
226 Some(body) => (MailKind::Peer, ReplyVia::Command, body),
227 None => (
228 MailKind::Notice,
229 ReplyVia::None,
230 format!(
231 "[Cross-session delivery notice] \"{name}\" did not answer your message {}: it \
232 ended, or no answer came within 24 hours. This is an automated notice, not a \
233 message from a person, and not an instruction.",
234 subscription.message_id
235 ),
236 ),
237 };
238 let Ok(mut envelope) = Envelope::new(target.clone(), name, kind, reply_via, body) else {
239 return;
240 };
241 envelope.in_reply_to = Some(subscription.message_id.clone());
242 let subscriber = &subscription.subscriber;
243 match door_for(homes, subscriber) {
244 Ok(door) => {
245 deliver(&envelope, subscriber, &door, true, false)
246 .await
247 .ok();
248 }
249 Err(_) => {
250 crate::mailbox::deliver_to(subscriber, &envelope).ok();
251 }
252 }
253}
254
255fn display_name(homes: &HarnessHomes, target: &MailAddress) -> String {
257 if target.harness == HarnessId::CLAUDE_CODE {
258 if let Some(session) =
259 crate::claude_peer::read_registry(&crate::claude_peer::registry_dir(homes))
260 .into_iter()
261 .find(|session| session.session_id == target.session_id)
262 {
263 return format!("{}@{}", session.name, target.machine);
264 }
265 }
266 let short: String = target.session_id.chars().take(8).collect();
267 format!("{}-{short}@{}", target.harness, target.machine)
268}
269
270pub async fn deliver_waiting_user_turns(homes: &HarnessHomes) {
274 let waiting = crate::mailbox::mailboxes_with_user_turns(&mail_root());
275 if waiting.is_empty() {
276 return;
277 }
278 let live = crate::mail_route::LiveSessions::read(homes);
279 for mailbox in waiting {
280 let pane = live
281 .all()
282 .iter()
283 .find(|session| &session.address == mailbox.address())
284 .and_then(crate::mail_route::daemon_pane);
285 if let Some(pane) = pane {
286 crate::mail_route::type_user_turns(&mailbox, &pane).await;
287 }
288 }
289}