From 075995b6087e978452afd56eff6ca51e0199b146 Mon Sep 17 00:00:00 2001 From: Cody Henthorne Date: Fri, 25 Sep 2026 11:47:14 -0400 Subject: [PATCH] Improve relability of large group sends on slower devices. --- .../api/SignalServiceMessageSender.java | 41 +++++++++++-------- 1 file changed, 24 insertions(+), 17 deletions(-) diff --git a/lib/libsignal-service/src/main/java/org/whispersystems/signalservice/api/SignalServiceMessageSender.java b/lib/libsignal-service/src/main/java/org/whispersystems/signalservice/api/SignalServiceMessageSender.java index f3bf588f06..2882d9fb11 100644 --- a/lib/libsignal-service/src/main/java/org/whispersystems/signalservice/api/SignalServiceMessageSender.java +++ b/lib/libsignal-service/src/main/java/org/whispersystems/signalservice/api/SignalServiceMessageSender.java @@ -172,7 +172,8 @@ public class SignalServiceMessageSender { private static final String TAG = SignalServiceMessageSender.class.getSimpleName().substring(0, 23); - private static final int RETRY_COUNT = 4; + private static final int RETRY_COUNT = 4; + private static final int MAX_CONCURRENT_SENDS = 32; private final PushServiceSocket socket; private final SignalServiceAccountDataStore aciStore; @@ -2195,22 +2196,25 @@ public class SignalServiceMessageSender { return kotlin.Unit.INSTANCE; }); - List> singleResults = new LinkedList<>(); - Iterator recipientIterator = recipients.iterator(); - Iterator sealedSenderAccessIterator = sealedSenderAccesses.iterator(); + Observable> singleResults = Observable.create(emitter -> { + Iterator recipientIterator = recipients.iterator(); + Iterator sealedSenderAccessIterator = sealedSenderAccesses.iterator(); - while (recipientIterator.hasNext()) { - SignalServiceAddress recipient = recipientIterator.next(); - SealedSenderAccess sealedSenderAccess = sealedSenderAccessIterator.next(); + while (recipientIterator.hasNext() && !emitter.isDisposed()) { + SignalServiceAddress recipient = recipientIterator.next(); + SealedSenderAccess sealedSenderAccess = sealedSenderAccessIterator.next(); - singleResults.add(sendMessageRx(recipient, sealedSenderAccess, timestamp, content, online, cancelationSignal, sendEvents, urgent, story, 0).toObservable()); - } + emitter.onNext(sendMessageRx(recipient, sealedSenderAccess, timestamp, content, online, cancelationSignal, sendEvents, urgent, story, 0).toObservable()); + } + + emitter.onComplete(); + }); List results; try { - results = Observable.mergeDelayError(singleResults, Integer.MAX_VALUE, 1) + results = Observable.mergeDelayError(singleResults, MAX_CONCURRENT_SENDS) .observeOn(scheduler, true) - .scan(new ArrayList(singleResults.size()), (state, result) -> { + .scan(new ArrayList(recipients.size()), (state, result) -> { state.add(result); if (partialListener != null) { partialListener.onPartialSendComplete(result); @@ -2260,9 +2264,9 @@ public class SignalServiceMessageSender { * Sends a message over the appropriate websocket, falls back to REST when unavailable, and emits a {@link SendMessageResult} for most business * logic error cases. *

- * Uses a "feature" or Rx where if no {@link Single#subscribeOn(Scheduler)} operator is used, the subscribing thread is used to perform the - * initial work. This allows the calling thread to do the starting of the send work (encryption and putting it on the wire) and can be called - * multiple times in a loop, but allow the network transit/processing/error retry logic to run on a background thread. + * Encryption is performed eagerly on the calling thread when this method is called, before the returned single is subscribed to. This keeps + * encryption on a single thread to avoid session lock contention, while allowing subscribers to limit how many sends are in flight. The + * network transit/processing/error retry logic runs on a background thread. *

* Processing happens on the background thread via an {@link Single#observeOn(Scheduler)} call after the encrypt and send. Error * handling operators are added after the observe so they will also run on a background thread. Retry logic during error handling @@ -2286,7 +2290,8 @@ public class SignalServiceMessageSender { long startTime = System.currentTimeMillis(); enforceMaxEnvelopeContentSize(content); - Single messagesSingle = Single.fromCallable(() -> { + Single messagesSingle; + try { OutgoingPushMessageList messages = getEncryptedMessages(recipient, sealedSenderAccess, timestamp, content, online, urgent, story); if (retryCount == 0 && sendEvents != null) { @@ -2299,8 +2304,10 @@ public class SignalServiceMessageSender { Log.d(TAG, "[sendMessage][" + timestamp + "] Sending a SKDM to " + messages.getDestination() + " for devices: " + messages.getDevices() + (content.getContent().get().dataMessage != null ? " (it's piggy-backing on a DataMessage) via Rx" : " via Rx")); } - return messages; - }); + messagesSingle = Single.just(messages); + } catch (Exception e) { + messagesSingle = Single.error(e); + } Single sendWithFallback = messagesSingle .flatMap(messages -> {