| 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(¬ification_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(¬ification_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(¬ification_config), |
| 206 | attention, |
| 207 | &mut std::io::sink(), |
| 208 | &mut |kind, bell| { |
| 209 | sound_policy::decide_configured( |
| 210 | ¬ification_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 |