mirror of
https://github.com/signalapp/Signal-Android.git
synced 2026-08-05 21:07:49 +01:00
Improve message backlog processing.
Moves SendDeliveryReceiptJob out of the critical path of message processing. EditMessageProcessor now uses the batch cache.
This commit is contained in:
committed by
Alex Hart
parent
b2b314b18e
commit
f1aefb1cdd
@@ -1224,8 +1224,8 @@ open class MessageTable(context: Context?, databaseHelper: SignalDatabase) : Dat
|
||||
}
|
||||
}
|
||||
|
||||
fun insertEditMessageInbox(mediaMessage: IncomingMessage, targetMessage: MmsMessageRecord): Optional<InsertResult> {
|
||||
val insertResult = insertMessageInbox(retrieved = mediaMessage, editedMessage = targetMessage, notifyObservers = false)
|
||||
fun insertEditMessageInbox(mediaMessage: IncomingMessage, targetMessage: MmsMessageRecord, skipThreadUpdate: Boolean = false): Optional<InsertResult> {
|
||||
val insertResult = insertMessageInbox(retrieved = mediaMessage, editedMessage = targetMessage, notifyObservers = false, skipThreadUpdate = skipThreadUpdate)
|
||||
|
||||
if (insertResult.isPresent) {
|
||||
val (messageId) = insertResult.get()
|
||||
|
||||
@@ -1,185 +0,0 @@
|
||||
package org.thoughtcrime.securesms.jobs;
|
||||
|
||||
|
||||
import androidx.annotation.NonNull;
|
||||
import androidx.annotation.Nullable;
|
||||
|
||||
import org.signal.core.util.logging.Log;
|
||||
import org.signal.libsignal.zkgroup.groupsend.GroupSendFullToken;
|
||||
import org.signal.network.exceptions.PushNetworkException;
|
||||
import org.thoughtcrime.securesms.crypto.SealedSenderAccessUtil;
|
||||
import org.thoughtcrime.securesms.database.RecipientTable.RegisteredState;
|
||||
import org.thoughtcrime.securesms.database.SignalDatabase;
|
||||
import org.thoughtcrime.securesms.database.model.MessageId;
|
||||
import org.thoughtcrime.securesms.database.model.RecipientRecord;
|
||||
import org.thoughtcrime.securesms.dependencies.AppDependencies;
|
||||
import org.thoughtcrime.securesms.jobmanager.Job;
|
||||
import org.thoughtcrime.securesms.jobmanager.JsonJobData;
|
||||
import org.thoughtcrime.securesms.jobmanager.impl.NetworkConstraint;
|
||||
import org.thoughtcrime.securesms.jobmanager.impl.SealedSenderConstraint;
|
||||
import org.thoughtcrime.securesms.net.NotPushRegisteredException;
|
||||
import org.thoughtcrime.securesms.recipients.Recipient;
|
||||
import org.thoughtcrime.securesms.recipients.RecipientId;
|
||||
import org.thoughtcrime.securesms.recipients.RecipientUtil;
|
||||
import org.thoughtcrime.securesms.transport.UndeliverableMessageException;
|
||||
import org.whispersystems.signalservice.api.SignalServiceMessageSender;
|
||||
import org.whispersystems.signalservice.api.crypto.ContentHint;
|
||||
import org.whispersystems.signalservice.api.crypto.UntrustedIdentityException;
|
||||
import org.whispersystems.signalservice.api.messages.SendMessageResult;
|
||||
import org.whispersystems.signalservice.api.messages.SignalServiceReceiptMessage;
|
||||
import org.whispersystems.signalservice.api.push.SignalServiceAddress;
|
||||
import org.whispersystems.signalservice.api.push.exceptions.ServerRejectedException;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.util.Collections;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
|
||||
public class SendDeliveryReceiptJob extends BaseJob {
|
||||
|
||||
public static final String KEY = "SendDeliveryReceiptJob";
|
||||
|
||||
private static final String KEY_RECIPIENT = "recipient";
|
||||
private static final String KEY_MESSAGE_SENT_TIMESTAMP = "message_id";
|
||||
private static final String KEY_TIMESTAMP = "timestamp";
|
||||
private static final String KEY_MESSAGE_ID = "message_db_id";
|
||||
|
||||
private static final String TAG = Log.tag(SendReadReceiptJob.class);
|
||||
|
||||
private final RecipientId recipientId;
|
||||
private final long messageSentTimestamp;
|
||||
private final long timestamp;
|
||||
|
||||
@Nullable
|
||||
private final MessageId messageId;
|
||||
|
||||
public SendDeliveryReceiptJob(@NonNull RecipientId recipientId, long messageSentTimestamp, @NonNull MessageId messageId) {
|
||||
this(new Job.Parameters.Builder()
|
||||
.addConstraint(NetworkConstraint.KEY)
|
||||
.addConstraint(SealedSenderConstraint.KEY)
|
||||
.setLifespan(TimeUnit.DAYS.toMillis(1))
|
||||
.setMaxAttempts(Parameters.UNLIMITED)
|
||||
.setQueue(recipientId.toQueueKey())
|
||||
.build(),
|
||||
recipientId,
|
||||
messageSentTimestamp,
|
||||
messageId,
|
||||
System.currentTimeMillis());
|
||||
}
|
||||
|
||||
private SendDeliveryReceiptJob(@NonNull Job.Parameters parameters,
|
||||
@NonNull RecipientId recipientId,
|
||||
long messageSentTimestamp,
|
||||
@Nullable MessageId messageId,
|
||||
long timestamp)
|
||||
{
|
||||
super(parameters);
|
||||
|
||||
this.recipientId = recipientId;
|
||||
this.messageSentTimestamp = messageSentTimestamp;
|
||||
this.messageId = messageId;
|
||||
this.timestamp = timestamp;
|
||||
}
|
||||
|
||||
@Override
|
||||
public @Nullable byte[] serialize() {
|
||||
JsonJobData.Builder builder = new JsonJobData.Builder().putString(KEY_RECIPIENT, recipientId.serialize())
|
||||
.putLong(KEY_MESSAGE_SENT_TIMESTAMP, messageSentTimestamp)
|
||||
.putLong(KEY_TIMESTAMP, timestamp);
|
||||
|
||||
if (messageId != null) {
|
||||
builder.putString(KEY_MESSAGE_ID, messageId.serialize());
|
||||
}
|
||||
|
||||
return builder.serialize();
|
||||
}
|
||||
|
||||
@Override
|
||||
public @NonNull String getFactoryKey() {
|
||||
return KEY;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onRun() throws IOException, UntrustedIdentityException, UndeliverableMessageException {
|
||||
if (!Recipient.self().isRegistered()) {
|
||||
throw new NotPushRegisteredException();
|
||||
}
|
||||
|
||||
SignalServiceMessageSender messageSender = AppDependencies.getSignalServiceMessageSender();
|
||||
RecipientRecord recipient = SignalDatabase.recipients().getRecord(recipientId);
|
||||
|
||||
if (recipient.getId().equals(Recipient.self().getId())) {
|
||||
Log.i(TAG, "Not sending to self, abort");
|
||||
return;
|
||||
}
|
||||
|
||||
if (recipient.getRegistered() == RegisteredState.NOT_REGISTERED) {
|
||||
Log.w(TAG, recipient.getId() + " is unregistered!");
|
||||
return;
|
||||
}
|
||||
|
||||
if (recipient.getServiceId() == null && recipient.getE164() == null) {
|
||||
Log.w(TAG, "No serviceId or e164!");
|
||||
return;
|
||||
}
|
||||
|
||||
SignalServiceAddress remoteAddress = RecipientUtil.toSignalServiceAddress(recipient);
|
||||
SignalServiceReceiptMessage receiptMessage = new SignalServiceReceiptMessage(SignalServiceReceiptMessage.Type.DELIVERY,
|
||||
Collections.singletonList(messageSentTimestamp),
|
||||
timestamp);
|
||||
|
||||
SendMessageResult result = ReceiptSender.sendWithSessionRepair(recipientId, () -> messageSender.sendReceipt(remoteAddress,
|
||||
SealedSenderAccessUtil.getSealedSenderAccessFor(recipient, this::getGroupSendFullToken),
|
||||
receiptMessage,
|
||||
recipient.needsPniSignature()));
|
||||
|
||||
if (result != null && messageId != null) {
|
||||
SignalDatabase.messageLog().insertIfPossible(recipientId, timestamp, result, ContentHint.IMPLICIT, messageId, false);
|
||||
}
|
||||
}
|
||||
|
||||
private @Nullable GroupSendFullToken getGroupSendFullToken() {
|
||||
if (messageId == null) {
|
||||
return null;
|
||||
}
|
||||
|
||||
long threadId = SignalDatabase.messages().getThreadIdForMessage(messageId.getId());
|
||||
if (threadId == -1) {
|
||||
return null;
|
||||
}
|
||||
|
||||
return SignalDatabase.groups().getGroupSendFullToken(threadId, recipientId);
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean onShouldRetry(@NonNull Exception e) {
|
||||
if (e instanceof ServerRejectedException) return false;
|
||||
if (e instanceof PushNetworkException) return true;
|
||||
if (e instanceof IOException) return true;
|
||||
|
||||
return false;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onFailure() {
|
||||
Log.w(TAG, "Failed to send delivery receipt to: " + recipientId);
|
||||
}
|
||||
|
||||
public static final class Factory implements Job.Factory<SendDeliveryReceiptJob> {
|
||||
@Override
|
||||
public @NonNull SendDeliveryReceiptJob create(@NonNull Parameters parameters, @Nullable byte[] serializedData) {
|
||||
JsonJobData data = JsonJobData.deserialize(serializedData);
|
||||
|
||||
MessageId messageId = null;
|
||||
|
||||
if (data.hasString(KEY_MESSAGE_ID)) {
|
||||
messageId = MessageId.deserialize(data.getString(KEY_MESSAGE_ID));
|
||||
}
|
||||
|
||||
return new SendDeliveryReceiptJob(parameters,
|
||||
RecipientId.from(data.getString(KEY_RECIPIENT)),
|
||||
data.getLong(KEY_MESSAGE_SENT_TIMESTAMP),
|
||||
messageId,
|
||||
data.getLong(KEY_TIMESTAMP));
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,190 @@
|
||||
package org.thoughtcrime.securesms.jobs
|
||||
|
||||
import org.signal.core.util.logging.Log
|
||||
import org.signal.core.util.logging.Log.tag
|
||||
import org.signal.libsignal.zkgroup.groupsend.GroupSendFullToken
|
||||
import org.thoughtcrime.securesms.crypto.SealedSenderAccessUtil
|
||||
import org.thoughtcrime.securesms.database.RecipientTable
|
||||
import org.thoughtcrime.securesms.database.SignalDatabase
|
||||
import org.thoughtcrime.securesms.database.model.MessageId
|
||||
import org.thoughtcrime.securesms.dependencies.AppDependencies
|
||||
import org.thoughtcrime.securesms.jobmanager.Job
|
||||
import org.thoughtcrime.securesms.jobmanager.JsonJobData
|
||||
import org.thoughtcrime.securesms.jobmanager.impl.DecryptionsDrainedConstraint
|
||||
import org.thoughtcrime.securesms.jobmanager.impl.NetworkConstraint
|
||||
import org.thoughtcrime.securesms.jobmanager.impl.SealedSenderConstraint
|
||||
import org.thoughtcrime.securesms.jobs.ReceiptSender.ReceiptSendOperation
|
||||
import org.thoughtcrime.securesms.jobs.ReceiptSender.sendWithSessionRepair
|
||||
import org.thoughtcrime.securesms.recipients.Recipient
|
||||
import org.thoughtcrime.securesms.recipients.RecipientId
|
||||
import org.thoughtcrime.securesms.recipients.RecipientUtil
|
||||
import org.thoughtcrime.securesms.transport.UndeliverableMessageException
|
||||
import org.whispersystems.signalservice.api.crypto.ContentHint
|
||||
import org.whispersystems.signalservice.api.crypto.UntrustedIdentityException
|
||||
import org.whispersystems.signalservice.api.messages.SignalServiceReceiptMessage
|
||||
import org.whispersystems.signalservice.api.push.exceptions.ServerRejectedException
|
||||
import java.io.IOException
|
||||
import java.util.concurrent.TimeUnit
|
||||
|
||||
class SendDeliveryReceiptJob private constructor(
|
||||
parameters: Parameters,
|
||||
private val recipientId: RecipientId,
|
||||
private val messageSentTimestamps: List<Long>,
|
||||
private val messageIds: List<MessageId>,
|
||||
private val receiptSentTimestamp: Long
|
||||
) : Job(parameters) {
|
||||
|
||||
companion object {
|
||||
const val KEY: String = "SendDeliveryReceiptJob"
|
||||
|
||||
private const val KEY_RECIPIENT = "recipient"
|
||||
private const val KEY_RECEIPT_TIMESTAMP = "timestamp"
|
||||
private const val KEY_MESSAGE_SENT_TIMESTAMPS = "message_sent_timestamps"
|
||||
private const val KEY_MESSAGE_IDS = "message_db_ids"
|
||||
|
||||
private const val KEY_LEGACY_MESSAGE_SENT_TIMESTAMP = "message_id"
|
||||
private const val KEY_LEGACY_MESSAGE_ID = "message_db_id"
|
||||
|
||||
private val TAG = tag(SendDeliveryReceiptJob::class.java)
|
||||
|
||||
fun create(
|
||||
recipientId: RecipientId,
|
||||
messageSentTimestamps: List<Long>,
|
||||
messageIds: List<MessageId>
|
||||
): List<SendDeliveryReceiptJob> {
|
||||
return messageSentTimestamps.zip(messageIds)
|
||||
.chunked(SendReadReceiptJob.MAX_TIMESTAMPS)
|
||||
.map { chunk -> SendDeliveryReceiptJob(recipientId, chunk.map { it.first }, chunk.map { it.second }) }
|
||||
}
|
||||
}
|
||||
|
||||
private constructor(recipientId: RecipientId, messageSentTimestamps: List<Long>, messageIds: List<MessageId>) : this(
|
||||
parameters = Parameters.Builder()
|
||||
.addConstraint(NetworkConstraint.KEY)
|
||||
.addConstraint(SealedSenderConstraint.KEY)
|
||||
.addConstraint(DecryptionsDrainedConstraint.KEY)
|
||||
.setLifespan(TimeUnit.DAYS.toMillis(1))
|
||||
.setMaxAttempts(Parameters.UNLIMITED)
|
||||
.setQueue(recipientId.toReceiptQueueKey())
|
||||
.build(),
|
||||
recipientId = recipientId,
|
||||
messageSentTimestamps = messageSentTimestamps,
|
||||
messageIds = messageIds,
|
||||
receiptSentTimestamp = System.currentTimeMillis()
|
||||
)
|
||||
|
||||
override fun serialize(): ByteArray? {
|
||||
return JsonJobData.Builder().putString(KEY_RECIPIENT, recipientId.serialize())
|
||||
.putLongArray(KEY_MESSAGE_SENT_TIMESTAMPS, messageSentTimestamps.toLongArray())
|
||||
.putStringListAsArray(KEY_MESSAGE_IDS, messageIds.map { it.serialize() })
|
||||
.putLong(KEY_RECEIPT_TIMESTAMP, receiptSentTimestamp)
|
||||
.serialize()
|
||||
}
|
||||
|
||||
override fun getFactoryKey(): String {
|
||||
return KEY
|
||||
}
|
||||
|
||||
override fun run(): Result {
|
||||
if (!Recipient.self().isRegistered) {
|
||||
return Result.failure()
|
||||
}
|
||||
|
||||
if (messageSentTimestamps.isEmpty()) {
|
||||
return Result.success()
|
||||
}
|
||||
|
||||
val messageSender = AppDependencies.signalServiceMessageSender
|
||||
val recipient = SignalDatabase.recipients.getRecord(recipientId)
|
||||
|
||||
if (recipient.id == Recipient.self().id) {
|
||||
Log.i(TAG, "Not sending to self, abort")
|
||||
return Result.success()
|
||||
}
|
||||
|
||||
if (recipient.registered == RecipientTable.RegisteredState.NOT_REGISTERED) {
|
||||
Log.w(TAG, "${recipient.id} is unregistered!")
|
||||
return Result.failure()
|
||||
}
|
||||
|
||||
if (recipient.serviceId == null && recipient.e164 == null) {
|
||||
Log.w(TAG, "No serviceId or e164!")
|
||||
return Result.failure()
|
||||
}
|
||||
|
||||
val result = try {
|
||||
val remoteAddress = RecipientUtil.toSignalServiceAddress(recipient)
|
||||
val receiptMessage = SignalServiceReceiptMessage(
|
||||
SignalServiceReceiptMessage.Type.DELIVERY,
|
||||
messageSentTimestamps,
|
||||
receiptSentTimestamp
|
||||
)
|
||||
|
||||
sendWithSessionRepair(
|
||||
recipientId,
|
||||
ReceiptSendOperation {
|
||||
messageSender.sendReceipt(
|
||||
remoteAddress,
|
||||
SealedSenderAccessUtil.getSealedSenderAccessFor(recipient) { this.getGroupSendFullToken() },
|
||||
receiptMessage,
|
||||
recipient.needsPniSignature
|
||||
)
|
||||
}
|
||||
)
|
||||
} catch (e: ServerRejectedException) {
|
||||
Log.i(TAG, "Send failed", e)
|
||||
return Result.failure()
|
||||
} catch (e: IOException) {
|
||||
Log.d(TAG, "Send failed, retrying", e)
|
||||
return Result.retry(defaultBackoff())
|
||||
} catch (e: UntrustedIdentityException) {
|
||||
Log.i(TAG, "Send failed", e)
|
||||
return Result.failure()
|
||||
} catch (e: UndeliverableMessageException) {
|
||||
Log.i(TAG, "Send failed", e)
|
||||
return Result.failure()
|
||||
}
|
||||
|
||||
if (result != null && messageIds.isNotEmpty()) {
|
||||
SignalDatabase.messageLog.insertIfPossible(recipientId, receiptSentTimestamp, result, ContentHint.IMPLICIT, messageIds, false)
|
||||
}
|
||||
|
||||
return Result.success()
|
||||
}
|
||||
|
||||
private fun getGroupSendFullToken(): GroupSendFullToken? {
|
||||
for (messageId in messageIds) {
|
||||
val threadId = SignalDatabase.messages.getThreadIdForMessage(messageId.id)
|
||||
if (threadId != -1L) {
|
||||
return SignalDatabase.groups.getGroupSendFullToken(threadId, recipientId)
|
||||
}
|
||||
}
|
||||
|
||||
return null
|
||||
}
|
||||
|
||||
override fun onFailure() {
|
||||
Log.w(TAG, "Failed to send delivery receipt to: $recipientId")
|
||||
}
|
||||
|
||||
class Factory : Job.Factory<SendDeliveryReceiptJob> {
|
||||
override fun create(parameters: Parameters, serializedData: ByteArray?): SendDeliveryReceiptJob {
|
||||
val data = JsonJobData.deserialize(serializedData)
|
||||
|
||||
val recipientId = RecipientId.from(data.getString(KEY_RECIPIENT))
|
||||
val timestamp = data.getLong(KEY_RECEIPT_TIMESTAMP)
|
||||
val sentTimestamps: List<Long>
|
||||
val messageIds: List<MessageId>
|
||||
|
||||
if (data.hasLongArray(KEY_MESSAGE_SENT_TIMESTAMPS)) {
|
||||
sentTimestamps = data.getLongArray(KEY_MESSAGE_SENT_TIMESTAMPS).toList()
|
||||
messageIds = if (data.hasStringArray(KEY_MESSAGE_IDS)) data.getStringArrayAsList(KEY_MESSAGE_IDS).map { MessageId.deserialize(it) } else emptyList()
|
||||
} else {
|
||||
sentTimestamps = listOf(data.getLong(KEY_LEGACY_MESSAGE_SENT_TIMESTAMP))
|
||||
messageIds = if (data.hasString(KEY_LEGACY_MESSAGE_ID)) listOf(MessageId.deserialize(data.getString(KEY_LEGACY_MESSAGE_ID))) else emptyList()
|
||||
}
|
||||
|
||||
return SendDeliveryReceiptJob(parameters, recipientId, sentTimestamps, messageIds, timestamp)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -10,9 +10,11 @@ import org.signal.libsignal.zkgroup.groups.GroupSecretParams
|
||||
import org.thoughtcrime.securesms.database.MessageTable
|
||||
import org.thoughtcrime.securesms.database.SignalDatabase
|
||||
import org.thoughtcrime.securesms.database.model.GroupRecord
|
||||
import org.thoughtcrime.securesms.database.model.MessageId
|
||||
import org.thoughtcrime.securesms.dependencies.AppDependencies
|
||||
import org.thoughtcrime.securesms.groups.GroupId
|
||||
import org.thoughtcrime.securesms.jobmanager.Job
|
||||
import org.thoughtcrime.securesms.jobs.SendDeliveryReceiptJob
|
||||
import org.thoughtcrime.securesms.messages.SignalServiceProtoUtil.groupMasterKey
|
||||
import org.thoughtcrime.securesms.messages.SignalServiceProtoUtil.hasGroupContext
|
||||
import org.thoughtcrime.securesms.recipients.RecipientId
|
||||
@@ -76,9 +78,14 @@ abstract class BatchCache {
|
||||
SignalDatabase.messageLog.deleteEntriesForRecipient(timestamps, recipientId, device)
|
||||
}
|
||||
|
||||
protected fun flushDeliveryReceipt(recipientId: RecipientId, timestamp: Long, messageId: MessageId) {
|
||||
SendDeliveryReceiptJob.create(recipientId, listOf(timestamp), listOf(messageId)).forEach { flushJob(it) }
|
||||
}
|
||||
|
||||
abstract fun addJob(job: Job)
|
||||
abstract fun addIncomingMessageInsertThreadUpdate(threadId: Long)
|
||||
abstract fun addMslDelete(recipientId: RecipientId, device: Int, timestamps: List<Long>)
|
||||
abstract fun addDeliveryReceipt(recipientId: RecipientId, groupId: GroupId.V2?, timestamp: Long, messageId: MessageId)
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -99,6 +106,10 @@ class OneTimeBatchCache : BatchCache() {
|
||||
override fun addMslDelete(recipientId: RecipientId, device: Int, timestamps: List<Long>) {
|
||||
flushMslDelete(recipientId, device, timestamps)
|
||||
}
|
||||
|
||||
override fun addDeliveryReceipt(recipientId: RecipientId, groupId: GroupId.V2?, timestamp: Long, messageId: MessageId) {
|
||||
flushDeliveryReceipt(recipientId, timestamp, messageId)
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -117,6 +128,7 @@ class ReusedBatchCache : BatchCache() {
|
||||
private val batchedJobs = ArrayList<Job>(BATCH_SIZE)
|
||||
private val threadUpdates = HashSet<Long>(BATCH_SIZE)
|
||||
private val mslDeletes = HashMap<Pair<RecipientId, Int>, MutableList<Long>>(BATCH_SIZE)
|
||||
private val deliveryReceipts = HashMap<Pair<RecipientId, GroupId.V2?>, DeliveryReceiptAccumulator>(BATCH_SIZE)
|
||||
|
||||
override fun addJob(job: Job) {
|
||||
batchedJobs += job
|
||||
@@ -130,9 +142,20 @@ class ReusedBatchCache : BatchCache() {
|
||||
mslDeletes.getOrPut(recipientId to device) { mutableListOf() } += timestamps
|
||||
}
|
||||
|
||||
override fun addDeliveryReceipt(recipientId: RecipientId, groupId: GroupId.V2?, timestamp: Long, messageId: MessageId) {
|
||||
val accumulator = deliveryReceipts.getOrPut(recipientId to groupId) { DeliveryReceiptAccumulator() }
|
||||
accumulator.timestamps += timestamp
|
||||
accumulator.messageIds += messageId
|
||||
}
|
||||
|
||||
override fun flushAndClear() {
|
||||
super.flushAndClear()
|
||||
|
||||
deliveryReceipts.flatMapTo(batchedJobs) { (key, accumulator) ->
|
||||
SendDeliveryReceiptJob.create(key.first, accumulator.timestamps, accumulator.messageIds)
|
||||
}
|
||||
deliveryReceipts.clear()
|
||||
|
||||
if (batchedJobs.isNotEmpty()) {
|
||||
AppDependencies.jobManager.addAll(batchedJobs)
|
||||
}
|
||||
@@ -147,4 +170,9 @@ class ReusedBatchCache : BatchCache() {
|
||||
threadUpdates.clear()
|
||||
mslDeletes.clear()
|
||||
}
|
||||
|
||||
private class DeliveryReceiptAccumulator {
|
||||
val timestamps = ArrayList<Long>(BATCH_SIZE)
|
||||
val messageIds = ArrayList<MessageId>(BATCH_SIZE)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -63,7 +63,6 @@ import org.thoughtcrime.securesms.jobs.PushProcessEarlyMessagesJob
|
||||
import org.thoughtcrime.securesms.jobs.PushProcessMessageJob
|
||||
import org.thoughtcrime.securesms.jobs.RefreshAttributesJob
|
||||
import org.thoughtcrime.securesms.jobs.RetrieveProfileJob
|
||||
import org.thoughtcrime.securesms.jobs.SendDeliveryReceiptJob
|
||||
import org.thoughtcrime.securesms.jobs.StorageSyncJob
|
||||
import org.thoughtcrime.securesms.jobs.TrimThreadJob
|
||||
import org.thoughtcrime.securesms.jobs.UploadAttachmentToArchiveJob
|
||||
@@ -217,7 +216,7 @@ object DataMessageProcessor {
|
||||
}
|
||||
|
||||
if (metadata.sealedSender && messageId != null) {
|
||||
batchCache.addJob(SendDeliveryReceiptJob(senderRecipient.id, message.timestamp!!, messageId))
|
||||
batchCache.addDeliveryReceipt(senderRecipient.id, groupId, message.timestamp!!, messageId)
|
||||
} else if (!metadata.sealedSender) {
|
||||
if (RecipientUtil.shouldHaveProfileKey(threadRecipient)) {
|
||||
Log.w(MessageContentProcessor.TAG, "Received an unsealed sender message from " + senderRecipient.id + ", but they should already have our profile key. Correcting.")
|
||||
|
||||
@@ -2,7 +2,6 @@ package org.thoughtcrime.securesms.messages
|
||||
|
||||
import android.content.Context
|
||||
import org.signal.core.util.UuidUtil
|
||||
import org.signal.core.util.concurrent.SignalExecutors
|
||||
import org.signal.core.util.orNull
|
||||
import org.thoughtcrime.securesms.database.MessageTable.InsertResult
|
||||
import org.thoughtcrime.securesms.database.MessageType
|
||||
@@ -16,7 +15,6 @@ import org.thoughtcrime.securesms.dependencies.AppDependencies
|
||||
import org.thoughtcrime.securesms.groups.GroupId
|
||||
import org.thoughtcrime.securesms.jobs.AttachmentDownloadJob
|
||||
import org.thoughtcrime.securesms.jobs.PushProcessEarlyMessagesJob
|
||||
import org.thoughtcrime.securesms.jobs.SendDeliveryReceiptJob
|
||||
import org.thoughtcrime.securesms.messages.MessageContentProcessor.Companion.log
|
||||
import org.thoughtcrime.securesms.messages.MessageContentProcessor.Companion.warn
|
||||
import org.thoughtcrime.securesms.messages.SignalServiceProtoUtil.groupId
|
||||
@@ -46,7 +44,8 @@ object EditMessageProcessor {
|
||||
envelope: Envelope,
|
||||
content: Content,
|
||||
metadata: EnvelopeMetadata,
|
||||
earlyMessageCacheEntry: EarlyMessageCacheEntry?
|
||||
earlyMessageCacheEntry: EarlyMessageCacheEntry?,
|
||||
batchCache: BatchCache
|
||||
) {
|
||||
val editMessage = content.editMessage!!
|
||||
|
||||
@@ -92,14 +91,16 @@ object EditMessageProcessor {
|
||||
targetMessage = targetMessage.withAttachments(SignalDatabase.attachments.getAttachmentsForMessage(targetMessage.id))
|
||||
|
||||
val insertResult: InsertResult? = if (isMediaMessage || targetMessage.quote != null || targetMessage.slideDeck.slides.isNotEmpty()) {
|
||||
handleEditMediaMessage(senderRecipient.id, groupId, envelope, metadata, message, targetMessage)
|
||||
handleEditMediaMessage(senderRecipient.id, groupId, envelope, metadata, message, targetMessage, batchCache)
|
||||
} else {
|
||||
handleEditTextMessage(senderRecipient.id, groupId, envelope, metadata, message, targetMessage)
|
||||
handleEditTextMessage(senderRecipient.id, groupId, envelope, metadata, message, targetMessage, batchCache)
|
||||
}
|
||||
|
||||
if (insertResult != null) {
|
||||
SignalExecutors.BOUNDED.execute {
|
||||
AppDependencies.jobManager.add(SendDeliveryReceiptJob(senderRecipient.id, message.timestamp!!, MessageId(insertResult.messageId)))
|
||||
batchCache.addDeliveryReceipt(senderRecipient.id, groupId, message.timestamp!!, MessageId(insertResult.messageId))
|
||||
|
||||
if (insertResult.needsThreadUpdate) {
|
||||
batchCache.addIncomingMessageInsertThreadUpdate(insertResult.threadId)
|
||||
}
|
||||
|
||||
if (targetMessage.expireStarted > 0) {
|
||||
@@ -122,7 +123,8 @@ object EditMessageProcessor {
|
||||
envelope: Envelope,
|
||||
metadata: EnvelopeMetadata,
|
||||
message: DataMessage,
|
||||
targetMessage: MmsMessageRecord
|
||||
targetMessage: MmsMessageRecord,
|
||||
batchCache: BatchCache
|
||||
): InsertResult? {
|
||||
val cappedBodyRanges = message.bodyRanges.take(DataMessageProcessor.BODY_RANGE_PROCESSING_LIMIT)
|
||||
val messageRanges: BodyRangeList? = cappedBodyRanges.filter { Util.allAreNull(it.mentionAci, it.mentionAciBinary) }.toList().toBodyRangeList()
|
||||
@@ -165,7 +167,7 @@ object EditMessageProcessor {
|
||||
messageRanges = messageRanges
|
||||
)
|
||||
|
||||
val insertResult = SignalDatabase.messages.insertEditMessageInbox(mediaMessage, targetMessage).orNull()
|
||||
val insertResult = SignalDatabase.messages.insertEditMessageInbox(mediaMessage, targetMessage, skipThreadUpdate = batchCache.batchThreadUpdates).orNull()
|
||||
if (insertResult?.insertedAttachments != null) {
|
||||
SignalDatabase.runPostSuccessfulTransaction {
|
||||
val downloadJobs: List<AttachmentDownloadJob> = insertResult.insertedAttachments.mapNotNull { (_, attachmentId) ->
|
||||
@@ -183,7 +185,8 @@ object EditMessageProcessor {
|
||||
envelope: Envelope,
|
||||
metadata: EnvelopeMetadata,
|
||||
message: DataMessage,
|
||||
targetMessage: MmsMessageRecord
|
||||
targetMessage: MmsMessageRecord,
|
||||
batchCache: BatchCache
|
||||
): InsertResult? {
|
||||
val textMessage = IncomingMessage(
|
||||
type = MessageType.NORMAL,
|
||||
@@ -199,6 +202,6 @@ object EditMessageProcessor {
|
||||
serverGuid = UuidUtil.getStringUUID(envelope.serverGuid, envelope.serverGuidBinary)
|
||||
)
|
||||
|
||||
return SignalDatabase.messages.insertEditMessageInbox(textMessage, targetMessage).orNull()
|
||||
return SignalDatabase.messages.insertEditMessageInbox(textMessage, targetMessage, skipThreadUpdate = batchCache.batchThreadUpdates).orNull()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -554,7 +554,8 @@ open class MessageContentProcessor(private val context: Context) {
|
||||
envelope,
|
||||
content,
|
||||
metadata,
|
||||
if (processingEarlyContent) null else EarlyMessageCacheEntry(envelope, content, metadata, serverDeliveredTimestamp)
|
||||
if (processingEarlyContent) null else EarlyMessageCacheEntry(envelope, content, metadata, serverDeliveredTimestamp),
|
||||
batchCache
|
||||
)
|
||||
}
|
||||
|
||||
|
||||
@@ -156,6 +156,8 @@ class RecipientId private constructor(private val id: Long) : Parcelable, Compar
|
||||
|
||||
fun toScheduledSendQueueKey(): String = "RecipientId::$id::SCHEDULED"
|
||||
|
||||
fun toReceiptQueueKey(): String = "RecipientId::$id::RECEIPT"
|
||||
|
||||
override fun equals(other: Any?): Boolean {
|
||||
if (this === other) return true
|
||||
if (javaClass != other?.javaClass) return false
|
||||
|
||||
Reference in New Issue
Block a user