Improve relability of large group sends on slower devices.

This commit is contained in:
Cody Henthorne
2026-09-30 13:34:03 -03:00
committed by Alex Hart
parent 7581b1774e
commit 075995b608
@@ -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<Observable<SendMessageResult>> singleResults = new LinkedList<>();
Iterator<SignalServiceAddress> recipientIterator = recipients.iterator();
Iterator<SealedSenderAccess> sealedSenderAccessIterator = sealedSenderAccesses.iterator();
Observable<Observable<SendMessageResult>> singleResults = Observable.create(emitter -> {
Iterator<SignalServiceAddress> recipientIterator = recipients.iterator();
Iterator<SealedSenderAccess> 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<SendMessageResult> results;
try {
results = Observable.mergeDelayError(singleResults, Integer.MAX_VALUE, 1)
results = Observable.mergeDelayError(singleResults, MAX_CONCURRENT_SENDS)
.observeOn(scheduler, true)
.scan(new ArrayList<SendMessageResult>(singleResults.size()), (state, result) -> {
.scan(new ArrayList<SendMessageResult>(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.
* <p>
* 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.
* <p>
* 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<OutgoingPushMessageList> messagesSingle = Single.fromCallable(() -> {
Single<OutgoingPushMessageList> 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<SendMessageResult> sendWithFallback = messagesSingle
.flatMap(messages -> {