Skip to content

Commit 619483a

Browse files
authored
perf: split notification handler, boxed-slice children, interest-based props (#553)
1 parent 0c506d5 commit 619483a

6 files changed

Lines changed: 327 additions & 397 deletions

File tree

‎src/client.rs‎

Lines changed: 5 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1783,17 +1783,19 @@ impl Client {
17831783
if response.delta_update {
17841784
debug!(
17851785
"Props delta update received ({} changed props)",
1786-
response.props.len()
1786+
response.experiment_props.len()
17871787
);
17881788
} else {
17891789
debug!(
17901790
"Props full update received ({} props, hash={:?})",
1791-
response.props.len(),
1791+
response.experiment_props.len(),
17921792
response.hash
17931793
);
17941794
}
17951795

1796-
self.ab_props.apply_response(&response).await;
1796+
self.ab_props
1797+
.apply_props(response.delta_update, response.experiment_props.into_iter())
1798+
.await;
17971799

17981800
if let Some(new_hash) = response.hash {
17991801
self.persistence_manager

‎src/handlers/notification.rs‎

Lines changed: 142 additions & 190 deletions
Original file line numberDiff line numberDiff line change
@@ -46,209 +46,36 @@ impl StanzaHandler for NotificationHandler {
4646
}
4747
}
4848

49+
/// Dispatch notification by type. Each arm calls a separate async fn so the
50+
/// compiler doesn't size this future for all arms simultaneously.
4951
async fn handle_notification_impl(client: &Arc<Client>, node: Arc<OwnedNodeRef>) {
5052
let nr = node.get();
5153
let notification_type = nr.attrs().optional_string("type");
52-
let notification_type = notification_type.as_deref().unwrap_or_default();
53-
54-
match notification_type {
55-
"encrypt" => {
56-
// Identity change: <notification type="encrypt" from="user@s.whatsapp.net">
57-
// <identity/>
58-
// </notification>
59-
// WA Web: WAWebHandleIdentityChange — clears device record, deletes sessions,
60-
// marks sender keys for rotation, re-establishes session.
61-
if nr.get_optional_child("identity").is_some() {
62-
handle_identity_change(client, nr).await;
63-
} else if nr
64-
.get_attr("from")
65-
.is_some_and(|v| v.as_str() == wacore_binary::SERVER_JID)
66-
{
67-
// Server-originated encrypt notifications:
68-
// "count" → handlePreKeyLow, "digest" → handleDigestKey
69-
let first_child_tag = nr
70-
.children()
71-
.and_then(|c| c.first().map(|n| n.tag.as_ref()));
72-
73-
match first_child_tag {
74-
Some("count") => {
75-
handle_prekey_low(client).await;
76-
}
77-
Some("digest") => {
78-
handle_digest_key(client);
79-
}
80-
other => {
81-
warn!("Unhandled encrypt notification child: {:?}", other);
82-
}
83-
}
84-
}
85-
}
86-
"server_sync" => {
87-
// Server sync notifications inform us of app state changes from other devices.
88-
// Matches WhatsApp Web's handleServerSyncNotification which calls
89-
// markCollectionsForSync() with the parsed collection names.
90-
use std::str::FromStr;
91-
use wacore::appstate::patch_decode::WAPatchName;
92-
93-
let mut collections = Vec::new();
94-
if let Some(children) = nr.children() {
95-
for collection_node in children.iter().filter(|c| c.tag == "collection") {
96-
let name_cow = collection_node.attrs().optional_string("name");
97-
let name_str = name_cow.as_deref().unwrap_or("<unknown>");
98-
let server_version =
99-
collection_node.attrs().optional_u64("version").unwrap_or(0);
100-
debug!(
101-
target: "Client/AppState",
102-
"Received server_sync for collection '{}' version {}",
103-
name_str, server_version
104-
);
105-
if let Ok(patch_name) = WAPatchName::from_str(name_str)
106-
&& !matches!(patch_name, WAPatchName::Unknown)
107-
{
108-
collections.push((patch_name, server_version));
109-
}
110-
}
111-
}
112-
113-
if !collections.is_empty() {
114-
let client_clone = client.clone();
115-
let generation = client
116-
.connection_generation
117-
.load(std::sync::atomic::Ordering::Acquire);
118-
client.runtime.spawn(Box::pin(async move {
119-
// Check if connection was replaced before starting sync
120-
if client_clone
121-
.connection_generation
122-
.load(std::sync::atomic::Ordering::Acquire)
123-
!= generation
124-
{
125-
log::debug!(target: "Client/AppState", "server_sync task cancelled: connection generation changed");
126-
return;
127-
}
128-
129-
// Filter by version comparison before syncing.
130-
// Matches WA Web's markCollectionsForSync version comparison filter.
131-
let backend = client_clone.persistence_manager.backend();
132-
let mut to_sync = Vec::new();
133-
for (name, server_version) in collections {
134-
if server_version > 0 {
135-
match backend.get_version(name.as_str()).await {
136-
Ok(state) if state.version >= server_version => {
137-
debug!(
138-
target: "Client/AppState",
139-
"Skipping server_sync for {:?}: local version {} >= server version {}",
140-
name, state.version, server_version
141-
);
142-
continue;
143-
}
144-
Ok(_) => {}
145-
Err(e) => {
146-
warn!(
147-
target: "Client/AppState",
148-
"Failed to get local version for {:?}: {e}, syncing anyway", name
149-
);
150-
}
151-
}
152-
}
153-
to_sync.push(name);
154-
}
15554

156-
if !to_sync.is_empty() {
157-
if client_clone.is_shutting_down() {
158-
log::debug!(target: "Client/AppState", "Skipping server_sync: client is shutting down");
159-
return;
160-
}
161-
// Re-check generation after version filtering to avoid syncing
162-
// against a stale connection after the awaited work above.
163-
if client_clone
164-
.connection_generation
165-
.load(std::sync::atomic::Ordering::Acquire)
166-
!= generation
167-
{
168-
log::debug!(target: "Client/AppState", "server_sync task cancelled: connection generation changed during version check");
169-
return;
170-
}
171-
if let Err(e) = client_clone.sync_collections_batched(to_sync).await
172-
&& !client_clone.is_shutting_down()
173-
{
174-
warn!(
175-
target: "Client/AppState",
176-
"Failed to batch sync app state from server_sync: {e}"
177-
);
178-
}
179-
}
180-
})).detach();
181-
}
182-
}
183-
"account_sync" => {
184-
// Handle push name updates
185-
if let Some(new_push_name) = nr.attrs().optional_string("pushname") {
186-
client
187-
.clone()
188-
.update_push_name_and_notify(new_push_name.to_string())
189-
.await;
190-
}
191-
192-
// Handle device list updates (when a new device is paired)
193-
// Matches WhatsApp Web's handleAccountSyncNotification for DEVICES type
194-
if let Some(devices_node) = nr.get_optional_child_by_tag(&["devices"]) {
195-
handle_account_sync_devices(client, nr, devices_node).await;
196-
}
197-
}
198-
"devices" => {
199-
// Handle device list change notifications (WhatsApp Web: handleDevicesNotification)
200-
// These are sent when a user adds, removes, or updates a device
201-
handle_devices_notification(client, nr).await;
202-
}
55+
match notification_type.as_deref().unwrap_or_default() {
56+
"encrypt" => handle_encrypt_notification(client, nr).await,
57+
"server_sync" => handle_server_sync_notification(client, nr),
58+
"account_sync" => handle_account_sync_notification(client, nr).await,
59+
"devices" => handle_devices_notification(client, nr).await,
20360
"link_code_companion_reg" => {
204-
// Handle pair code notification (stage 2 of pair code authentication)
205-
// This is sent when the user enters the code on their phone
20661
crate::pair_code::handle_pair_code_notification(client, nr).await;
20762
}
208-
"business" => {
209-
// Handle business notification (WhatsApp Web: handleBusinessNotification)
210-
// Notifies about business account status changes: verified name, profile, removal
211-
handle_business_notification(client, nr).await;
212-
}
213-
"picture" => {
214-
// Handle profile picture change notifications (WhatsApp Web: WAWebHandleProfilePicNotification)
215-
handle_picture_notification(client, nr);
216-
}
217-
"privacy_token" => {
218-
// Handle incoming trusted contact privacy token notifications.
219-
// Matches WhatsApp Web's WAWebHandlePrivacyTokenNotification.
220-
handle_privacy_token_notification(client, nr).await;
221-
}
222-
"status" => {
223-
// Handle status/about text change notifications (WhatsApp Web: WAWebHandleAboutNotification)
224-
handle_status_notification(client, nr);
225-
}
226-
"contacts" => {
227-
handle_contacts_notification(client, nr).await;
228-
}
229-
"w:gp2" => {
230-
handle_group_notification(client, Arc::clone(&node)).await;
231-
}
232-
"disappearing_mode" => {
233-
// WA Web: WAWebHandleDisappearingModeNotification →
234-
// WAWebUpdateDisappearingModeForContact.
235-
// Parses <disappearing_mode duration="..." t="..."/> child,
236-
// updates the contact's default ephemeral setting.
237-
handle_disappearing_mode_notification(client, nr);
238-
}
239-
"newsletter" => {
240-
handle_newsletter_notification(client, Arc::clone(&node));
241-
}
63+
"business" => handle_business_notification(client, nr).await,
64+
"picture" => handle_picture_notification(client, nr),
65+
"privacy_token" => handle_privacy_token_notification(client, nr).await,
66+
"status" => handle_status_notification(client, nr),
67+
"contacts" => handle_contacts_notification(client, nr).await,
68+
"w:gp2" => handle_group_notification(client, Arc::clone(&node)).await,
69+
"disappearing_mode" => handle_disappearing_mode_notification(client, nr),
70+
"newsletter" => handle_newsletter_notification(client, Arc::clone(&node)),
24271
"mediaretry" => {
243-
// Handled by wait_for_node waiter in MediaReupload::request().
244-
// Ack is sent automatically by the stanza dispatch loop.
24572
debug!(
24673
"Received mediaretry notification for msg {}",
24774
nr.attrs().optional_string("id").unwrap_or_default()
24875
);
24976
}
250-
_ => {
251-
debug!("Unhandled notification type '{notification_type}', dispatching raw event");
77+
other => {
78+
debug!("Unhandled notification type '{other}', dispatching raw event");
25279
client
25380
.core
25481
.event_bus
@@ -257,6 +84,131 @@ async fn handle_notification_impl(client: &Arc<Client>, node: Arc<OwnedNodeRef>)
25784
}
25885
}
25986

87+
async fn handle_encrypt_notification(client: &Arc<Client>, nr: &wacore_binary::NodeRef<'_>) {
88+
if nr.get_optional_child("identity").is_some() {
89+
handle_identity_change(client, nr).await;
90+
} else if nr
91+
.get_attr("from")
92+
.is_some_and(|v| v.as_str() == wacore_binary::SERVER_JID)
93+
{
94+
let first_child_tag = nr
95+
.children()
96+
.and_then(|c| c.first().map(|n| n.tag.as_ref()));
97+
match first_child_tag {
98+
Some("count") => handle_prekey_low(client).await,
99+
Some("digest") => handle_digest_key(client),
100+
other => warn!("Unhandled encrypt notification child: {:?}", other),
101+
}
102+
}
103+
}
104+
105+
/// Sync is fire-and-forget (spawned), so this is not async -- it parses
106+
/// collection nodes synchronously and spawns the async sync task.
107+
fn handle_server_sync_notification(client: &Arc<Client>, nr: &wacore_binary::NodeRef<'_>) {
108+
use std::str::FromStr;
109+
use wacore::appstate::patch_decode::WAPatchName;
110+
111+
let mut collections = Vec::new();
112+
if let Some(children) = nr.children() {
113+
for collection_node in children.iter().filter(|c| c.tag == "collection") {
114+
let name_cow = collection_node.attrs().optional_string("name");
115+
let name_str = name_cow.as_deref().unwrap_or("<unknown>");
116+
let server_version = collection_node.attrs().optional_u64("version").unwrap_or(0);
117+
debug!(
118+
target: "Client/AppState",
119+
"Received server_sync for collection '{}' version {}",
120+
name_str, server_version
121+
);
122+
if let Ok(patch_name) = WAPatchName::from_str(name_str)
123+
&& !matches!(patch_name, WAPatchName::Unknown)
124+
{
125+
collections.push((patch_name, server_version));
126+
}
127+
}
128+
}
129+
130+
if !collections.is_empty() {
131+
let client_clone = client.clone();
132+
let generation = client
133+
.connection_generation
134+
.load(std::sync::atomic::Ordering::Acquire);
135+
client
136+
.runtime
137+
.spawn(Box::pin(async move {
138+
if client_clone
139+
.connection_generation
140+
.load(std::sync::atomic::Ordering::Acquire)
141+
!= generation
142+
{
143+
log::debug!(target: "Client/AppState", "server_sync task cancelled: connection generation changed");
144+
return;
145+
}
146+
147+
let backend = client_clone.persistence_manager.backend();
148+
let mut to_sync = Vec::new();
149+
for (name, server_version) in collections {
150+
if server_version > 0 {
151+
match backend.get_version(name.as_str()).await {
152+
Ok(state) if state.version >= server_version => {
153+
debug!(
154+
target: "Client/AppState",
155+
"Skipping server_sync for {:?}: local version {} >= server version {}",
156+
name, state.version, server_version
157+
);
158+
continue;
159+
}
160+
Ok(_) => {}
161+
Err(e) => {
162+
warn!(
163+
target: "Client/AppState",
164+
"Failed to get local version for {:?}: {e}, syncing anyway",
165+
name
166+
);
167+
}
168+
}
169+
}
170+
to_sync.push(name);
171+
}
172+
173+
if !to_sync.is_empty() {
174+
if client_clone.is_shutting_down() {
175+
log::debug!(target: "Client/AppState", "Skipping server_sync: client is shutting down");
176+
return;
177+
}
178+
if client_clone
179+
.connection_generation
180+
.load(std::sync::atomic::Ordering::Acquire)
181+
!= generation
182+
{
183+
log::debug!(target: "Client/AppState", "server_sync task cancelled: connection generation changed during version check");
184+
return;
185+
}
186+
if let Err(e) = client_clone.sync_collections_batched(to_sync).await
187+
&& !client_clone.is_shutting_down()
188+
{
189+
warn!(
190+
target: "Client/AppState",
191+
"Failed to batch sync app state from server_sync: {e}"
192+
);
193+
}
194+
}
195+
}))
196+
.detach();
197+
}
198+
}
199+
200+
async fn handle_account_sync_notification(client: &Arc<Client>, nr: &wacore_binary::NodeRef<'_>) {
201+
if let Some(new_push_name) = nr.attrs().optional_string("pushname") {
202+
client
203+
.clone()
204+
.update_push_name_and_notify(new_push_name.to_string())
205+
.await;
206+
}
207+
if let Some(devices_node) = nr.get_optional_child_by_tag(&["devices"]) {
208+
handle_account_sync_devices(client, nr, devices_node).await;
209+
}
210+
}
211+
260212
/// Handle encrypt/count notification (PreKey Low).
261213
///
262214
/// Matches WA Web's `WAWebHandlePreKeyLow`:

0 commit comments

Comments
 (0)