Improve message processing perf by delaying early processing job.

This commit is contained in:
Cody Henthorne
2026-08-19 19:05:49 -04:00
parent 6086781d68
commit 746de7d165
14 changed files with 155 additions and 72 deletions
@@ -103,7 +103,8 @@ class DataMessageProcessorTest_polls {
metadata = EnvelopeMetadata(alice.requireServiceId(), null, 1, false, null, harness.self.requireServiceId(), CiphertextMessage.WHISPER_TYPE),
threadRecipient = Recipient.resolved(groupRecipientId),
groupId = groupId,
receivedTime = 200
receivedTime = 200,
batchCache = OneTimeBatchCache()
)
assert(insertResult?.messageId != null)
@@ -124,7 +125,8 @@ class DataMessageProcessorTest_polls {
metadata = EnvelopeMetadata(alice.requireServiceId(), null, 1, false, null, harness.self.requireServiceId(), CiphertextMessage.WHISPER_TYPE),
threadRecipient = bob,
groupId = groupId,
receivedTime = 200
receivedTime = 200,
batchCache = OneTimeBatchCache()
)
assert(insertResult == null)
@@ -142,7 +144,8 @@ class DataMessageProcessorTest_polls {
metadata = EnvelopeMetadata(alice.requireServiceId(), null, 1, false, null, harness.self.requireServiceId(), CiphertextMessage.WHISPER_TYPE),
threadRecipient = bob,
groupId = groupId,
receivedTime = 200
receivedTime = 200,
batchCache = OneTimeBatchCache()
)
assert(insertResult == null)
@@ -158,7 +161,8 @@ class DataMessageProcessorTest_polls {
metadata = EnvelopeMetadata(alice.requireServiceId(), null, 1, false, null, harness.self.requireServiceId(), CiphertextMessage.WHISPER_TYPE),
threadRecipient = bob,
groupId = groupId,
receivedTime = 200
receivedTime = 200,
batchCache = OneTimeBatchCache()
)
assert(insertResult == null)
@@ -311,7 +315,8 @@ class DataMessageProcessorTest_polls {
message = DataMessage(pollVote = pollVote),
senderRecipient = senderRecipient,
threadRecipient = Recipient.resolved(groupRecipientId),
earlyMessageCacheEntry = null
earlyMessageCacheEntry = null,
batchCache = OneTimeBatchCache()
)
}
@@ -297,6 +297,11 @@ class JobController {
return jobStorage.getAllMatchingFilter(predicate);
}
@WorkerThread
synchronized List<MinimalJobSpec> findMinimalJobs(@NonNull Predicate<MinimalJobSpec> predicate) {
return jobStorage.getAllMinimalJobSpecsMatchingFilter(predicate);
}
@WorkerThread
synchronized void onRetry(@NonNull Job job, long backoffInterval) {
if (backoffInterval < 0) {
@@ -298,6 +298,19 @@ public class JobManager implements ConstraintObserver.Notifier {
return jobController.findJobs(predicate);
}
/**
* Search through the list of pending jobs and find all that match a given predicate. Unlike {@link #find(Predicate)}, this reads from the in-memory job list,
* making it dramatically cheaper when there are many jobs enqueued. The tradeoff is that the predicate can only consider the properties present on a
* {@link MinimalJobSpec}.
*
* Note that there will always be races here, and the result you get back may not be valid anymore by the time you get it. Use with caution.
*/
@WorkerThread
public @NonNull List<MinimalJobSpec> findMinimalJobs(@NonNull Predicate<MinimalJobSpec> predicate) {
waitUntilInitialized();
return jobController.findMinimalJobs(predicate);
}
/**
* Runs the specified job synchronously. Beware: All normal dependencies are respected, meaning
* you must take great care where you call this. It could take a very long time to complete!
@@ -17,6 +17,9 @@ interface JobStorage {
@WorkerThread
fun getAllMatchingFilter(predicate: Predicate<JobSpec>): List<JobSpec>
@WorkerThread
fun getAllMinimalJobSpecsMatchingFilter(predicate: Predicate<MinimalJobSpec>): List<MinimalJobSpec>
@WorkerThread
fun getNextEligibleJob(currentTime: Long, filter: (MinimalJobSpec) -> Boolean): JobSpec?
@@ -135,6 +135,11 @@ class FastJobStorage(private val jobDatabase: JobDatabase) : JobStorage {
return jobDatabase.getAllMatchingFilter(predicate)
}
@Synchronized
override fun getAllMinimalJobSpecsMatchingFilter(predicate: Predicate<MinimalJobSpec>): List<MinimalJobSpec> {
return minimalJobs.filter { predicate.test(it) }
}
@Synchronized
override fun getNextEligibleJob(currentTime: Long, filter: (MinimalJobSpec) -> Boolean): JobSpec? {
val stopwatch = debugStopwatch("get-pending")
@@ -77,12 +77,16 @@ class PushProcessEarlyMessagesJob private constructor(parameters: Parameters) :
/**
* Enqueues a job to run after the most-recently-enqueued [PushProcessMessageJob].
*
* Prefer [org.thoughtcrime.securesms.messages.BatchCache.requiresEarlyMessageProcessing] over calling this directly, so that a batch of messages only results in
* a single enqueue.
*/
@JvmStatic
fun enqueue() {
val jobManger = AppDependencies.jobManager
val youngestProcessJobId: String? = jobManger.find { it.factoryKey == PushProcessMessageJob.KEY }
val youngestProcessJobId: String? = jobManger
.findMinimalJobs { it.factoryKey == PushProcessMessageJob.KEY }
.maxByOrNull { it.createTime }
?.id
@@ -14,6 +14,7 @@ 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.PushProcessEarlyMessagesJob
import org.thoughtcrime.securesms.jobs.SendDeliveryReceiptJob
import org.thoughtcrime.securesms.messages.SignalServiceProtoUtil.groupMasterKey
import org.thoughtcrime.securesms.messages.SignalServiceProtoUtil.hasGroupContext
@@ -70,6 +71,10 @@ abstract class BatchCache {
AppDependencies.jobManager.add(job)
}
protected fun flushEarlyMessageProcessing() {
PushProcessEarlyMessagesJob.enqueue()
}
protected fun flushIncomingMessageInsertThreadUpdate(threadId: Long) {
SignalDatabase.threads.updateForMessageInsert(threadId, unarchive = true)
}
@@ -86,6 +91,7 @@ abstract class BatchCache {
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)
abstract fun requiresEarlyMessageProcessing()
}
/**
@@ -110,6 +116,10 @@ class OneTimeBatchCache : BatchCache() {
override fun addDeliveryReceipt(recipientId: RecipientId, groupId: GroupId.V2?, timestamp: Long, messageId: MessageId) {
flushDeliveryReceipt(recipientId, timestamp, messageId)
}
override fun requiresEarlyMessageProcessing() {
flushEarlyMessageProcessing()
}
}
/**
@@ -129,6 +139,7 @@ class ReusedBatchCache : BatchCache() {
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)
private var earlyMessageProcessingNeeded = false
override fun addJob(job: Job) {
batchedJobs += job
@@ -148,6 +159,10 @@ class ReusedBatchCache : BatchCache() {
accumulator.messageIds += messageId
}
override fun requiresEarlyMessageProcessing() {
earlyMessageProcessingNeeded = true
}
override fun flushAndClear() {
super.flushAndClear()
@@ -161,6 +176,11 @@ class ReusedBatchCache : BatchCache() {
}
batchedJobs.clear()
if (earlyMessageProcessingNeeded) {
flushEarlyMessageProcessing()
}
earlyMessageProcessingNeeded = false
if (threadUpdates.isNotEmpty() || mslDeletes.isNotEmpty()) {
SignalDatabase.runInTransaction {
threadUpdates.forEach { flushIncomingMessageInsertThreadUpdate(it) }
@@ -59,7 +59,6 @@ import org.thoughtcrime.securesms.jobs.GroupV2UpdateSelfProfileKeyJob
import org.thoughtcrime.securesms.jobs.PaymentLedgerUpdateJob
import org.thoughtcrime.securesms.jobs.PaymentTransactionCheckJob
import org.thoughtcrime.securesms.jobs.ProfileKeySendJob
import org.thoughtcrime.securesms.jobs.PushProcessEarlyMessagesJob
import org.thoughtcrime.securesms.jobs.PushProcessMessageJob
import org.thoughtcrime.securesms.jobs.RefreshAttributesJob
import org.thoughtcrime.securesms.jobs.RetrieveProfileJob
@@ -172,8 +171,8 @@ object DataMessageProcessor {
message.isInvalid -> handleInvalidMessage(context, senderRecipient.id, groupId, envelope.clientTimestamp!!)
message.isExpirationUpdate -> insertResult = handleExpirationUpdate(envelope, metadata, senderRecipient, threadRecipient.id, groupId, message.expireTimerDuration, message.expireTimerVersion, receivedTime, false)
message.isStoryReaction -> insertResult = handleStoryReaction(context, envelope, metadata, message, senderRecipient.id, groupId)
message.reaction != null -> messageId = handleReaction(context, envelope, message, senderRecipient.id, earlyMessageCacheEntry)
message.hasRemoteDelete -> messageId = handleRemoteDelete(context, envelope, message, senderRecipient.id, earlyMessageCacheEntry)
message.reaction != null -> messageId = handleReaction(context, envelope, message, senderRecipient.id, earlyMessageCacheEntry, batchCache)
message.hasRemoteDelete -> messageId = handleRemoteDelete(context, envelope, message, senderRecipient.id, earlyMessageCacheEntry, batchCache)
message.isPaymentActivationRequest -> insertResult = handlePaymentActivation(envelope, metadata, message, senderRecipient.id, receivedTime, isActivatePaymentsRequest = true, isPaymentsActivated = false)
message.isPaymentActivated -> insertResult = handlePaymentActivation(envelope, metadata, message, senderRecipient.id, receivedTime, isActivatePaymentsRequest = false, isPaymentsActivated = true)
message.payment != null -> insertResult = handlePayment(context, envelope, metadata, message, senderRecipient.id, receivedTime)
@@ -183,11 +182,11 @@ object DataMessageProcessor {
message.body != null -> insertResult = handleTextMessage(context, envelope, metadata, message, senderRecipient, threadRecipient, groupId, receivedTime, localMetrics, batchCache)
message.groupCallUpdate != null -> handleGroupCallUpdateMessage(envelope, senderRecipient.id, groupId)
message.pollCreate != null -> insertResult = handlePollCreate(context, envelope, metadata, message, senderRecipient, threadRecipient, groupId, receivedTime)
message.pollTerminate != null -> insertResult = handlePollTerminate(context, envelope, metadata, message, senderRecipient, earlyMessageCacheEntry, threadRecipient, groupId, receivedTime)
message.pollVote != null -> messageId = handlePollVote(context, envelope, message, senderRecipient, threadRecipient, earlyMessageCacheEntry)
message.pinMessage != null -> insertResult = handlePinMessage(envelope, metadata, message, senderRecipient, threadRecipient, groupId, receivedTime, earlyMessageCacheEntry)
message.unpinMessage != null -> messageId = handleUnpinMessage(envelope, message, senderRecipient, threadRecipient, earlyMessageCacheEntry)
message.adminDelete != null -> messageId = handleAdminRemoteDelete(context, envelope, message, senderRecipient, threadRecipient, earlyMessageCacheEntry)
message.pollTerminate != null -> insertResult = handlePollTerminate(context, envelope, metadata, message, senderRecipient, earlyMessageCacheEntry, threadRecipient, groupId, receivedTime, batchCache)
message.pollVote != null -> messageId = handlePollVote(context, envelope, message, senderRecipient, threadRecipient, earlyMessageCacheEntry, batchCache)
message.pinMessage != null -> insertResult = handlePinMessage(envelope, metadata, message, senderRecipient, threadRecipient, groupId, receivedTime, earlyMessageCacheEntry, batchCache)
message.unpinMessage != null -> messageId = handleUnpinMessage(envelope, message, senderRecipient, threadRecipient, earlyMessageCacheEntry, batchCache)
message.adminDelete != null -> messageId = handleAdminRemoteDelete(context, envelope, message, senderRecipient, threadRecipient, earlyMessageCacheEntry, batchCache)
}
SignalTrace.endSection()
@@ -509,7 +508,8 @@ object DataMessageProcessor {
envelope: Envelope,
message: DataMessage,
senderRecipientId: RecipientId,
earlyMessageCacheEntry: EarlyMessageCacheEntry?
earlyMessageCacheEntry: EarlyMessageCacheEntry?,
batchCache: BatchCache
): MessageId? {
val reaction: DataMessage.Reaction = message.reaction!!
@@ -536,7 +536,7 @@ object DataMessageProcessor {
warn(envelope.clientTimestamp!!, "[handleReaction] Could not find matching message! Putting it in the early message cache. timestamp: " + targetSentTimestamp + " author: " + targetAuthor.id)
if (earlyMessageCacheEntry != null) {
AppDependencies.earlyMessageCache.store(targetAuthor.id, targetSentTimestamp, earlyMessageCacheEntry)
PushProcessEarlyMessagesJob.enqueue()
batchCache.requiresEarlyMessageProcessing()
}
return null
}
@@ -577,7 +577,7 @@ object DataMessageProcessor {
return targetMessageId
}
fun handleRemoteDelete(context: Context, envelope: Envelope, message: DataMessage, senderRecipientId: RecipientId, earlyMessageCacheEntry: EarlyMessageCacheEntry?): MessageId? {
fun handleRemoteDelete(context: Context, envelope: Envelope, message: DataMessage, senderRecipientId: RecipientId, earlyMessageCacheEntry: EarlyMessageCacheEntry?, batchCache: BatchCache): MessageId? {
val delete = message.delete!!
log(envelope.clientTimestamp!!, "Remote delete for message ${delete.targetSentTimestamp}")
@@ -598,7 +598,7 @@ object DataMessageProcessor {
warn(envelope.clientTimestamp!!, "[handleRemoteDelete] Could not find matching message! timestamp: $targetSentTimestamp author: $senderRecipientId")
if (earlyMessageCacheEntry != null) {
AppDependencies.earlyMessageCache.store(senderRecipientId, targetSentTimestamp, earlyMessageCacheEntry)
PushProcessEarlyMessagesJob.enqueue()
batchCache.requiresEarlyMessageProcessing()
}
null
@@ -1142,7 +1142,8 @@ object DataMessageProcessor {
earlyMessageCacheEntry: EarlyMessageCacheEntry? = null,
threadRecipient: Recipient,
groupId: GroupId.V2?,
receivedTime: Long
receivedTime: Long,
batchCache: BatchCache
): InsertResult? {
val pollTerminate: DataMessage.PollTerminate = message.pollTerminate!!
val targetSentTimestamp = pollTerminate.targetSentTimestamp!!
@@ -1151,7 +1152,7 @@ object DataMessageProcessor {
handlePossibleExpirationUpdate(envelope, metadata, senderRecipient, threadRecipient, groupId, message.expireTimerDuration, message.expireTimerVersion, receivedTime)
val messageId = handlePollValidation(envelope = envelope, targetSentTimestamp = targetSentTimestamp, senderRecipient = senderRecipient, earlyMessageCacheEntry = earlyMessageCacheEntry, targetAuthor = senderRecipient, threadRecipient = threadRecipient)
val messageId = handlePollValidation(envelope = envelope, targetSentTimestamp = targetSentTimestamp, senderRecipient = senderRecipient, earlyMessageCacheEntry = earlyMessageCacheEntry, targetAuthor = senderRecipient, threadRecipient = threadRecipient, batchCache = batchCache)
if (messageId == null) {
return null
}
@@ -1191,7 +1192,8 @@ object DataMessageProcessor {
message: DataMessage,
senderRecipient: Recipient,
threadRecipient: Recipient,
earlyMessageCacheEntry: EarlyMessageCacheEntry?
earlyMessageCacheEntry: EarlyMessageCacheEntry?,
batchCache: BatchCache
): MessageId? {
val pollVote: DataMessage.PollVote = message.pollVote!!
val targetSentTimestamp = pollVote.targetSentTimestamp!!
@@ -1204,7 +1206,7 @@ object DataMessageProcessor {
return null
}
val messageId = handlePollValidation(envelope, targetSentTimestamp, senderRecipient, earlyMessageCacheEntry, Recipient.externalPush(targetAuthorServiceId), threadRecipient)
val messageId = handlePollValidation(envelope, targetSentTimestamp, senderRecipient, earlyMessageCacheEntry, Recipient.externalPush(targetAuthorServiceId), threadRecipient, batchCache)
if (messageId == null) {
return null
}
@@ -1260,7 +1262,8 @@ object DataMessageProcessor {
threadRecipient: Recipient,
groupId: GroupId.V2?,
receivedTime: Long,
earlyMessageCacheEntry: EarlyMessageCacheEntry? = null
earlyMessageCacheEntry: EarlyMessageCacheEntry? = null,
batchCache: BatchCache
): InsertResult? {
val pinMessage = message.pinMessage!!
log(envelope.clientTimestamp!!, "[handlePinMessage] Pin message for " + pinMessage.targetSentTimestamp)
@@ -1280,7 +1283,7 @@ object DataMessageProcessor {
warn(envelope.clientTimestamp!!, "[handlePinMessage] Could not find matching message! Putting it in the early message cache. timestamp: ${pinMessage.targetSentTimestamp}")
if (earlyMessageCacheEntry != null) {
AppDependencies.earlyMessageCache.store(targetAuthor.id, pinMessage.targetSentTimestamp!!, earlyMessageCacheEntry)
PushProcessEarlyMessagesJob.enqueue()
batchCache.requiresEarlyMessageProcessing()
}
return null
}
@@ -1354,7 +1357,8 @@ object DataMessageProcessor {
message: DataMessage,
senderRecipient: Recipient,
threadRecipient: Recipient,
earlyMessageCacheEntry: EarlyMessageCacheEntry? = null
earlyMessageCacheEntry: EarlyMessageCacheEntry? = null,
batchCache: BatchCache
): MessageId? {
val unpinMessage = message.unpinMessage!!
log(envelope.clientTimestamp!!, "[handleUnpinMessage] Unpin message for ${unpinMessage.targetSentTimestamp}")
@@ -1372,7 +1376,7 @@ object DataMessageProcessor {
warn(envelope.clientTimestamp!!, "[handleUnpinMessage] Could not find matching message! Putting it in the early message cache. timestamp: ${unpinMessage.targetSentTimestamp}")
if (earlyMessageCacheEntry != null) {
AppDependencies.earlyMessageCache.store(targetAuthor.id, unpinMessage.targetSentTimestamp!!, earlyMessageCacheEntry)
PushProcessEarlyMessagesJob.enqueue()
batchCache.requiresEarlyMessageProcessing()
}
return null
}
@@ -1419,7 +1423,7 @@ object DataMessageProcessor {
return MessageId(targetMessageId)
}
fun handleAdminRemoteDelete(context: Context, envelope: Envelope, message: DataMessage, senderRecipient: Recipient, threadRecipient: Recipient, earlyMessageCacheEntry: EarlyMessageCacheEntry?): MessageId? {
fun handleAdminRemoteDelete(context: Context, envelope: Envelope, message: DataMessage, senderRecipient: Recipient, threadRecipient: Recipient, earlyMessageCacheEntry: EarlyMessageCacheEntry?, batchCache: BatchCache): MessageId? {
val delete = message.adminDelete!!
log(envelope.clientTimestamp!!, "Admin delete for message ${delete.targetSentTimestamp}")
@@ -1437,7 +1441,7 @@ object DataMessageProcessor {
warn(envelope.clientTimestamp!!, "[handleAdminRemoteDelete] Could not find matching message! timestamp: $targetSentTimestamp")
if (earlyMessageCacheEntry != null) {
AppDependencies.earlyMessageCache.store(targetAuthor.id, targetSentTimestamp, earlyMessageCacheEntry)
PushProcessEarlyMessagesJob.enqueue()
batchCache.requiresEarlyMessageProcessing()
}
return null
}
@@ -1614,14 +1618,15 @@ object DataMessageProcessor {
senderRecipient: Recipient,
earlyMessageCacheEntry: EarlyMessageCacheEntry?,
targetAuthor: Recipient,
threadRecipient: Recipient
threadRecipient: Recipient,
batchCache: BatchCache
): MessageId? {
val targetMessage = SignalDatabase.messages.getMessageFor(targetSentTimestamp, targetAuthor.id)
if (targetMessage == null) {
warn(envelope.clientTimestamp!!, "[handlePollValidation] Could not find matching message! Putting it in the early message cache. timestamp: $targetSentTimestamp author: ${targetAuthor.id}")
if (earlyMessageCacheEntry != null) {
AppDependencies.earlyMessageCache.store(senderRecipient.id, targetSentTimestamp, earlyMessageCacheEntry)
PushProcessEarlyMessagesJob.enqueue()
batchCache.requiresEarlyMessageProcessing()
}
return null
}
@@ -14,7 +14,6 @@ import org.thoughtcrime.securesms.database.withAttachments
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.messages.MessageContentProcessor.Companion.log
import org.thoughtcrime.securesms.messages.MessageContentProcessor.Companion.warn
import org.thoughtcrime.securesms.messages.SignalServiceProtoUtil.groupId
@@ -59,7 +58,7 @@ object EditMessageProcessor {
if (earlyMessageCacheEntry != null) {
AppDependencies.earlyMessageCache.store(senderRecipient.id, editMessage.targetSentTimestamp!!, earlyMessageCacheEntry)
PushProcessEarlyMessagesJob.enqueue()
batchCache.requiresEarlyMessageProcessing()
}
return
@@ -583,7 +583,7 @@ class IncomingMessageObserver(
/**
* Attempts to process the entire batch in a single transaction for performance.
*
* @return true if the transaction committed, false if it the batch was rolled back.
* @return true if the transaction committed, false if the batch was rolled back.
*/
private fun processBatchInTransaction(batch: List<EnvelopeResponse>): Boolean {
val allFollowUpOperations = mutableListOf<FollowUpOperation>()
@@ -499,7 +499,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
)
}
@@ -5,7 +5,6 @@ import android.content.Context
import org.signal.core.util.Stopwatch
import org.thoughtcrime.securesms.database.SignalDatabase
import org.thoughtcrime.securesms.dependencies.AppDependencies
import org.thoughtcrime.securesms.jobs.PushProcessEarlyMessagesJob
import org.thoughtcrime.securesms.keyvalue.SignalStore
import org.thoughtcrime.securesms.messages.MessageContentProcessor.Companion.log
import org.thoughtcrime.securesms.messages.MessageContentProcessor.Companion.warn
@@ -30,7 +29,7 @@ object ReceiptMessageProcessor {
when (receiptMessage.type) {
ReceiptMessage.Type.DELIVERY -> handleDeliveryReceipt(envelope, metadata, receiptMessage, senderRecipient.id, batchCache)
ReceiptMessage.Type.READ -> handleReadReceipt(context, senderRecipient.id, envelope, metadata, receiptMessage, earlyMessageCacheEntry, batchCache)
ReceiptMessage.Type.VIEWED -> handleViewedReceipt(context, envelope, metadata, receiptMessage, senderRecipient.id, earlyMessageCacheEntry)
ReceiptMessage.Type.VIEWED -> handleViewedReceipt(context, envelope, metadata, receiptMessage, senderRecipient.id, earlyMessageCacheEntry, batchCache)
else -> warn(envelope.clientTimestamp!!, "Unknown recipient message type ${receiptMessage.type}")
}
}
@@ -56,7 +55,7 @@ object ReceiptMessageProcessor {
}
if (missingTargetTimestamps.isNotEmpty()) {
PushProcessEarlyMessagesJob.enqueue()
batchCache.requiresEarlyMessageProcessing()
}
SignalDatabase.pendingPniSignatureMessages.acknowledgeReceipts(senderRecipientId, deliveryReceipt.timestamp, metadata.sourceDeviceId)
@@ -101,7 +100,7 @@ object ReceiptMessageProcessor {
}
if (missingTargetTimestamps.isNotEmpty() && earlyMessageCacheEntry != null) {
PushProcessEarlyMessagesJob.enqueue()
batchCache.requiresEarlyMessageProcessing()
}
}
@@ -111,7 +110,8 @@ object ReceiptMessageProcessor {
metadata: EnvelopeMetadata,
viewedReceipt: ReceiptMessage,
senderRecipientId: RecipientId,
earlyMessageCacheEntry: EarlyMessageCacheEntry?
earlyMessageCacheEntry: EarlyMessageCacheEntry?,
batchCache: BatchCache
) {
val readReceipts = TextSecurePreferences.isReadReceiptsEnabled(context)
val storyViewedReceipts = SignalStore.story.viewedReceiptsEnabled
@@ -146,7 +146,7 @@ object ReceiptMessageProcessor {
}
if (missingTargetTimestamps.isNotEmpty() && earlyMessageCacheEntry != null) {
PushProcessEarlyMessagesJob.enqueue()
batchCache.requiresEarlyMessageProcessing()
}
}
}
@@ -77,7 +77,6 @@ import org.thoughtcrime.securesms.jobs.MultiDeviceContactUpdateJob
import org.thoughtcrime.securesms.jobs.MultiDeviceKeysUpdateJob
import org.thoughtcrime.securesms.jobs.MultiDeviceStickerPackSyncJob
import org.thoughtcrime.securesms.jobs.PreKeysSyncJob
import org.thoughtcrime.securesms.jobs.PushProcessEarlyMessagesJob
import org.thoughtcrime.securesms.jobs.RefreshCallLinkDetailsJob
import org.thoughtcrime.securesms.jobs.RefreshDonationSubscriptionStatusJob
import org.thoughtcrime.securesms.jobs.RefreshOwnProfileJob
@@ -168,16 +167,17 @@ object SyncMessageProcessor {
envelope: Envelope,
content: Content,
metadata: EnvelopeMetadata,
earlyMessageCacheEntry: EarlyMessageCacheEntry?
earlyMessageCacheEntry: EarlyMessageCacheEntry?,
batchCache: BatchCache
) {
val syncMessage = content.syncMessage!!
when {
syncMessage.sent != null -> handleSynchronizeSentMessage(context, envelope, content, metadata, syncMessage.sent!!, senderRecipient, threadRecipient, earlyMessageCacheEntry)
syncMessage.sent != null -> handleSynchronizeSentMessage(context, envelope, content, metadata, syncMessage.sent!!, senderRecipient, threadRecipient, earlyMessageCacheEntry, batchCache)
syncMessage.request != null -> handleSynchronizeRequestMessage(context, syncMessage.request!!, envelope.clientTimestamp!!)
syncMessage.read.isNotEmpty() -> handleSynchronizeReadMessage(context, syncMessage.read, envelope.clientTimestamp!!, earlyMessageCacheEntry)
syncMessage.read.isNotEmpty() -> handleSynchronizeReadMessage(context, syncMessage.read, envelope.clientTimestamp!!, earlyMessageCacheEntry, batchCache)
syncMessage.viewed.isNotEmpty() -> handleSynchronizeViewedMessage(context, syncMessage.viewed, envelope.clientTimestamp!!)
syncMessage.viewOnceOpen != null -> handleSynchronizeViewOnceOpenMessage(context, syncMessage.viewOnceOpen!!, envelope.clientTimestamp!!, earlyMessageCacheEntry)
syncMessage.viewOnceOpen != null -> handleSynchronizeViewOnceOpenMessage(context, syncMessage.viewOnceOpen!!, envelope.clientTimestamp!!, earlyMessageCacheEntry, batchCache)
syncMessage.verified != null -> handleSynchronizeVerifiedMessage(context, syncMessage.verified!!)
syncMessage.stickerPackOperation.isNotEmpty() -> handleSynchronizeStickerPackOperation(syncMessage.stickerPackOperation, envelope.clientTimestamp!!)
syncMessage.configuration != null -> handleSynchronizeConfigurationMessage(context, syncMessage.configuration!!, envelope.clientTimestamp!!)
@@ -191,7 +191,7 @@ object SyncMessageProcessor {
syncMessage.callEvent != null -> handleSynchronizeCallEvent(syncMessage.callEvent!!, envelope.clientTimestamp!!)
syncMessage.callLinkUpdate != null -> handleSynchronizeCallLink(syncMessage.callLinkUpdate!!, envelope.clientTimestamp!!)
syncMessage.callLogEvent != null -> handleSynchronizeCallLogEvent(syncMessage.callLogEvent!!, envelope.clientTimestamp!!)
syncMessage.deleteForMe != null -> handleSynchronizeDeleteForMe(context, syncMessage.deleteForMe!!, envelope.clientTimestamp!!, earlyMessageCacheEntry)
syncMessage.deleteForMe != null -> handleSynchronizeDeleteForMe(context, syncMessage.deleteForMe!!, envelope.clientTimestamp!!, earlyMessageCacheEntry, batchCache)
syncMessage.attachmentBackfillRequest != null -> handleSynchronizeAttachmentBackfillRequest(syncMessage.attachmentBackfillRequest!!, envelope.clientTimestamp!!)
syncMessage.attachmentBackfillResponse != null -> handleSynchronizeAttachmentBackfillResponse(syncMessage.attachmentBackfillResponse!!, envelope.clientTimestamp!!)
syncMessage.usernameChange != null -> handleSynchronizeUsernameChange(envelope.clientTimestamp!!)
@@ -208,7 +208,8 @@ object SyncMessageProcessor {
sent: Sent,
senderRecipient: Recipient,
threadRecipient: Recipient,
earlyMessageCacheEntry: EarlyMessageCacheEntry?
earlyMessageCacheEntry: EarlyMessageCacheEntry?,
batchCache: BatchCache
) {
log(envelope.clientTimestamp!!, "Processing sent transcript for message with ID ${sent.timestamp!!}")
@@ -222,7 +223,7 @@ object SyncMessageProcessor {
}
if (sent.editMessage != null) {
handleSynchronizeSentEditMessage(context, envelope, sent, senderRecipient, earlyMessageCacheEntry)
handleSynchronizeSentEditMessage(context, envelope, sent, senderRecipient, earlyMessageCacheEntry, batchCache)
SignalStore.misc.lastSyncMessageSeenTimeMs = System.currentTimeMillis()
return
}
@@ -258,27 +259,27 @@ object SyncMessageProcessor {
dataMessage.isExpirationUpdate -> threadId = handleSynchronizeSentExpirationUpdate(sent)
dataMessage.storyContext != null -> threadId = handleSynchronizeSentStoryReply(sent, envelope.clientTimestamp!!)
dataMessage.reaction != null -> {
DataMessageProcessor.handleReaction(context, envelope, dataMessage, senderRecipient.id, earlyMessageCacheEntry)
DataMessageProcessor.handleReaction(context, envelope, dataMessage, senderRecipient.id, earlyMessageCacheEntry, batchCache)
threadId = SignalDatabase.threads.getOrCreateThreadIdFor(getSyncMessageDestination(sent))
}
dataMessage.hasRemoteDelete -> DataMessageProcessor.handleRemoteDelete(context, envelope, dataMessage, senderRecipient.id, earlyMessageCacheEntry)
dataMessage.hasRemoteDelete -> DataMessageProcessor.handleRemoteDelete(context, envelope, dataMessage, senderRecipient.id, earlyMessageCacheEntry, batchCache)
dataMessage.payment != null -> log(envelope.clientTimestamp!!, "Ignoring payment notification/activation from sync transcript; payment row arrives via SyncMessage.OutgoingPayment.")
dataMessage.isMediaMessage -> threadId = handleSynchronizeSentMediaMessage(context, sent, envelope.clientTimestamp!!, senderRecipient)
dataMessage.pollCreate != null -> threadId = handleSynchronizedPollCreate(envelope, dataMessage, sent, senderRecipient)
dataMessage.pollVote != null -> {
val destination = getSyncMessageDestination(sent)
DataMessageProcessor.handlePollVote(context, envelope, dataMessage, senderRecipient, destination, earlyMessageCacheEntry)
DataMessageProcessor.handlePollVote(context, envelope, dataMessage, senderRecipient, destination, earlyMessageCacheEntry, batchCache)
threadId = SignalDatabase.threads.getOrCreateThreadIdFor(getSyncMessageDestination(sent))
}
dataMessage.pollTerminate != null -> threadId = handleSynchronizedPollEnd(envelope, dataMessage, sent, senderRecipient, earlyMessageCacheEntry)
dataMessage.pinMessage != null -> threadId = handleSynchronizedPinMessage(envelope, dataMessage, sent, senderRecipient, earlyMessageCacheEntry)
dataMessage.pollTerminate != null -> threadId = handleSynchronizedPollEnd(envelope, dataMessage, sent, senderRecipient, earlyMessageCacheEntry, batchCache)
dataMessage.pinMessage != null -> threadId = handleSynchronizedPinMessage(envelope, dataMessage, sent, senderRecipient, earlyMessageCacheEntry, batchCache)
dataMessage.unpinMessage != null -> {
val destination = getSyncMessageDestination(sent)
DataMessageProcessor.handleUnpinMessage(envelope, dataMessage, senderRecipient, destination, earlyMessageCacheEntry)
DataMessageProcessor.handleUnpinMessage(envelope, dataMessage, senderRecipient, destination, earlyMessageCacheEntry, batchCache)
threadId = SignalDatabase.threads.getOrCreateThreadIdFor(destination)
}
dataMessage.adminDelete != null -> {
DataMessageProcessor.handleAdminRemoteDelete(context, envelope, dataMessage, senderRecipient, threadRecipient, earlyMessageCacheEntry)
DataMessageProcessor.handleAdminRemoteDelete(context, envelope, dataMessage, senderRecipient, threadRecipient, earlyMessageCacheEntry, batchCache)
threadId = SignalDatabase.threads.getOrCreateThreadIdFor(getSyncMessageDestination(sent))
}
else -> threadId = handleSynchronizeSentTextMessage(sent, envelope.clientTimestamp!!)
@@ -353,7 +354,8 @@ object SyncMessageProcessor {
envelope: Envelope,
sent: Sent,
senderRecipient: Recipient,
earlyMessageCacheEntry: EarlyMessageCacheEntry?
earlyMessageCacheEntry: EarlyMessageCacheEntry?,
batchCache: BatchCache
) {
val editMessage: EditMessage = sent.editMessage!!
val targetSentTimestamp: Long = editMessage.targetSentTimestamp!!
@@ -364,7 +366,7 @@ object SyncMessageProcessor {
warn(envelope.clientTimestamp!!, "[handleSynchronizeSentEditMessage] Could not find matching message! targetTimestamp: $targetSentTimestamp author: $senderRecipientId")
if (earlyMessageCacheEntry != null) {
AppDependencies.earlyMessageCache.store(senderRecipientId, targetSentTimestamp, earlyMessageCacheEntry)
PushProcessEarlyMessagesJob.enqueue()
batchCache.requiresEarlyMessageProcessing()
}
} else if (MessageConstraintsUtil.isValidEditMessageReceive(targetMessage, senderRecipient, envelope.serverTimestamp!!)) {
val message: DataMessage = editMessage.dataMessage!!
@@ -1005,7 +1007,8 @@ object SyncMessageProcessor {
context: Context,
readMessages: List<Read>,
envelopeTimestamp: Long,
earlyMessageCacheEntry: EarlyMessageCacheEntry?
earlyMessageCacheEntry: EarlyMessageCacheEntry?,
batchCache: BatchCache
) {
log(envelopeTimestamp, "Synchronize read message. Count: ${readMessages.size}, Timestamps: ${readMessages.map { it.timestamp }}")
@@ -1026,7 +1029,7 @@ object SyncMessageProcessor {
}
if (unhandled.isNotEmpty() && earlyMessageCacheEntry != null) {
PushProcessEarlyMessagesJob.enqueue()
batchCache.requiresEarlyMessageProcessing()
}
SignalStore.misc.lastSyncMessageSeenTimeMs = System.currentTimeMillis()
@@ -1074,7 +1077,7 @@ object SyncMessageProcessor {
}
}
private fun handleSynchronizeViewOnceOpenMessage(context: Context, openMessage: ViewOnceOpen, envelopeTimestamp: Long, earlyMessageCacheEntry: EarlyMessageCacheEntry?) {
private fun handleSynchronizeViewOnceOpenMessage(context: Context, openMessage: ViewOnceOpen, envelopeTimestamp: Long, earlyMessageCacheEntry: EarlyMessageCacheEntry?, batchCache: BatchCache) {
log(envelopeTimestamp, "Handling a view-once open for message: " + openMessage.timestamp)
val author: RecipientId = Recipient.externalPush(ACI.parseOrThrow(openMessage.senderAci, openMessage.senderAciBinary)).id
@@ -1092,7 +1095,7 @@ object SyncMessageProcessor {
warn(envelopeTimestamp.toString(), "Got a view-once open message for a message we don't have!")
if (earlyMessageCacheEntry != null) {
AppDependencies.earlyMessageCache.store(author, timestamp, earlyMessageCacheEntry)
PushProcessEarlyMessagesJob.enqueue()
batchCache.requiresEarlyMessageProcessing()
}
}
@@ -1603,11 +1606,11 @@ object SyncMessageProcessor {
return Recipient.resolved(callLink.recipientId)
}
private fun handleSynchronizeDeleteForMe(context: Context, deleteForMe: SyncMessage.DeleteForMe, envelopeTimestamp: Long, earlyMessageCacheEntry: EarlyMessageCacheEntry?) {
private fun handleSynchronizeDeleteForMe(context: Context, deleteForMe: SyncMessage.DeleteForMe, envelopeTimestamp: Long, earlyMessageCacheEntry: EarlyMessageCacheEntry?, batchCache: BatchCache) {
log(envelopeTimestamp, "Synchronize delete message messageDeletes=${deleteForMe.messageDeletes.size} conversationDeletes=${deleteForMe.conversationDeletes.size} localOnlyConversationDeletes=${deleteForMe.localOnlyConversationDeletes.size}")
if (deleteForMe.messageDeletes.isNotEmpty()) {
handleSynchronizeMessageDeletes(deleteForMe.messageDeletes, envelopeTimestamp, earlyMessageCacheEntry)
handleSynchronizeMessageDeletes(deleteForMe.messageDeletes, envelopeTimestamp, earlyMessageCacheEntry, batchCache)
}
if (deleteForMe.conversationDeletes.isNotEmpty()) {
@@ -1619,13 +1622,13 @@ object SyncMessageProcessor {
}
if (deleteForMe.attachmentDeletes.isNotEmpty()) {
handleSynchronizeAttachmentDeletes(deleteForMe.attachmentDeletes, envelopeTimestamp, earlyMessageCacheEntry)
handleSynchronizeAttachmentDeletes(deleteForMe.attachmentDeletes, envelopeTimestamp, earlyMessageCacheEntry, batchCache)
}
AppDependencies.messageNotifier.updateNotification(context)
}
private fun handleSynchronizeMessageDeletes(messageDeletes: List<SyncMessage.DeleteForMe.MessageDeletes>, envelopeTimestamp: Long, earlyMessageCacheEntry: EarlyMessageCacheEntry?) {
private fun handleSynchronizeMessageDeletes(messageDeletes: List<SyncMessage.DeleteForMe.MessageDeletes>, envelopeTimestamp: Long, earlyMessageCacheEntry: EarlyMessageCacheEntry?, batchCache: BatchCache) {
val messagesToDelete: List<MessageTable.SyncMessageId> = messageDeletes
.asSequence()
.map { it.messages }
@@ -1643,7 +1646,7 @@ object SyncMessageProcessor {
}
if (unhandled.isNotEmpty() && earlyMessageCacheEntry != null) {
PushProcessEarlyMessagesJob.enqueue()
batchCache.requiresEarlyMessageProcessing()
}
}
@@ -1707,7 +1710,7 @@ object SyncMessageProcessor {
}
}
private fun handleSynchronizeAttachmentDeletes(attachmentDeletes: List<SyncMessage.DeleteForMe.AttachmentDelete>, envelopeTimestamp: Long, earlyMessageCacheEntry: EarlyMessageCacheEntry?) {
private fun handleSynchronizeAttachmentDeletes(attachmentDeletes: List<SyncMessage.DeleteForMe.AttachmentDelete>, envelopeTimestamp: Long, earlyMessageCacheEntry: EarlyMessageCacheEntry?, batchCache: BatchCache) {
val toDelete: List<AttachmentTable.SyncAttachmentId> = attachmentDeletes
.mapNotNull { delete ->
delete.toSyncAttachmentId(delete.targetMessage?.toSyncMessageId(envelopeTimestamp), envelopeTimestamp)
@@ -1723,7 +1726,7 @@ object SyncMessageProcessor {
}
if (unhandled.isNotEmpty() && earlyMessageCacheEntry != null) {
PushProcessEarlyMessagesJob.enqueue()
batchCache.requiresEarlyMessageProcessing()
}
}
@@ -2063,7 +2066,8 @@ object SyncMessageProcessor {
message: DataMessage,
sent: Sent,
senderRecipient: Recipient,
earlyMessageCacheEntry: EarlyMessageCacheEntry?
earlyMessageCacheEntry: EarlyMessageCacheEntry?,
batchCache: BatchCache
): Long {
log(envelope.clientTimestamp!!, "Synchronize sent poll terminate message")
@@ -2081,7 +2085,7 @@ object SyncMessageProcessor {
warn(envelope.clientTimestamp!!, "Unable to find target message for poll termination. Putting in early message cache.")
if (earlyMessageCacheEntry != null) {
AppDependencies.earlyMessageCache.store(senderRecipient.id, pollTerminate.targetSentTimestamp!!, earlyMessageCacheEntry)
PushProcessEarlyMessagesJob.enqueue()
batchCache.requiresEarlyMessageProcessing()
}
return -1
}
@@ -2130,7 +2134,8 @@ object SyncMessageProcessor {
message: DataMessage,
sent: Sent,
senderRecipient: Recipient,
earlyMessageCacheEntry: EarlyMessageCacheEntry?
earlyMessageCacheEntry: EarlyMessageCacheEntry?,
batchCache: BatchCache
): Long {
log(envelope.clientTimestamp!!, "Synchronize pinned message")
@@ -2155,7 +2160,7 @@ object SyncMessageProcessor {
warn(envelope.clientTimestamp!!, "Unable to find target message for sync message. Putting in early message cache.")
if (earlyMessageCacheEntry != null) {
AppDependencies.earlyMessageCache.store(senderRecipient.id, pinMessage.targetSentTimestamp!!, earlyMessageCacheEntry)
PushProcessEarlyMessagesJob.enqueue()
batchCache.requiresEarlyMessageProcessing()
}
return -1
}
@@ -1,6 +1,7 @@
package org.thoughtcrime.securesms.jobs
import assertk.assertThat
import assertk.assertions.isEmpty
import assertk.assertions.isEqualTo
import assertk.assertions.isNotEqualTo
import assertk.assertions.isNotNull
@@ -1056,6 +1057,23 @@ class FastJobStorageTest {
assertThat(subject.getJobCountForFactoryAndQueue("f1", "does-not-exist")).isEqualTo(0)
}
@Test
fun `getAllMinimalJobSpecsMatchingFilter - general`() {
val fullSpec1 = FullSpec(jobSpec(id = "1", factoryKey = "f1", createTime = 1), emptyList(), emptyList())
val fullSpec2 = FullSpec(jobSpec(id = "2", factoryKey = "f1", createTime = 3), emptyList(), emptyList())
val fullSpec3 = FullSpec(jobSpec(id = "3", factoryKey = "f2", queueKey = "q1", createTime = 2), emptyList(), emptyList())
val subject = FastJobStorage(mockDatabase(listOf(fullSpec1, fullSpec2, fullSpec3)))
subject.init()
val f1Jobs = subject.getAllMinimalJobSpecsMatchingFilter { it.factoryKey == "f1" }
assertThat(f1Jobs.map { it.id }.toSet()).isEqualTo(setOf("1", "2"))
assertThat(f1Jobs.maxByOrNull { it.createTime }?.id).isEqualTo("2")
assertThat(subject.getAllMinimalJobSpecsMatchingFilter { it.queueKey == "q1" }.map { it.id }).isEqualTo(listOf("3"))
assertThat(subject.getAllMinimalJobSpecsMatchingFilter { it.factoryKey == "does-not-exist" }).isEmpty()
}
@Test
fun `areQueuesEmpty - all non-empty`() {
val subject = FastJobStorage(mockDatabase(DataSet1.FULL_SPECS))