diff --git a/crates/core/src/call.rs b/crates/core/src/call.rs index 783278fd..f6107409 100644 --- a/crates/core/src/call.rs +++ b/crates/core/src/call.rs @@ -99,7 +99,7 @@ impl IncomingCall { caller_jid, is_video, is_offline: offer.offline, - received_at: wacore::time::now_utc(), + received_at: offer.timestamp, } } diff --git a/crates/core/src/calls.rs b/crates/core/src/calls.rs index 6cef44ca..55261565 100644 --- a/crates/core/src/calls.rs +++ b/crates/core/src/calls.rs @@ -325,6 +325,22 @@ impl CallState { /// The first offer wins the slot rather than the last, because that is the /// one the user has already been shown. pub fn set_incoming(&mut self, call: IncomingCall) -> Admission { + // Redelivered signaling must not park a call behind itself, including + // after it was accepted and its stage became active. + if self + .stage + .as_ref() + .is_some_and(|stage| stage.call_id() == call.call_id) + { + return Admission::Ringing; + } + if self + .waiting + .as_ref() + .is_some_and(|waiting| waiting.call_id() == &call.call_id) + { + return Admission::Parked; + } let Some(stage) = &self.stage else { self.remote_before_connect = None; self.stage = Some(Stage::Incoming(call)); @@ -1258,6 +1274,38 @@ mod tests { ); } + #[test] + fn duplicate_offers_never_park_a_call_behind_itself() { + let mut state = CallState::default(); + state.set_incoming(incoming("FIRST")); + let before = state.clone(); + assert_eq!(state.set_incoming(incoming("FIRST")), Admission::Ringing); + assert_eq!(state, before); + state.connect(&"FIRST".into()); + state.set_incoming(incoming("SECOND")); + let before = state.clone(); + assert_eq!(state.set_incoming(incoming("FIRST")), Admission::Ringing); + assert_eq!(state.set_incoming(incoming("SECOND")), Admission::Parked); + assert_eq!(state, before); + state.end(&"FIRST".into()); + assert_eq!(state.incoming().unwrap().call_id, "SECOND"); + assert!(state.waiting().is_none()); + } + + #[test] + fn incoming_calls_keep_the_protocol_timestamp() { + let mut wire = offer("OLD"); + wire.timestamp = chrono::DateTime::from_timestamp(1_700_000_000, 0).unwrap(); + let call = IncomingCall::new( + "OLD", + "Example".into(), + "a@s.whatsapp.net".into(), + false, + &wire, + ); + assert_eq!(call.received_at, wire.timestamp); + } + /// The same promotion, on the other way a stage can empty for good: a /// call this device never managed to place. `take_outgoing` cleared the /// stage and left the parked caller ringing behind an empty card. diff --git a/crates/daemon/src/server/requests.rs b/crates/daemon/src/server/requests.rs index 3d3a6d08..9f674d9d 100644 --- a/crates/daemon/src/server/requests.rs +++ b/crates/daemon/src/server/requests.rs @@ -173,6 +173,9 @@ pub(super) async fn handle_request( .await } ClientRequest::ReloadHistory => acted(dispatch(hub, commands, Action::ReloadHistory).await), + ClientRequest::ReconnectSession => { + acted(dispatch(hub, commands, Action::ReconnectSession).await) + } // Answered with the page under this id, like a download and for the // same reason. ClientRequest::LoadMessages(request) => { diff --git a/crates/daemon/src/server/tests.rs b/crates/daemon/src/server/tests.rs index f4e0f544..b58c1ba6 100644 --- a/crates/daemon/src/server/tests.rs +++ b/crates/daemon/src/server/tests.rs @@ -613,6 +613,48 @@ fn a_mismatched_hello_is_rejected_with_both_versions() { } } +#[test] +fn a_daemon_before_manual_whatsapp_retry_is_rejected() { + let rejection = check_hello(&hello(39, false)).expect_err("v39 cannot carry ReconnectSession"); + assert!(matches!( + serde_json::from_str::(&rejection.unwrap()).unwrap(), + DaemonMessage::Error { + error: ProtocolError::VersionMismatch { + client: 39, + daemon: 40 + }, + .. + } + )); +} + +#[tokio::test] +async fn manual_whatsapp_retry_reaches_the_session_while_offline() { + let hub = connected_hub(); + hub.apply(crate::state::Change::live( + oxidezap_ipc::DaemonEvent::ConnectionChanged(oxidezap_ipc::ConnectionState::Disconnected { + reason: "synthetic drop".into(), + }), + )); + let (commands, taken) = bridge(CommandOutcome::Accepted); + let answer = handle_request( + bare(ClientRequest::ReconnectSession), + &hub, + &no_plugins(), + &commands, + &outbox(), + ) + .await; + assert!(matches!( + parse(answer.frame), + DaemonMessage::Accepted { .. } + )); + assert!(matches!( + taken.await.unwrap(), + Some(Action::ReconnectSession) + )); +} + #[test] fn state_is_not_served_before_a_hello() { let line = serde_json::to_string(&ClientRequest::Snapshot).unwrap(); diff --git a/crates/daemon/src/session_bridge/act.rs b/crates/daemon/src/session_bridge/act.rs index 33e367ed..e6f36e92 100644 --- a/crates/daemon/src/session_bridge/act.rs +++ b/crates/daemon/src/session_bridge/act.rs @@ -1760,6 +1760,23 @@ impl Bridge { client.reload_history(); CommandOutcome::Accepted } + Action::ReconnectSession => { + if matches!( + self.hub.connection(), + oxidezap_ipc::ConnectionState::LoggedOut { .. } + ) { + return CommandOutcome::Refused( + "this account must be paired again".to_string(), + ); + } + match client.retry_connection().await { + Ok(Ok(())) => CommandOutcome::Accepted, + Ok(Err(detail)) => CommandOutcome::NoSession(detail), + Err(_) => CommandOutcome::NoSession( + "the session stopped during reconnection".to_string(), + ), + } + } Action::RefreshAvatars => { client.reset_avatar_cache(); CommandOutcome::Accepted diff --git a/crates/daemon/src/session_bridge/action.rs b/crates/daemon/src/session_bridge/action.rs index a4044dee..a0bf004a 100644 --- a/crates/daemon/src/session_bridge/action.rs +++ b/crates/daemon/src/session_bridge/action.rs @@ -55,6 +55,7 @@ pub enum Action { /// Reload the whole history, for a front end that has just attached and /// holds nothing. ReloadHistory, + ReconnectSession, /// Forget every resolved picture and resolve them again. /// /// For a cleared media cache: the metadata may be unchanged, but the bytes @@ -224,6 +225,7 @@ impl Action { _ => !matches!( self, Self::ReloadHistory + | Self::ReconnectSession | Self::RefreshVideo | Self::ForgetSession(_) | Self::MarkStatusWatched(_) diff --git a/crates/gui/src/app/body.rs b/crates/gui/src/app/body.rs index 1007a0eb..cef6f2eb 100644 --- a/crates/gui/src/app/body.rs +++ b/crates/gui/src/app/body.rs @@ -433,6 +433,34 @@ pub(super) mod tests { (cx, app) } + fn start_outgoing_call( + cx: &mut gpui::VisualTestContext, + app: &Entity, + call_id: &str, + ) { + app.update(cx, |app, cx| { + app.calls.update(cx, |calls, cx| { + calls.place_outgoing( + OutgoingCall::new( + call_id, + "peer@example.invalid".into(), + "Test peer".into(), + false, + ), + cx, + ); + }); + }); + cx.run_until_parked(); + cx.update(|window, cx| window.draw(cx).clear(cx)); + } + + fn outgoing_call_cancel_point(cx: &mut gpui::VisualTestContext) -> gpui::Point { + cx.debug_bounds("call-cancel") + .expect("the outgoing call must render its Cancel control") + .center() + } + #[gpui::test] fn pasted_image_waits_in_a_visible_preview(cx: &mut gpui::TestAppContext) { let (mut cx, app) = paste_preview_fixture(cx); @@ -511,6 +539,182 @@ pub(super) mod tests { }); } + #[gpui::test] + fn paste_preview_background_does_not_activate_overlapping_call_control( + cx: &mut gpui::TestAppContext, + ) { + let (mut cx, app) = connected_app_fixture(cx); + cx.simulate_resize(gpui::size(gpui::px(1600.), gpui::px(1100.))); + cx.run_until_parked(); + start_outgoing_call(&mut cx, &app, "test-call-before-preview"); + let baseline_point = outgoing_call_cancel_point(&mut cx); + cx.simulate_click(baseline_point, gpui::Modifiers::default()); + cx.run_until_parked(); + cx.read(|cx| assert!(app.read(cx).call_state(cx).stage().is_none())); + + start_outgoing_call(&mut cx, &app, "test-call-under-preview"); + let call_cancel_point = outgoing_call_cancel_point(&mut cx); + + cx.update(|_window, cx| { + app.update(cx, |app, cx| { + assert!(app.open_confirmation( + "peer@example.invalid".into(), + None, + crate::platform::picker::Chosen { + files: vec![crate::platform::picker::Picked::automatic( + "pasted.png".into(), + "image/png".into(), + one_pixel_png(), + )], + refused: Vec::new(), + }, + cx, + )); + }); + }); + cx.run_until_parked(); + cx.update(|window, cx| window.draw(cx).clear(cx)); + + cx.read(|cx| { + let app = app.read(cx); + assert!(app.call_state(cx).stage().is_some()); + assert!(app.paste_preview.is_some()); + assert!( + !app.keyboard_surfaces.call_card, + "a confirmation modal must temporarily hide the interactive call card" + ); + assert_eq!(app.keyboard_owner, Some(KeyboardOwner::PastePreview)); + }); + let preview = cx + .debug_bounds("paste-preview") + .expect("the modal preview must cover the window"); + assert!( + preview.contains(&call_cancel_point), + "the modal covers the previously verified call control point" + ); + + cx.simulate_click(call_cancel_point, gpui::Modifiers::default()); + cx.run_until_parked(); + cx.update(|window, cx| window.draw(cx).clear(cx)); + + cx.read(|cx| { + let app = app.read(cx); + assert!( + app.paste_preview.is_some(), + "background click keeps preview open" + ); + assert!( + app.call_state(cx).stage().is_some(), + "the overlapping call control must not receive the modal click" + ); + assert!( + !app.keyboard_surfaces.call_card, + "the card stays suppressed while the modal is open" + ); + }); + + let cancel = cx + .debug_bounds("paste-preview-cancel") + .expect("cancel control"); + cx.simulate_click(cancel.center(), gpui::Modifiers::default()); + cx.run_until_parked(); + cx.update(|window, cx| window.draw(cx).clear(cx)); + cx.read(|cx| { + let app = app.read(cx); + assert!(app.paste_preview.is_none()); + assert!(app.call_state(cx).stage().is_some()); + assert_eq!( + app.keyboard_owner, + Some(KeyboardOwner::RingingCall("test-call-under-preview".into())) + ); + assert!( + app.keyboard_surfaces.call_card, + "the card returns after close" + ); + }); + + cx.simulate_click(call_cancel_point, gpui::Modifiers::default()); + cx.run_until_parked(); + cx.read(|cx| assert!(app.read(cx).call_state(cx).stage().is_none())); + } + + #[gpui::test] + fn media_viewer_scrim_does_not_activate_overlapping_call_control( + cx: &mut gpui::TestAppContext, + ) { + let (mut cx, app) = connected_app_fixture(cx); + cx.simulate_resize(gpui::size(gpui::px(1600.), gpui::px(1100.))); + cx.run_until_parked(); + + start_outgoing_call(&mut cx, &app, "test-call-before-viewer"); + let baseline_point = outgoing_call_cancel_point(&mut cx); + cx.simulate_click(baseline_point, gpui::Modifiers::default()); + cx.run_until_parked(); + cx.read(|cx| assert!(app.read(cx).call_state(cx).stage().is_none())); + + start_outgoing_call(&mut cx, &app, "test-call-under-viewer"); + let call_cancel_point = outgoing_call_cancel_point(&mut cx); + cx.update(|window, cx| { + app.update(cx, |app, cx| { + let mut photo = ChatMessage::new_incoming( + "viewer-photo".into(), + "peer@example.invalid".into(), + String::new(), + ); + photo.media = Some( + oxidezap_core::MediaContent::image( + std::sync::Arc::new(one_pixel_png()), + "image/png".into(), + false, + ) + .with_size(Some(1), Some(1)), + ); + Arc::make_mut(&mut app.chats[0]).messages.push(photo); + app.open_media_viewer("viewer-photo", window, cx); + }); + }); + cx.run_until_parked(); + cx.update(|window, cx| window.draw(cx).clear(cx)); + + cx.read(|cx| { + let app = app.read(cx); + assert!(app.media_viewer(cx).is_some()); + assert!(app.call_state(cx).stage().is_some()); + assert!( + !app.keyboard_surfaces.call_card, + "a fullscreen viewer must temporarily hide the interactive call card" + ); + assert_eq!(app.keyboard_owner, Some(KeyboardOwner::Viewer)); + }); + + cx.simulate_click(call_cancel_point, gpui::Modifiers::default()); + cx.run_until_parked(); + cx.update(|window, cx| window.draw(cx).clear(cx)); + cx.read(|cx| { + let app = app.read(cx); + assert!( + app.media_viewer(cx).is_none(), + "the scrim keeps its close action" + ); + assert!( + app.call_state(cx).stage().is_some(), + "the call control behind the scrim must not receive its click" + ); + assert_eq!( + app.keyboard_owner, + Some(KeyboardOwner::RingingCall("test-call-under-viewer".into())) + ); + assert!( + app.keyboard_surfaces.call_card, + "the card returns after close" + ); + }); + + cx.simulate_click(call_cancel_point, gpui::Modifiers::default()); + cx.run_until_parked(); + cx.read(|cx| assert!(app.read(cx).call_state(cx).stage().is_none())); + } + #[gpui::test] fn escape_cancels_preview_without_sending_and_restores_composer(cx: &mut gpui::TestAppContext) { let (mut cx, app) = paste_preview_fixture(cx); diff --git a/crates/gui/src/app/calls_ctl.rs b/crates/gui/src/app/calls_ctl.rs index ad1c4752..0658f024 100644 --- a/crates/gui/src/app/calls_ctl.rs +++ b/crates/gui/src/app/calls_ctl.rs @@ -962,6 +962,9 @@ impl WhatsAppApp { /// Take the daemon's call state as authoritative, and write down whatever /// ended on the way. pub(super) fn adopt_calls(&mut self, mut calls: CallState, cx: &mut Context) { + if self.is_offline() { + calls.end_all(); + } // Named here rather than inside the entity, because a caller's name // comes from this window's chat list. self.name_callers(&mut calls); diff --git a/crates/gui/src/app/events.rs b/crates/gui/src/app/events.rs index bc728e83..05953515 100644 --- a/crates/gui/src/app/events.rs +++ b/crates/gui/src/app/events.rs @@ -227,6 +227,7 @@ impl WhatsAppApp { cx.notify(); } UiEvent::Connected => { + self.recovery.update(cx, |recovering, _| recovering.stop()); self.app_state = AppState::Connected; cx.notify(); } @@ -237,9 +238,16 @@ impl WhatsAppApp { cx.notify(); } UiEvent::Disconnected(reason) => { - // Nothing diagnosed it, so it is the outage the screen was - // written for. - self.connection_ended(oxidezap_core::Fault::unreachable(reason), cx); + // WhatsApp disconnected, but the daemon connection carrying + // this event is still alive. Its supervisor owns recovery. + warn!("WhatsApp disconnected: {reason}"); + self.leave_connected_view(None, cx); + self.recovery.update(cx, |recovering, _| recovering.stop()); + self.app_state = AppState::Offline; + let mut calls = self.calls.read(cx).state().clone(); + calls.end_all(); + self.adopt_calls(calls, cx); + cx.notify(); } UiEvent::Error(msg) => { self.connection_ended(oxidezap_core::Fault::unreachable(msg), cx); diff --git a/crates/gui/src/app/media_ctl.rs b/crates/gui/src/app/media_ctl.rs index b16cf5b9..5c285023 100644 --- a/crates/gui/src/app/media_ctl.rs +++ b/crates/gui/src/app/media_ctl.rs @@ -28,6 +28,20 @@ impl PendingAudioPreparation { fn resume_snapshot(&self) -> (f32, bool) { (self.position, self.was_playing) } + + /// Toggle the saved intent and return it when the browser player has + /// already accepted the clip and is still decoding it. + fn toggle_playback(&mut self, player_loading: bool) -> Option { + self.was_playing = !self.was_playing; + player_loading.then_some(self.was_playing) + } + + /// Update the saved seek position and return it when the player can bank + /// the seek while its browser decode is in flight. + fn seek(&mut self, fraction: f32, player_loading: bool) -> Option { + self.position = fraction.clamp(0.0, 1.0); + player_loading.then_some(self.position) + } } fn media_identity(data: &Arc>) -> usize { @@ -324,7 +338,10 @@ impl WhatsAppApp { .as_mut() .filter(|pending| pending.message_id == message_id) { - pending.position = fraction.clamp(0.0, 1.0); + let seek = pending.seek(fraction, self.audio_player.is_loading()); + if let Some(position) = seek { + self.audio_player.seek(position); + } cx.notify(); return; } @@ -606,10 +623,17 @@ impl WhatsAppApp { .as_mut() .filter(|pending| pending.message_id == message_id) { - // A native preparation has no stream to pause yet. Flip the - // intent stored beside the worker; completion will install the - // samples at the saved position and honour this pause. - pending.was_playing = !pending.was_playing; + // A native preparation has no stream to pause yet, so completion + // will honour the updated snapshot. WebAudio has already accepted + // the clip while its decode is pending; forward the same intent so + // the player applies it as soon as the buffer arrives. + if let Some(should_play) = pending.toggle_playback(self.audio_player.is_loading()) { + if should_play { + self.audio_player.resume(); + } else { + self.audio_player.pause(); + } + } cx.notify(); return; } @@ -645,7 +669,13 @@ impl WhatsAppApp { .as_mut() .filter(|pending| pending.message_id == message_id) { - pending.was_playing = !pending.was_playing; + if let Some(should_play) = pending.toggle_playback(self.audio_player.is_loading()) { + if should_play { + self.audio_player.resume(); + } else { + self.audio_player.pause(); + } + } cx.notify(); return; } @@ -1638,6 +1668,21 @@ mod tests { assert!(pending.matches("voice-1", 7)); } + #[test] + fn pending_audio_controls_forward_only_while_player_is_loading() { + let mut pending = pending(7, 0.25, true); + + assert_eq!(pending.toggle_playback(false), None); + assert_eq!(pending.resume_snapshot(), (0.25, false)); + assert_eq!(pending.toggle_playback(true), Some(true)); + assert_eq!(pending.toggle_playback(true), Some(false)); + + assert_eq!(pending.seek(-0.5, false), None); + assert_eq!(pending.resume_snapshot(), (0.0, false)); + assert_eq!(pending.seek(1.5, true), Some(1.0)); + assert_eq!(pending.resume_snapshot(), (1.0, false)); + } + #[test] fn an_old_preparation_cannot_install_after_a_new_epoch() { let pending = pending(8, 0.25, true); diff --git a/crates/gui/src/app/mod.rs b/crates/gui/src/app/mod.rs index 003e1eac..b993f8d0 100644 --- a/crates/gui/src/app/mod.rs +++ b/crates/gui/src/app/mod.rs @@ -2340,17 +2340,15 @@ impl WhatsAppApp { self.search.update(cx, |search, cx| search.forget(cx)); } - /// Everything that has to stop when the connected view goes away. + /// Stop transient controls tied to a live connection when it goes away. /// /// The controls that stop them are drawn by that view, so anything still /// running when it is replaced has no way to be stopped: a recording /// holds the microphone open, ticks every 100ms and grows a buffer, and /// a voice note plays on over a screen that is now an error message. - /// Three transitions leave that view — a disconnect, an error and a - /// logout — which is three chances to forget, so it is one method. - /// - /// Not [`AppState::Offline`]: that keeps the conversation on screen and - /// only refuses to send. + /// A disconnect keeps cached history on screen in [`AppState::Offline`], + /// while an error or logout replaces it. In all three cases, controls + /// whose work belongs to the live connection must stop here. fn leave_connected_view(&mut self, window: Option<&mut Window>, cx: &mut Context) { // The window-scoped text selection outlives the conversation's rows. // Clear it with the rest of the connected view so Copy cannot expose @@ -4316,15 +4314,6 @@ impl Render for WhatsAppApp { // it, computed before the borrow below. let connected = matches!(self.app_state, AppState::Connected | AppState::Offline); - // Outside the Settings-versus-conversation branch on purpose. The - // card and the focus it takes were built by the conversation view - // alone, so a call arriving while Settings was open rang at the far - // end with nothing on screen to answer or refuse it — and no working - // shortcut either — until the user happened to close Settings. - let call_overlay = connected - .then(|| render_call_overlay(self, window, cx)) - .flatten(); - let paste_preview = self.paste_preview.as_ref().map(|preview| { render_paste_preview( PastePreviewProps { @@ -4350,11 +4339,25 @@ impl Render for WhatsAppApp { .into_any_element() }); + // Outside the Settings-versus-conversation branch on purpose. The + // card and the focus it takes were built by the conversation view + // alone, so a call arriving while Settings was open rang at the far + // end with nothing on screen to answer or refuse it — and no working + // shortcut either — until the user happened to close Settings. + // Modal surfaces own mouse and keyboard input while open. Keep the + // call state alive underneath them, but do not render an interactive + // call card above a fullscreen viewer or confirmation surface. + let paste_preview_open = paste_preview.is_some(); + let message_delete_open = message_delete.is_some(); + let media_viewer_open = self.viewer.read(cx).showing().is_some(); + let modal_open = paste_preview_open || message_delete_open || media_viewer_open; + let call_overlay = (connected && !modal_open) + .then(|| render_call_overlay(self, window, cx)) + .flatten(); + // The card is the one surface the root draws itself, so it is the // one the root answers for. let call_card = call_overlay.is_some(); - let paste_preview_open = paste_preview.is_some(); - let message_delete_open = message_delete.is_some(); // Above the call card as well as the body: a notice raised by // something the call did is about the call, and a card that covered @@ -4362,32 +4365,50 @@ impl Render for WhatsAppApp { // out of here, so the clock taking a line down repaints the stack and // asks nothing of the conversation underneath it; the stack draws // nothing at all while it is empty. - root.child(body.cached(gpui::StyleRefinement::default().size_full())) - .children(paste_preview) - .children(message_delete) - .children(call_overlay) - .child(self.notices().clone()) - // Cached views report their surfaces in prepaint. Move focus after - // drawing, when focus can schedule the frame that delivers blur/focus. - .child( - gpui::canvas( - move |_, window, cx| { - entity.update(cx, |app, _| { - app.keyboard_surfaces.call_card = call_card; - app.keyboard_surfaces.paste_preview = paste_preview_open; - app.keyboard_surfaces.message_delete = message_delete_open; - }); - let entity = entity.downgrade(); - window.defer(cx, move |window, cx| { - let _ = entity.update(cx, |app, cx| { - app.sync_overlay_focus(window, cx); - }); - }); - }, - |_, (), _, _| {}, + let conversation_strip = !self.showing_settings(cx) + && self.destination == Destination::Chats + && self.responsive_layout(window, cx).show_chat_area() + && self.selected_chat_jid().is_some(); + let offline_banner = (self.is_offline() && !conversation_strip) + .then(|| crate::views::render_offline_strip(entity.clone(), cx.product().metrics, cx)); + root.child( + div() + .size_full() + .flex() + .flex_col() + .child( + div() + .flex_1() + .min_h_0() + .child(body.cached(gpui::StyleRefinement::default().size_full())), ) - .absolute(), + .children(offline_banner), + ) + .children(paste_preview) + .children(message_delete) + .children(call_overlay) + .child(self.notices().clone()) + // Cached views report their surfaces in prepaint. Move focus after + // drawing, when focus can schedule the frame that delivers blur/focus. + .child( + gpui::canvas( + move |_, window, cx| { + entity.update(cx, |app, _| { + app.keyboard_surfaces.call_card = call_card; + app.keyboard_surfaces.paste_preview = paste_preview_open; + app.keyboard_surfaces.message_delete = message_delete_open; + }); + let entity = entity.downgrade(); + window.defer(cx, move |window, cx| { + let _ = entity.update(cx, |app, cx| { + app.sync_overlay_focus(window, cx); + }); + }); + }, + |_, (), _, _| {}, ) + .absolute(), + ) } } diff --git a/crates/gui/src/app/recovery.rs b/crates/gui/src/app/recovery.rs index e74987d1..fb09dedc 100644 --- a/crates/gui/src/app/recovery.rs +++ b/crates/gui/src/app/recovery.rs @@ -136,6 +136,40 @@ impl Recovering { } impl WhatsAppApp { + /// Retry WhatsApp through the live IPC connection. Only failure to reach + /// that service uses the separate front-end attach recovery. + pub fn retry_whatsapp_connection(&mut self, cx: &mut Context) { + let Some(client) = self.client.as_ref() else { + self.retry_connection(cx); + return; + }; + let answer = client.reconnect_whatsapp(); + cx.spawn(async move |app: WeakEntity, cx| { + let result = answer.await; + let _ = app.update(cx, |app, cx| app.finish_whatsapp_retry(result, cx)); + }) + .detach(); + } + + fn finish_whatsapp_retry( + &mut self, + result: Result, tokio::sync::oneshot::error::RecvError>, + cx: &mut Context, + ) { + if !self.is_offline() { + return; + } + match result { + // Accepted means the existing supervisor was woken, not that + // WhatsApp is online yet. Only its Connected event says that. + Ok(Ok(())) => {} + // A daemon refusal is an answer over healthy IPC, not a reason + // to attach another reader to the same ended account session. + Ok(Err(failure)) => self.notify_user(failure.detail, notices::Tone::Problem, cx), + Err(_) => self.retry_connection(cx), + } + } + /// Seconds until the automatic retry, or `None` when none is scheduled. pub fn retry_countdown(&self, cx: &App) -> Option { self.recovery.read(cx).countdown_secs() @@ -207,6 +241,63 @@ impl WhatsAppApp { mod tests { use super::*; + #[gpui::test] + fn a_rejected_whatsapp_retry_does_not_reattach_healthy_ipc(cx: &mut gpui::TestAppContext) { + let app = cx.new(|cx| { + let mut app = WhatsAppApp::new(cx); + app.app_state = AppState::Offline; + app + }); + app.update(cx, |app, cx| { + app.finish_whatsapp_retry(Ok(Err(Failure::permanent("session ended"))), cx); + assert!(app.is_offline()); + assert!(app.reconnect_task.is_none()); + assert!(app.notices.read(cx).has_problem("session ended")); + app.finish_whatsapp_retry(Ok(Ok(())), cx); + assert!(app.is_offline(), "only Connected can restore sending"); + }); + } + + #[gpui::test] + fn whatsapp_loss_preserves_history_and_waits_on_the_existing_session( + cx: &mut gpui::TestAppContext, + ) { + let app = cx.new(|cx| { + let mut app = WhatsAppApp::new(cx); + app.app_state = AppState::Connected; + app.chats + .push(Arc::new(Chat::new("999000@s.whatsapp.net".into()))); + app + }); + app.update(cx, |app, cx| { + let mut ringing = CallState::default(); + ringing.set_incoming(oxidezap_core::IncomingCall { + call_id: "BEFORE-DROP".into(), + caller_name: "Example".into(), + caller_jid: "999000@s.whatsapp.net".into(), + is_video: false, + is_offline: false, + received_at: wacore::time::now_utc(), + }); + app.adopt_calls(ringing.clone(), cx); + assert!(app.calls.read(cx).state().incoming().is_some()); + app.handle_event(UiEvent::Disconnected("synthetic drop".into()), cx); + assert!(app.is_offline()); + assert!(!app.can_send()); + assert_eq!(app.chats.len(), 1); + assert!(app.calls.read(cx).state().stage().is_none()); + assert!(app.retry_countdown(cx).is_none()); + assert!(app.reconnect_task.is_none()); + app.adopt_calls(ringing, cx); + assert!( + app.calls.read(cx).state().stage().is_none(), + "offline replay cannot reopen the card" + ); + app.handle_event(UiEvent::Connected, cx); + assert!(app.can_send()); + }); + } + /// The countdown is derived from a deadline rather than decremented, so a /// tick that runs late lands on the right number instead of the number of /// ticks that happened to run. diff --git a/crates/gui/src/components/call_card/ringing.rs b/crates/gui/src/components/call_card/ringing.rs index 892a6855..a9bd856f 100644 --- a/crates/gui/src/components/call_card/ringing.rs +++ b/crates/gui/src/components/call_card/ringing.rs @@ -7,7 +7,7 @@ //! gave both a red `Cancel`, which said the same word for leaving a call //! unanswered and for hanging up on someone. -use gpui::{App, Entity, IntoElement, ParentElement, Styled, div}; +use gpui::{App, Entity, InteractiveElement, IntoElement, ParentElement, Styled, div}; use gpui_component::ActiveTheme as _; use gpui_component::Icon; use gpui_component::button::{Button, ButtonVariants as _}; @@ -116,6 +116,7 @@ pub fn outgoing( )) .child( div() + .debug_selector(|| "call-cancel".into()) .w_full() .flex() .justify_center() diff --git a/crates/gui/src/components/media_viewer.rs b/crates/gui/src/components/media_viewer.rs index e3007f70..8d028eb3 100644 --- a/crates/gui/src/components/media_viewer.rs +++ b/crates/gui/src/components/media_viewer.rs @@ -83,6 +83,17 @@ pub fn render_media_viewer( // conversation nobody could see. Swallowing it here is what makes the // picture the only thing on screen that responds. .on_scroll_wheel(|_, _window, cx| cx.stop_propagation()) + // Mouse-up is where GPUI synthesizes click handlers. Stop both halves + // of the pointer sequence at the modal, after a child control has had + // its own chance to act, and occlude overlapping hitboxes behind it. + .on_mouse_down(gpui::MouseButton::Left, |_, _window, cx| { + cx.stop_propagation(); + }) + .on_mouse_up(gpui::MouseButton::Left, |_, _window, cx| { + cx.stop_propagation(); + }) + .on_click(|_, _window, cx| cx.stop_propagation()) + .occlude() .on_action(move |_: &ViewerPrev, _window, cx| { key_prev_entity.update(cx, |app, cx| app.step_media_viewer(false, cx)); }) diff --git a/crates/gui/src/components/paste_preview.rs b/crates/gui/src/components/paste_preview.rs index dd6d39c9..44d78b14 100644 --- a/crates/gui/src/components/paste_preview.rs +++ b/crates/gui/src/components/paste_preview.rs @@ -88,6 +88,11 @@ pub fn render_paste_preview(props: PastePreviewProps<'_>, cx: &App) -> impl Into .on_mouse_down(gpui::MouseButton::Left, |_, _window, cx| { cx.stop_propagation(); }) + .on_mouse_up(gpui::MouseButton::Left, |_, _window, cx| { + cx.stop_propagation(); + }) + .on_click(|_, _window, cx| cx.stop_propagation()) + .occlude() .child( div() .text_size(metrics.text_title()) diff --git a/crates/gui/src/platform/picker.rs b/crates/gui/src/platform/picker.rs index dfe627ad..a3fd115b 100644 --- a/crates/gui/src/platform/picker.rs +++ b/crates/gui/src/platform/picker.rs @@ -251,7 +251,14 @@ fn selected_kind_and_mime( let inferred_from_name = declared_mime.is_empty() || (category == Some(AttachmentCategory::PhotosVideos) && generic_mime); let candidate = if inferred_from_name { - mime_for_name(file_name) + let named_mime = mime_for_name(file_name); + if category == Some(AttachmentCategory::PhotosVideos) + && named_mime == "application/octet-stream" + { + image_mime_from_bytes(bytes).unwrap_or(named_mime) + } else { + named_mime + } } else { declared_mime }; @@ -578,35 +585,43 @@ mod imp { ); let photo = directory.join(format!("{marker}.png")); let heic = directory.join(format!("{marker}.heic")); + let extensionless = directory.join(format!("{marker}-without-extension")); std::fs::write(&photo, b"\x89PNG\r\n\x1a\nrest").expect("write photo fixture"); std::fs::write(&heic, b"original HEIC bytes").expect("write HEIC fixture"); + std::fs::write(&extensionless, b"\x89PNG\r\n\x1a\nrest") + .expect("write extensionless photo fixture"); let media = read_all( - &[photo.clone(), heic.clone()], + &[photo.clone(), extensionless.clone(), heic.clone()], Some(AttachmentCategory::PhotosVideos), &mut super::super::Budget::default(), ); let documents = read_all( - &[photo.clone(), heic.clone()], + &[photo.clone(), extensionless.clone(), heic.clone()], Some(AttachmentCategory::Document), &mut super::super::Budget::default(), ); std::fs::remove_file(photo).expect("remove photo fixture"); std::fs::remove_file(heic).expect("remove HEIC fixture"); + std::fs::remove_file(extensionless).expect("remove extensionless fixture"); - assert_eq!(media.files.len(), 1); + assert_eq!(media.files.len(), 2); assert_eq!(media.files[0].kind, oxidezap_core::OutgoingMedia::Image); + assert_eq!(media.files[1].kind, oxidezap_core::OutgoingMedia::Image); + assert_eq!(media.files[1].mime_type, "image/png"); assert_eq!(media.refused.len(), 1); assert!(media.refused[0].contains(".heic")); - assert_eq!(documents.files.len(), 2); + assert_eq!(documents.files.len(), 3); assert!( documents .files .iter() .all(|file| file.kind == oxidezap_core::OutgoingMedia::Document) ); - assert_eq!(documents.files[1].bytes, b"original HEIC bytes"); - assert_eq!(documents.files[1].mime_type, "image/heic"); + assert_eq!(documents.files[1].bytes, b"\x89PNG\r\n\x1a\nrest"); + assert_eq!(documents.files[1].mime_type, "application/octet-stream"); + assert_eq!(documents.files[2].bytes, b"original HEIC bytes"); + assert_eq!(documents.files[2].mime_type, "image/heic"); } } } @@ -1020,6 +1035,17 @@ mod tests { ), Ok((OutgoingMedia::Image, "image/png".into())) ); + for declared_mime in ["", "application/octet-stream"] { + assert_eq!( + selected_kind_and_mime( + "photo-without-extension", + declared_mime, + png, + Some(AttachmentCategory::PhotosVideos), + ), + Ok((OutgoingMedia::Image, "image/png".into())) + ); + } assert!( selected_kind_and_mime( "false.png", @@ -1057,6 +1083,15 @@ mod tests { ), Ok((OutgoingMedia::Document, "application/octet-stream".into())) ); + assert_eq!( + selected_kind_and_mime( + "photo-without-extension", + "application/octet-stream", + png, + Some(AttachmentCategory::Document), + ), + Ok((OutgoingMedia::Document, "application/octet-stream".into())) + ); assert_eq!( selected_kind_and_mime("photo.png", "", png, Some(AttachmentCategory::Document)), Ok((OutgoingMedia::Document, "image/png".into())) diff --git a/crates/gui/src/session/frames.rs b/crates/gui/src/session/frames.rs index 101c0751..0b614e0e 100644 --- a/crates/gui/src/session/frames.rs +++ b/crates/gui/src/session/frames.rs @@ -170,7 +170,7 @@ impl<'a> Frames<'a> { // nothing to compare against. DaemonMessage::Hello { protocol, snapshot } if protocol == PROTOCOL_VERSION => { self.applied = snapshot.version; - self.follow_calls(&snapshot.calls); + self.follow_calls(&snapshot_calls(&snapshot)); for event in catch_up(&snapshot) { self.publish(event)?; } @@ -745,11 +745,19 @@ pub(super) fn catch_up(snapshot: &StateSnapshot) -> Vec { // Whatever is happening on the call front, as state. The offer for a // ringing call went out before this window existed, and a call this // account placed was never an event at all. - events.push(FromDaemon::Calls(Box::new(snapshot.calls.clone()))); + events.push(FromDaemon::Calls(Box::new(snapshot_calls(snapshot)))); events.push(FromDaemon::Account(snapshot.account.clone())); events } +fn snapshot_calls(snapshot: &StateSnapshot) -> CallState { + let mut calls = snapshot.calls.clone(); + if !snapshot.connection.is_connected() { + calls.end_all(); + } + calls +} + /// How many rows a snapshot may paint. /// /// The window the session's own load fills (`HISTORY_CHAT_LIMIT`). The daemon @@ -1599,6 +1607,19 @@ mod tests { } } + #[test] + fn an_offline_snapshot_cannot_resurrect_a_call_card() { + let mut snapshot = snapshot_of(Vec::new()); + snapshot.connection = ConnectionState::Disconnected { + reason: "synthetic drop".into(), + }; + snapshot.calls = video_call("OLD"); + let events = catch_up(&snapshot); + assert!(events.iter().any(|event| matches!(event, FromDaemon::Session(event) if matches!(&**event, UiEvent::Disconnected(_))))); + assert!(events.iter().any(|event| matches!(event, FromDaemon::Calls(calls) if calls.stage().is_none() && calls.waiting().is_none()))); + assert!(snapshot_calls(&snapshot).stage().is_none()); + } + fn summary(jid: &str, name: &str, unread: u32) -> ChatSummary { ChatSummary { jid: jid.to_string(), diff --git a/crates/gui/src/session/mod.rs b/crates/gui/src/session/mod.rs index 4a23fef4..7728fa62 100644 --- a/crates/gui/src/session/mod.rs +++ b/crates/gui/src/session/mod.rs @@ -1366,6 +1366,14 @@ impl SessionHandle { self.tell(ClientRequest::ReloadHistory); } + /// Ask the existing daemon session to retry WhatsApp, rather than opening + /// another IPC connection onto the same reconnect wait. + pub fn reconnect_whatsapp(&self) -> oneshot::Receiver> { + let (tx, rx) = oneshot::channel(); + self.ask(ClientRequest::ReconnectSession, Awaiting::Mutation(tx)); + rx + } + /// Tell the daemon these status updates have been watched. /// /// Not `mark_chat_read` on the broadcast: that clears one chat, and the diff --git a/crates/gui/src/views/chat.rs b/crates/gui/src/views/chat.rs index d81e5af7..af56141e 100644 --- a/crates/gui/src/views/chat.rs +++ b/crates/gui/src/views/chat.rs @@ -500,7 +500,7 @@ fn render_chat_area( /// /// Not a disabled field: a greyed-out composer says "you cannot type here" /// and stops. What the reader needs is why, and the way back. -fn render_offline_strip( +pub(crate) fn render_offline_strip( entity: Entity, metrics: Metrics, cx: &App, @@ -526,7 +526,7 @@ fn render_offline_strip( .min_w_0() .text_size(metrics.text_small()) .text_color(cx.theme().muted_foreground) - .child("Offline. You can read this conversation, but not send in it."), + .child("Offline. Your messages stay readable. Reconnect to send."), ) .child( Button::new("reconnect") @@ -534,7 +534,7 @@ fn render_offline_strip( .ghost() .cursor_pointer() .on_click(move |_, _window, cx| { - entity.update(cx, |app, cx| app.retry_connection(cx)); + entity.update(cx, |app, cx| app.retry_whatsapp_connection(cx)); }), ) } diff --git a/crates/gui/src/views/mod.rs b/crates/gui/src/views/mod.rs index 35f5f1e5..50b50187 100644 --- a/crates/gui/src/views/mod.rs +++ b/crates/gui/src/views/mod.rs @@ -13,6 +13,7 @@ mod logged_out; pub mod pairing; mod settings; +pub(crate) use chat::render_offline_strip; pub use chat::{render_call_overlay, render_connected_view}; pub use error::{render_error_view, render_refused_view}; pub use loading::{render_connecting_view, render_loading_view, render_syncing_view}; diff --git a/crates/ipc/src/protocol.rs b/crates/ipc/src/protocol.rs index 56648160..4a143b54 100644 --- a/crates/ipc/src/protocol.rs +++ b/crates/ipc/src/protocol.rs @@ -1010,6 +1010,9 @@ pub enum ClientRequest { /// client can, so it starts over. Sent for it automatically when it /// attaches, which is the same situation. ReloadHistory, + /// Retry the account's WhatsApp connection now, preserving its identity, + /// store and supervised session. Valid while WhatsApp is disconnected. + ReconnectSession, /// Wipe local state and pair again. /// /// A server 401 means the stored credentials are dead and reconnecting diff --git a/crates/ipc/src/transport.rs b/crates/ipc/src/transport.rs index 1a045fe9..903c27f5 100644 --- a/crates/ipc/src/transport.rs +++ b/crates/ipc/src/transport.rs @@ -5,6 +5,9 @@ use std::path::PathBuf; /// Bumped whenever a frame changes shape in a way an older peer would /// misread. The daemon refuses a mismatch rather than guessing. /// +/// 40: `ReconnectSession` interrupts the WhatsApp reconnect wait without +/// replacing the account session. Older daemons cannot act on this request. +/// /// 39: `ClientRequest::RevokeMessage` adds an addressed delete completion /// alongside v38 edits, retaining the v37 community hierarchy fields. /// A v38 daemon does not know the delete request. @@ -283,7 +286,7 @@ use std::path::PathBuf; /// would misparse the first three and not recognise the rest. /// /// [`PairingCode`]: crate::PairingCode -pub const PROTOCOL_VERSION: u32 = 39; +pub const PROTOCOL_VERSION: u32 = 40; /// Where the daemon's web bridge listens when nobody says otherwise. /// diff --git a/crates/ipc/tests/session_frames.rs b/crates/ipc/tests/session_frames.rs index 5d9e765c..4b32a56b 100644 --- a/crates/ipc/tests/session_frames.rs +++ b/crates/ipc/tests/session_frames.rs @@ -8,6 +8,19 @@ use oxidezap_core::{ReceiptType, UiEvent}; use oxidezap_ipc::DaemonMessage; +#[test] +fn reconnect_session_is_an_account_request_and_round_trips() { + let request = oxidezap_ipc::ClientRequest::ReconnectSession; + assert!(request.is_account_request()); + assert!(!request.is_control_request()); + let encoded = serde_json::to_string(&request).unwrap(); + assert_eq!(encoded, "{\"request\":\"reconnect_session\"}"); + assert!(matches!( + serde_json::from_str::(&encoded).unwrap(), + oxidezap_ipc::ClientRequest::ReconnectSession + )); +} + /// One of each, so a variant added to the library is a variant this covers. fn every_receipt_type() -> Vec { vec![ diff --git a/crates/session/src/whatsapp/calls/registry.rs b/crates/session/src/whatsapp/calls/registry.rs index 57278dd3..6d117de4 100644 --- a/crates/session/src/whatsapp/calls/registry.rs +++ b/crates/session/src/whatsapp/calls/registry.rs @@ -33,6 +33,17 @@ fn offered_video(offer: &WaIncomingCall) -> bool { matches!(&offer.action, CallAction::Offer { is_video, .. } if *is_video) } +/// Offers normally stop ringing after about 45 seconds. This deliberately +/// wider window leaves room for delivery latency and modest clock skew while +/// preventing a reconnect backlog from ringing calls from minutes ago. It is +/// a freshness policy, not a protocol TTL: future timestamps are not refused. +const MAX_RINGABLE_OFFER_AGE_MS: i64 = 120_000; + +fn ringable_offer(offer: &WaIncomingCall, now_ms: i64) -> bool { + !offer.offline + && now_ms.saturating_sub(offer.timestamp.timestamp_millis()) <= MAX_RINGABLE_OFFER_AGE_MS +} + fn direct_peer_video(state: VideoState) -> Option { match state { VideoState::Enabled => Some(true), @@ -225,6 +236,9 @@ const ANNOUNCED_ENDINGS: usize = 256; #[derive(Default)] struct Calls { + /// Captured before an event queues, and invalidated on transport loss. + connection_epoch: u64, + disconnected: bool, peer_lanes: HashMap>>, starting: HashMap, next_start: u64, @@ -310,13 +324,122 @@ impl Drop for StartGuard { } impl CallRegistry { - /// Record a ringing offer, so accept and decline have something to act on. - pub(in crate::whatsapp) fn offer(&self, call_id: String, call: Arc) { - self.calls - .lock() - .expect("call registry poisoned") - .pending - .insert(call_id, call); + /// Apply connection boundaries before any asynchronous event lane. Sending + /// the disconnect under this lock orders it against offer publication. + pub(in crate::whatsapp) fn connection_event(&self, event: &Event, ui: &UiEventSender) -> u64 { + let mut calls = self.calls.lock().expect("call registry poisoned"); + let mut retired = Vec::new(); + match event { + Event::Disconnected(_) | Event::LoggedOut(_) => { + calls.connection_epoch = calls + .connection_epoch + .checked_add(1) + .expect("call epoch exhausted"); + calls.disconnected = true; + let retired_ids: Vec<_> = calls + .pending + .keys() + .chain(calls.active.keys()) + .chain(calls.outgoing.keys()) + .chain(calls.in_flight.iter()) + .cloned() + .collect(); + for id in retired_ids { + if calls.announced.insert(id.clone()) { + calls.announced_order.push_back(id); + } + } + while calls.announced_order.len() > ANNOUNCED_ENDINGS { + if let Some(id) = calls.announced_order.pop_front() { + calls.announced.remove(&id); + } + } + calls.pending.clear(); + let accepting: Vec<_> = calls.in_flight.iter().cloned().collect(); + for id in accepting { + calls.cancelled.insert(id, Ending::Remote); + } + retired.extend(calls.active.drain().map(|(_, handle)| handle)); + calls.outgoing.clear(); + calls.peer_lanes.clear(); + if let Event::Disconnected(disconnected) = event { + let _ = ui.send(UiEvent::Disconnected(disconnected.reason.to_string())); + } + } + Event::Connected(_) => { + calls.disconnected = false; + // Offers run on different lanes. Publish readiness here, + // before presence/name refresh awaits, so a valid new call + // cannot reach an interface that still considers us offline. + let _ = ui.send(UiEvent::Connected); + } + _ => {} + } + let epoch = calls.connection_epoch; + drop(calls); + if matches!(event, Event::Disconnected(_) | Event::LoggedOut(_)) { + self.registration.notify_waiters(); + } + for handle in retired { + drop(crate::exec::spawn( + async move { handle.hangup_local().await }, + )); + } + epoch + } + + pub(in crate::whatsapp) fn offer_in_epoch( + &self, + call_id: String, + call: Arc, + epoch: u64, + ) -> bool { + let mut calls = self.calls.lock().expect("call registry poisoned"); + if calls.disconnected + || calls.connection_epoch != epoch + || !ringable_offer(&call, wacore::time::now_millis()) + || calls.announced.contains(&call_id) + || calls.pending.contains_key(&call_id) + || calls.active.contains_key(&call_id) + || calls.in_flight.contains(&call_id) + { + return false; + } + calls.pending.insert(call_id, call); + true + } + + pub(in crate::whatsapp) fn connection_is_current(&self, epoch: u64) -> bool { + let calls = self.calls.lock().expect("call registry poisoned"); + !calls.disconnected && calls.connection_epoch == epoch + } + + /// Recheck after caller-name lookup, with the send in the same section as + /// connection retirement. A stale offer cannot follow its disconnect. + pub(in crate::whatsapp) fn publish_offer( + &self, + offer: &Arc, + call: IncomingCall, + epoch: u64, + ui: &UiEventSender, + ) -> bool { + let mut calls = self.calls.lock().expect("call registry poisoned"); + if calls.disconnected + || calls.connection_epoch != epoch + || !calls + .pending + .get(&call.call_id) + .is_some_and(|pending| Arc::ptr_eq(pending, offer)) + { + return false; + } + // Caller lookup can itself wait through an outage. Recheck age as + // well as generation, and do not leave an unannounced offer to accept. + if !ringable_offer(offer, wacore::time::now_millis()) { + calls.pending.remove(&call.call_id); + return false; + } + ui.send(UiEvent::IncomingCall(call)).is_ok() } /// Take a ringing offer and mark its acceptance as in flight, together. @@ -2876,6 +2999,170 @@ fn log_termination(call_id: &str, outcome: CallTermination) { mod tests { use super::*; + fn replay_offer(id: &str) -> Arc { + let jid: Jid = "999000@s.whatsapp.net".parse().unwrap(); + Arc::new( + WaIncomingCall::builder() + .from(jid.clone()) + .stanza_id(id.to_string()) + .timestamp(wacore::time::now_utc()) + .offline(false) + .action(CallAction::Offer { + call_id: id.into(), + call_creator: jid, + caller_pn: None, + caller_country_code: None, + device_class: None, + joinable: true, + is_video: false, + audio: Vec::new(), + group_jid: None, + }) + .build(), + ) + } + + #[test] + fn offer_freshness_allows_delivery_margin_and_one_sided_clock_skew() { + let now = wacore::time::now_millis(); + let mut offer = (*replay_offer("SKEW")).clone(); + offer.timestamp = wacore::time::from_millis(now - 120_000).unwrap(); + assert!(ringable_offer(&offer, now)); + offer.timestamp = wacore::time::from_millis(now - 121_000).unwrap(); + assert!(!ringable_offer(&offer, now)); + offer.timestamp = wacore::time::from_millis(now + 300_000).unwrap(); + assert!(ringable_offer(&offer, now), "future time is not stale"); + offer.offline = true; + assert!(!ringable_offer(&offer, now)); + } + + #[test] + fn a_historic_offer_first_delivered_after_reconnect_does_not_ring() { + let calls = CallRegistry::default(); + let (ui, mut events) = ui_queue::channel( + Arc::new(tokio::sync::Notify::new()), + Arc::new(ui_queue::HistoryBudget::new()), + ); + let epoch = calls.connection_event( + &Event::Disconnected( + wacore::types::events::Disconnected::builder() + .reason(wacore::net::DisconnectReason::ReadError("offline".into())) + .build(), + ), + &ui, + ); + calls.connection_event( + &Event::Connected(wacore::types::events::Connected::builder().build()), + &ui, + ); + let mut old = (*replay_offer("BACKLOG")).clone(); + old.timestamp = + wacore::time::from_millis(old.timestamp.timestamp_millis() - 300_000).unwrap(); + assert!( + !old.offline, + "an offline flag is not required for freshness" + ); + assert!(!calls.offer_in_epoch("BACKLOG".into(), Arc::new(old), epoch)); + let fresh = replay_offer("CURRENT"); + assert!(calls.offer_in_epoch("CURRENT".into(), fresh.clone(), epoch)); + let current = IncomingCall::new( + "CURRENT", + "Example".into(), + fresh.from.to_string(), + false, + &fresh, + ); + assert!(calls.publish_offer(&fresh, current, epoch, &ui)); + assert!(matches!(events.try_recv(), Ok(UiEvent::Disconnected(_)))); + assert!(matches!(events.try_recv(), Ok(UiEvent::Connected))); + assert!( + matches!(events.try_recv(), Ok(UiEvent::IncomingCall(call)) if call.call_id == "CURRENT") + ); + assert!(events.try_recv().is_err()); + assert!(calls.begin_accept("BACKLOG").is_none()); + } + + #[test] + fn an_offer_that_aged_during_caller_lookup_is_not_published_or_acceptible() { + let calls = CallRegistry::default(); + let (ui, mut events) = ui_queue::channel( + Arc::new(tokio::sync::Notify::new()), + Arc::new(ui_queue::HistoryBudget::new()), + ); + let mut old = (*replay_offer("LOOKUP")).clone(); + old.timestamp = + wacore::time::from_millis(old.timestamp.timestamp_millis() - 300_000).unwrap(); + let old = Arc::new(old); + // The registered offer is the one held across the name lookup; age + // can change without either its identity or its epoch changing. + calls + .calls + .lock() + .unwrap() + .pending + .insert("LOOKUP".into(), old.clone()); + let call = IncomingCall::new( + "LOOKUP", + "Example".into(), + old.from.to_string(), + false, + &old, + ); + assert!(!calls.publish_offer(&old, call, 0, &ui)); + assert!(events.try_recv().is_err()); + assert!(calls.begin_accept("LOOKUP").is_none()); + } + + #[test] + fn an_offer_delayed_across_reconnect_cannot_ring_again() { + let calls = CallRegistry::default(); + let (ui, mut events) = ui_queue::channel( + Arc::new(tokio::sync::Notify::new()), + Arc::new(ui_queue::HistoryBudget::new()), + ); + let old = replay_offer("OLD"); + let queued = replay_offer("QUEUED"); + assert!(calls.offer_in_epoch("OLD".into(), old.clone(), 0)); + calls.mark_accepting("ACCEPTING"); + let disconnected = Event::Disconnected( + wacore::types::events::Disconnected::builder() + .reason(wacore::net::DisconnectReason::ReadError( + "synthetic drop".into(), + )) + .build(), + ); + let next = calls.connection_event(&disconnected, &ui); + assert!(calls.calls.lock().unwrap().pending.is_empty()); + assert_eq!(calls.ending_for("ACCEPTING"), Some(Ending::Remote)); + calls.connection_event( + &Event::Connected(wacore::types::events::Connected::builder().build()), + &ui, + ); + let old_ui = IncomingCall::new("OLD", "Example".into(), old.from.to_string(), false, &old); + assert!(!calls.publish_offer(&old, old_ui, 0, &ui)); + assert!(!calls.connection_is_current(0)); + assert!(calls.connection_is_current(next)); + assert!(!calls.offer_in_epoch("QUEUED".into(), queued, 0)); + assert!(!calls.offer_in_epoch("OLD".into(), old, next)); + assert!(!calls.offer_in_epoch("ACCEPTING".into(), replay_offer("ACCEPTING"), next)); + let fresh = replay_offer("FRESH"); + assert!(calls.offer_in_epoch("FRESH".into(), fresh.clone(), next)); + let fresh_ui = IncomingCall::new( + "FRESH", + "Example".into(), + fresh.from.to_string(), + false, + &fresh, + ); + assert!(calls.publish_offer(&fresh, fresh_ui, next, &ui)); + assert!(matches!(events.try_recv(), Ok(UiEvent::Disconnected(_)))); + assert!(matches!(events.try_recv(), Ok(UiEvent::Connected))); + assert!( + matches!(events.try_recv(), Ok(UiEvent::IncomingCall(call)) if call.call_id == "FRESH") + ); + assert!(events.try_recv().is_err()); + } + #[cfg_attr(not(target_family = "wasm"), test)] #[cfg_attr(target_family = "wasm", wasm_bindgen_test::wasm_bindgen_test)] fn source_video_event_keeps_the_operation_and_rejects_legacy() { diff --git a/crates/session/src/whatsapp/calls/registry/acceptance_fixture.rs b/crates/session/src/whatsapp/calls/registry/acceptance_fixture.rs index 0d48c0d9..a027cf54 100644 --- a/crates/session/src/whatsapp/calls/registry/acceptance_fixture.rs +++ b/crates/session/src/whatsapp/calls/registry/acceptance_fixture.rs @@ -127,20 +127,24 @@ async fn dispatch( let calls = calls.clone(); let ui = ui.clone(); let names = names.clone(); - move |event| { + move |event, epoch| { let client = client.clone(); let calls = calls.clone(); let ui = ui.clone(); let names = names.clone(); let processed = processed.clone(); async move { - WhatsAppClient::handle_event(event, client, ui, calls, names, None, None, None) - .await; + WhatsAppClient::handle_event( + event, client, ui, calls, names, None, None, None, epoch, + ) + .await; processed.send(()).await.unwrap(); } } }, stopping, + calls.clone(), + ui.clone(), ); let count = relevant.len(); for event in relevant { diff --git a/crates/session/src/whatsapp/lanes.rs b/crates/session/src/whatsapp/lanes.rs index 58fdb309..d7067ad8 100644 --- a/crates/session/src/whatsapp/lanes.rs +++ b/crates/session/src/whatsapp/lanes.rs @@ -14,6 +14,7 @@ use whatsapp_rust::wacore::types::events::Event; use whatsapp_rust::wacore_binary::JidExt; use whatsapp_rust::wacore_binary::jid::Jid; +use super::{UiEventSender, calls::CallRegistry}; use crate::names::NameBook; /// How many lanes events about a subject are spread across. @@ -36,19 +37,26 @@ const LANE_CAPACITY: usize = 64; /// naming neither is session-wide and gets a lane of its own, so a pairing /// code never waits behind a conversation. pub(super) struct EventLanes { - lanes: Vec>>, + lanes: Vec, u64)>>, stopping: tokio::sync::watch::Receiver<()>, + calls: CallRegistry, + ui: UiEventSender, } impl EventLanes { - pub(super) fn new(handle: F, stopping: tokio::sync::watch::Receiver<()>) -> Self + pub(super) fn new( + handle: F, + stopping: tokio::sync::watch::Receiver<()>, + calls: CallRegistry, + ui: UiEventSender, + ) -> Self where - F: Fn(Arc) -> Fut + Clone + crate::exec::MaybeSend + 'static, + F: Fn(Arc, u64) -> Fut + Clone + crate::exec::MaybeSend + 'static, Fut: Future + crate::exec::MaybeSend + 'static, { let lanes = (0..=EVENT_LANES) .map(|_| { - let (tx, mut rx) = mpsc::channel::>(LANE_CAPACITY); + let (tx, mut rx) = mpsc::channel::<(Arc, u64)>(LANE_CAPACITY); let handle = handle.clone(); let mut stopping = stopping.clone(); crate::exec::spawn_owned(async move { @@ -68,13 +76,18 @@ impl EventLanes { // client and its store alive. _ = stopping.changed() => return, }; - handle(event).await; + handle(event.0, event.1).await; } }); tx }) .collect(); - Self { lanes, stopping } + Self { + lanes, + stopping, + calls, + ui, + } } /// Dispatch an event and report recoverable events dropped because their @@ -89,6 +102,7 @@ impl EventLanes { event: Arc, ) -> DispatchOutcome { let mut outcome = DispatchOutcome::default(); + let epoch = self.calls.connection_event(&event, &self.ui); // A batch may span chats, and a lane is one chat's order: sent whole // on the first message's lane, a receipt for a later chat in it runs // on that chat's own lane and can overtake the message it answers. @@ -98,7 +112,10 @@ impl EventLanes { for event in split_by_subject(&event) { let lane = lane_for(client, names, &event).await; if recoverable(&event) { - if self.lanes[lane].try_send(Arc::clone(&event)).is_err() { + if self.lanes[lane] + .try_send((Arc::clone(&event), epoch)) + .is_err() + { outcome.dropped_recoverable = true; if let Event::Messages(batch) = &*event { for inbound in batch.iter() { @@ -121,7 +138,7 @@ impl EventLanes { } else { let mut stopping = self.stopping.clone(); tokio::select! { - result = self.lanes[lane].send(event) => { + result = self.lanes[lane].send((event, epoch)) => { if result.is_err() { return outcome; } diff --git a/crates/session/src/whatsapp/mod.rs b/crates/session/src/whatsapp/mod.rs index 0eaee983..87a4a303 100644 --- a/crates/session/src/whatsapp/mod.rs +++ b/crates/session/src/whatsapp/mod.rs @@ -199,6 +199,7 @@ const CONTROL_EVENT_KINDS: &[EventKind] = &[ EventKind::PairingCode, EventKind::PairSuccess, EventKind::Connected, + EventKind::Disconnected, EventKind::LoggedOut, EventKind::SelfPushNameUpdated, EventKind::IncomingCall, @@ -1347,8 +1348,10 @@ impl WhatsAppClient { // PN/LID pairing that decides the lane. let dispatch_client = client.clone(); let dispatch_names = names.clone(); + let dispatch_calls = calls.clone(); + let dispatch_ui = ui_tx.clone(); let mut lanes = EventLanes::new( - move |event| { + move |event, epoch| { let client = client.clone(); let ui_tx = ui_tx.clone(); let calls = calls.clone(); @@ -1366,11 +1369,14 @@ impl WhatsAppClient { Some(reload), Some(resolve_avatars), Some(resolve_chat_names), + epoch, ) .await; } }, stopping.clone(), + dispatch_calls, + dispatch_ui, ); let mut control_open = true; let mut data_open = true; @@ -1632,8 +1638,12 @@ impl WhatsAppClient { reload: Option>, resolve_avatars: Option, resolve_chat_names: Option, + epoch: u64, ) { match &*event { + // The dispatch boundary already retired calls and published this + // before any older call lane can resume its asynchronous work. + Event::Disconnected(_) => {} Event::RawNode(node) => calls.accept_advertisement(node).await, Event::PairingQrCode(qr) => { info!("QR code received"); @@ -1692,7 +1702,9 @@ impl WhatsAppClient { if let Some(resolve) = resolve_chat_names { resolve.new_connection(); } - let _ = ui_tx.send(UiEvent::Connected); + if !calls.connection_is_current(epoch) { + return; + } // Who this device is linked as. Read from the device store // rather than remembered from pairing: a client attaching // after a restart never saw that, and the account row was @@ -1752,9 +1764,11 @@ impl WhatsAppClient { info!("Ignoring offline call {} (stale)", call_id); return; } - info!("Incoming call from {}", call.from.observe()); let offer = Arc::new((**call).clone()); - calls.offer(call_id.clone(), offer.clone()); + if !calls.offer_in_epoch(call_id.clone(), offer.clone(), epoch) { + return; + } + info!("Incoming call from {}", call.from.observe()); // A call has to ring even when nobody can say which of // the caller's two addresses their chat is filed under: // the address as written is no worse than the one the @@ -1777,7 +1791,7 @@ impl WhatsAppClient { *is_video, &offer, ); - let _ = ui_tx.send(UiEvent::IncomingCall(ui_call)); + calls.publish_offer(&offer, ui_call, epoch, &ui_tx); } CallAction::Accept { call_id, .. } => { info!("Call {} accepted by peer", call_id); diff --git a/crates/session/src/whatsapp/mutations.rs b/crates/session/src/whatsapp/mutations.rs index b270b339..e48bb9ed 100644 --- a/crates/session/src/whatsapp/mutations.rs +++ b/crates/session/src/whatsapp/mutations.rs @@ -17,6 +17,25 @@ use crate::exec::Task; /// What a send-like mutation produced: the server-assigned message id. pub type SendId = String; +/// Once the remote delete succeeds, a local writer failure must not invite a +/// retry of a mutation that already happened. Keep the local update ordered +/// after the remote answer, and report its failures through diagnostics. +pub(super) async fn delete_and_record( + delete: impl std::future::Future>, + record: impl FnOnce() -> oxidezap_chat_store::Result<()>, + flush: impl std::future::Future>, +) -> Result<(), String> { + delete.await.map_err(|e| e.to_string())?; + if let Err(e) = record() { + log::warn!("delete succeeded remotely but could not be queued locally: {e}"); + return Ok(()); + } + if let Err(e) = flush.await { + log::warn!("delete succeeded remotely but could not be saved locally: {e}"); + } + Ok(()) +} + fn edit_text_content(original: Option<&wa::Message>, text: String) -> wa::Message { let context = original .map(oxidezap_chat_store::normalized_message) @@ -213,22 +232,19 @@ impl WhatsAppClient { .to_string(), ); } - live.client - .revoke_message( + delete_and_record( + live.client.revoke_message( chat.clone(), message_id.clone(), whatsapp_rust::send::RevokeType::Sender, - ) - .await - .map_err(|e| e.to_string())?; - live.chat_store - .record_revoke(&chat, &message_id, wacore::time::now_utc()) - .map_err(|e| format!("delete was sent but could not be saved locally: {e}"))?; - live.chat_store - .flush() - .await - .map_err(|e| format!("delete was sent but could not be saved locally: {e}"))?; - Ok(()) + ), + || { + live.chat_store + .record_revoke(&chat, &message_id, wacore::time::now_utc()) + }, + live.chat_store.flush(), + ) + .await } else { // Local delete rides the app-state path, so linked devices // converge on it; the store catches up through the event @@ -242,33 +258,28 @@ impl WhatsAppClient { ), None => return Err("message not found".to_string()), }; - live.client - .chat_actions() - .delete_message_for_me( + delete_and_record( + live.client.chat_actions().delete_message_for_me( &chat, participant.as_ref(), &message_id, from_me, false, Some(timestamp), - ) - .await - .map_err(|e| e.to_string())?; - live.chat_store - .record_delete_for_me( - &chat, - &message_id, - from_me, - participant, - timestamp, - wacore::time::now_utc(), - ) - .map_err(|e| format!("delete was sent but could not be saved locally: {e}"))?; - live.chat_store - .flush() - .await - .map_err(|e| format!("delete was sent but could not be saved locally: {e}"))?; - Ok(()) + ), + || { + live.chat_store.record_delete_for_me( + &chat, + &message_id, + from_me, + participant.clone(), + timestamp, + wacore::time::now_utc(), + ) + }, + live.chat_store.flush(), + ) + .await } }) } diff --git a/crates/session/src/whatsapp/ops.rs b/crates/session/src/whatsapp/ops.rs index b677b55b..9655bb23 100644 --- a/crates/session/src/whatsapp/ops.rs +++ b/crates/session/src/whatsapp/ops.rs @@ -116,6 +116,32 @@ pub struct BackfillReport { } impl WhatsAppClient { + /// Interrupt the library's current reconnect wait. Pause/resume wakes the + /// running supervisor; it does not construct another bot or open a store. + pub fn retry_connection(&self) -> Task> { + let session = self.session.clone(); + self.exec.spawn(async move { + let Some(live) = session.lock().await.clone() else { + return Err("no session yet".to_string()); + }; + let client = &live.client; + if !client + .enable_auto_reconnect + .load(std::sync::atomic::Ordering::Relaxed) + { + return Err("the WhatsApp session has ended".to_string()); + } + // A click that arrived just after recovery must not drop a healthy + // connection, or the calls running over it. + if client.is_connected() { + return Ok(()); + } + client.pause().await; + client.resume(); + Ok(()) + }) + } + /// Create a poll. Returns the server-assigned message id and timestamp. pub fn create_poll( &self, diff --git a/crates/session/src/whatsapp/tests.rs b/crates/session/src/whatsapp/tests.rs index 59243d38..2a617eda 100644 --- a/crates/session/src/whatsapp/tests.rs +++ b/crates/session/src/whatsapp/tests.rs @@ -22,6 +22,89 @@ use whatsapp_rust::wacore::types::events::{Event, ServerAck}; use whatsapp_rust::wacore_binary::Jid; use whatsapp_rust::waproto::whatsapp as wa; +#[tokio::test] +async fn a_remote_delete_failure_skips_local_writes() { + let result = super::mutations::delete_and_record( + async { Err::<(), _>("synthetic remote failure") }, + || panic!("a refused remote delete must not update the local store"), + async { panic!("a refused remote delete must not flush the local store") }, + ) + .await; + + assert_eq!(result, Err("synthetic remote failure".to_string())); +} + +#[tokio::test] +async fn a_sent_delete_succeeds_when_the_local_writer_has_stopped() { + let (chat_store, _) = test_session("sent-delete-stopped-writer").await; + chat_store.close().await.expect("stop local writer"); + let chat: Jid = TEST_PEER.parse().expect("synthetic chat"); + + let result = super::mutations::delete_and_record( + async { Ok::<_, String>(()) }, + || chat_store.record_revoke(&chat, "SENT-DELETE", wacore::time::now_utc()), + async { panic!("an unqueued delete must not await a flush") }, + ) + .await; + + assert_eq!(result, Ok(())); +} + +#[tokio::test] +async fn a_sent_delete_succeeds_when_the_local_batch_fails() { + use std::cell::Cell; + + let recorded = Cell::new(false); + let result = super::mutations::delete_and_record( + async { Ok::<_, String>(()) }, + || { + recorded.set(true); + Ok(()) + }, + async { + assert!(recorded.get(), "flush must follow the queued local update"); + Err(oxidezap_chat_store::ChatStoreError::WriteBatchFailed( + "synthetic rollback".to_string(), + )) + }, + ) + .await; + + assert_eq!(result, Ok(())); +} + +#[tokio::test] +async fn a_delete_waits_for_the_remote_answer_then_the_local_flush() { + use std::cell::Cell; + + let remotely_deleted = Cell::new(false); + let flushed = Cell::new(false); + let result = super::mutations::delete_and_record( + async { + remotely_deleted.set(true); + Ok::<_, String>(()) + }, + || { + assert!( + remotely_deleted.get(), + "local writes must follow remote success" + ); + Ok(()) + }, + async { + flushed.set(true); + Ok(()) + }, + ) + .await; + + assert_eq!(result, Ok(())); + assert!( + flushed.get(), + "successful local persistence must finish before returning" + ); +} + #[test] fn live_nonrenderable_fallback_uses_the_shared_message_kind() { let poll = wa::Message { @@ -107,6 +190,23 @@ fn the_live_data_lane_subscribes_to_server_acks() { ); } +#[test] +fn transport_disconnections_reach_the_control_feed() { + use whatsapp_rust::wacore::types::events::{Disconnected, EventHandler}; + let (handler, receiver, _) = super::interested_channel(super::CONTROL_EVENT_KINDS, 8); + handler.handle_event(Arc::new(Event::Disconnected( + Disconnected::builder() + .reason(wacore::net::DisconnectReason::ReadError( + "synthetic network failure".into(), + )) + .build(), + ))); + assert!(matches!( + &*receiver.try_recv().unwrap(), + Event::Disconnected(_) + )); +} + #[test] fn the_identity_lane_subscribes_to_contact_updates() { assert!( @@ -195,6 +295,7 @@ async fn a_server_ack_publishes_a_live_sent_receipt() { None, None, None, + 0, ) .await; @@ -236,6 +337,7 @@ async fn a_nack_or_non_message_ack_does_not_publish_sent() { None, None, None, + 0, ) .await; @@ -272,6 +374,7 @@ async fn a_chatless_ack_does_not_guess_a_live_destination() { None, None, None, + 0, ) .await; @@ -2838,6 +2941,7 @@ async fn offline_replays_keep_phone_read_flags_without_alerting() { None, None, None, + 0, ) .await; for expected_read in [true, false] { diff --git a/docs/gotchas.md b/docs/gotchas.md index 235ccff3..9c19c9ce 100644 --- a/docs/gotchas.md +++ b/docs/gotchas.md @@ -345,6 +345,22 @@ Non-obvious behaviour, and the reasoning behind it. Read the entry before changi lists rather than a `Result`: picking four photos and one film sends the four and says what happened to the fifth. +- **A WhatsApp outage is not an IPC outage.** Subscribe to `Disconnected` on + the session control feed; the library keeps its automatic backoff and its + single supervisor/store writer. The offline Retry action uses the existing + client's `pause`/`resume` lifecycle to interrupt that backoff, not a second + front-end attachment or a rebuilt session. A rejected retry stays offline + and reports its reason; only losing IPC starts attachment recovery. +- **A recovered offer is not necessarily a new call.** Call work carries the + connection epoch from intake and rechecks it after caller lookup, with offer + publication ordered against disconnect under the registry lock. Offers + explicitly marked offline never ring. An unmarked offer can still be a + backlog item first delivered after reconnect, so its original stanza time + is checked both before registration and before publication: at most two + minutes old, deliberately wider than normal 45-second ringing to allow + delivery delay and modest clock skew. Future timestamps are accepted; this + is a conservative freshness policy, not a claimed server TTL. Offline + snapshots also retire their call cards rather than replaying them. - **An ending is claimed, not owned.** Two places want to publish `UiEvent::CallEnded` — the arm handling the peer's ``, and the watcher parked on `wait_ended` that the resulting hangup resolves — and in a