返回 CodeWhale
notification_delivery.rs
根目录 / crates / tui / src / runtime_api / notification_delivery.rs
1 //! Native-client projection of the existing event, payload and notification
2 //! policy owners. This endpoint prepares an attempt; it never submits a banner.
3
4 use super::*;
5 use crate::notify::payload::NotificationPayload;
6 use crate::notify::{
7 DeliveryOutcome, Method, NotificationGate, audio as notification_audio, sound_policy,
8 };
9 use crate::runtime_threads::{RuntimeEventRecord, RuntimeTurnStatus};
10 use codewhale_localization::{Locale, MessageId, tr};
11
12 #[derive(Deserialize)]
13 #[serde(deny_unknown_fields)]
14 pub(super) struct PrepareRequest {
15 seq: u64,
16 focused: bool,
17 unfocused_for_ms: u64,
18 locale: String,
19 }
20
21 #[derive(Debug, Serialize)]
22 pub(super) struct PreparedNotification {
23 status: &'static str,
24 #[serde(skip_serializing_if = "Option::is_none")]
25 headline: Option<String>,
26 #[serde(skip_serializing_if = "Option::is_none")]
27 body: Option<String>,
28 /// A selection, never an audio receipt. No file paths cross this boundary.
29 sound: &'static str,
30 }
31
32 impl PreparedNotification {
33 fn suppressed(status: &'static str) -> Self {
34 Self {
35 status,
36 headline: None,
37 body: None,
38 sound: "off",
39 }
40 }
41 }
42
43 pub(super) async fn prepare(
44 State(state): State<RuntimeApiState>,
45 Path(id): Path<String>,
46 Json(request): Json<PrepareRequest>,
47 ) -> Result<Json<PreparedNotification>, ApiError> {
48 let Some(previous) = request.seq.checked_sub(1) else {
49 return Err(ApiError::bad_request(
50 "notification sequence must be positive",
51 ));
52 };
53 // Read the selected thread's durable event, never a renderer-supplied copy
54 // or title. Dropping the replay receiver stops the bounded reader.
55 let mut replay = state
56 .runtime_threads
57 .replay_events(&id, Some(previous), None)
58 .await
59 .map_err(|error| ApiError::internal(error.to_string()))?;
60 let mut selected = None;
61 while let Some(batch) = replay.batches.recv().await {
62 let batch = batch.map_err(ApiError::internal)?;
63 if let Some(event) = batch.into_iter().next() {
64 if event.seq == request.seq {
65 selected = Some(event);
66 }
67 break;
68 }
69 }
70 let event = selected.ok_or_else(|| {
71 ApiError::bad_request("notification event does not belong to this thread")
72 })?;
73 let locale = Locale::shipped()
74 .iter()
75 .copied()
76 .find(|locale| locale.tag().eq_ignore_ascii_case(&request.locale))
77 .or_else(|| matches!(request.locale.as_str(), "zh" | "zh-CN").then_some(Locale::ZhHans))
78 .ok_or_else(|| ApiError::bad_request("unsupported notification locale"))?;
79 // Replay can wait on disk while the same turn's request settles. Validate
80 // the selected record against a fresh snapshot, with no later await.
81 let detail = state
82 .runtime_threads
83 .get_thread_detail(&id)
84 .await
85 .map_err(map_thread_err)?;
86 let config = state.config.read().clone();
87 Ok(Json(prepare_record(
88 &config,
89 &detail,
90 &event,
91 request.focused,
92 Duration::from_millis(request.unfocused_for_ms),
93 locale,
94 Utc::now(),
95 )))
96 }
97
98 #[allow(clippy::too_many_arguments)] // Canonical snapshot plus observed host facts; no second policy object.
99 fn prepare_record(
100 config: &Config,
101 detail: &ThreadDetail,
102 event: &RuntimeEventRecord,
103 focused: bool,
104 unfocused_for: Duration,
105 locale: Locale,
106 now: chrono::DateTime<Utc>,
107 ) -> PreparedNotification {
108 // Old/recovered records remain visible in the work history but cannot
109 // become a fresh OS interruption. The host also anchors its replay cursor.
110 let age = now.signed_duration_since(event.timestamp).num_seconds();
111 if !(0..=60).contains(&age) || event.payload.get("recovered") == Some(&json!(true)) {
112 return PreparedNotification::suppressed("expired");
113 }
114 let Some(turn) = detail
115 .turns
116 .iter()
117 .find(|turn| Some(&turn.id) == event.turn_id.as_ref())
118 else {
119 return PreparedNotification::suppressed("settled");
120 };
121 if detail.thread.latest_turn_id.as_deref() != Some(turn.id.as_str()) {
122 return PreparedNotification::suppressed("settled");
123 }
124 if matches!(
125 event.event.as_str(),
126 "approval.required" | "user_input.required"
127 ) && !matches!(
128 turn.status,
129 RuntimeTurnStatus::Queued | RuntimeTurnStatus::InProgress
130 ) {
131 return PreparedNotification::suppressed("settled");
132 }
133 let payload = match event.event.as_str() {
134 "turn.completed" if turn.status == RuntimeTurnStatus::Completed => {
135 NotificationPayload::turn_complete(&tr(locale, MessageId::NotificationTurnComplete))
136 }
137 "approval.required" => {
138 let Some(pending) = detail.pending_approvals.iter().find(|pending| {
139 event.payload.get("id").and_then(Value::as_str) == Some(pending.id.as_str())
140 && pending.turn_id == turn.id
141 }) else {
142 return PreparedNotification::suppressed("settled");
143 };
144 NotificationPayload::approval_needed(
145 &tr(locale, MessageId::ConfigLabelNotificationApprovalNeeded),
146 &pending.tool_name,
147 )
148 }
149 "user_input.required"
150 if detail.pending_user_inputs.iter().any(|pending| {
151 event.payload.get("id").and_then(Value::as_str) == Some(pending.id.as_str())
152 && pending.turn_id == turn.id
153 }) =>
154 {
155 NotificationPayload::input_needed(&tr(
156 locale,
157 MessageId::ConfigLabelNotificationInputNeeded,
158 ))
159 }
160 _ => return PreparedNotification::suppressed("unsupported_event"),
161 };
162 prepare_payload(
163 config,
164 &payload,
165 Duration::from_millis(turn.duration_ms.unwrap_or(0)),
166 focused,
167 unfocused_for,
168 )
169 }
170
171 fn prepare_payload(
172 config: &Config,
173 payload: &NotificationPayload,
174 elapsed: Duration,
175 focused: bool,
176 unfocused_for: Duration,
177 ) -> PreparedNotification {
178 let notification_config = config.notifications_config();
179 let Some((method, threshold, _)) = crate::notify::settings_projection(&notification_config)
180 else {
181 return PreparedNotification::suppressed("suppressed");
182 };
183 let method = match method {
184 Method::Auto => Method::MacOS,
185 Method::Off => Method::Off,
186 _ => return PreparedNotification::suppressed("unsupported_method"),
187 };
188 let attention =
189 crate::notify::native_attention_allowed(&notification_config, focused, unfocused_for);
190 let threshold = if payload.kind() == crate::notify::payload::NotificationKind::TurnComplete {
191 threshold
192 } else {
193 Duration::ZERO
194 };
195 let mut sound = "off";
196 // Capture the shared policy's intended sinks. Neither closure performs IO;
197 // the response says prepared, never dispatched/delivered. The native host
198 // owns the subsequent permission check and submission receipt.
199 let outcome = crate::tui::notifications::notify_with_sinks(
200 method,
201 false,
202 payload,
203 threshold,
204 elapsed,
205 NotificationGate::from_config(&notification_config),
206 attention,
207 &mut std::io::sink(),
208 &mut |kind, bell| {
209 sound_policy::decide_configured(
210 &notification_config,
211 kind,
212 sound_policy::epoch_millis_now(),
213 bell,
214 )
215 },
216 &mut |cue, _| {
217 sound = match cue {
218 sound_policy::SoundCue::Beep => "beep",
219 sound_policy::SoundCue::Whale => "whale",
220 sound_policy::SoundCue::File(_) => "file",
221 sound_policy::SoundCue::Bell => "bell",
222 sound_policy::SoundCue::DoubleBell => "double-bell",
223 };
224 notification_audio::AudioOutcome::Dispatched
225 },
226 &mut |_| DeliveryOutcome::Dispatched(Method::MacOS),
227 );
228 if !matches!(outcome, DeliveryOutcome::Dispatched(_)) {
229 return PreparedNotification::suppressed("suppressed");
230 }
231 PreparedNotification {
232 status: "prepared",
233 headline: Some(payload.headline().to_string()),
234 body: Some(payload.body()),
235 sound,
236 }
237 }
238
238 lines RUST