Do not immediately reset IMO on group lock acquire failure.

This commit is contained in:
Cody Henthorne
2026-09-30 13:34:03 -03:00
committed by Alex Hart
parent 0b37d63907
commit 83be8338a9
@@ -22,6 +22,7 @@ import org.signal.core.util.SafeForegroundService
import org.signal.core.util.SleepTimer
import org.signal.core.util.UptimeSleepTimer
import org.signal.core.util.concurrent.SignalExecutors
import org.signal.core.util.groups.GroupChangeBusyException
import org.signal.core.util.logging.Log
import org.signal.network.config.ProxyConfig
import org.signal.network.config.SignalServiceConfiguration
@@ -57,6 +58,7 @@ import org.whispersystems.signalservice.api.websocket.SignalWebSocket
import org.whispersystems.signalservice.api.websocket.WebSocketConnectionState
import org.whispersystems.signalservice.api.websocket.WebSocketUnavailableException
import org.whispersystems.signalservice.internal.push.Envelope
import java.io.Closeable
import java.util.concurrent.CopyOnWriteArrayList
import java.util.concurrent.TimeUnit
import java.util.concurrent.TimeoutException
@@ -85,6 +87,8 @@ class IncomingMessageObserver(
private const val WEB_SOCKET_KEEP_ALIVE_TOKEN = "MessageRetrieval"
private const val MAX_GROUP_LOCK_ATTEMPTS = 6
/** How long we wait for the websocket to time out before we try to connect again. */
private val websocketReadTimeout: Long
get() = if (censored) 30.seconds.inWholeMilliseconds else 1.minutes.inWholeMilliseconds
@@ -531,7 +535,7 @@ class IncomingMessageObserver(
Log.i(TAG, "Retrieved ${batch.size} envelopes!")
val startTime = System.currentTimeMillis()
GroupsV2ProcessingLock.acquireGroupProcessingLock().use {
acquireGroupProcessingLock().use {
ReentrantSessionLock.INSTANCE.acquire().use {
val batchCommitted = processBatchInTransaction(batch)
@@ -602,6 +606,24 @@ class IncomingMessageObserver(
Log.w(TAG, "Terminated! (${this.hashCode()})")
}
/**
* Retries the group processing lock so a slow group job does not immediately tear down the websocket and force the batch to be redelivered.
*/
private fun acquireGroupProcessingLock(): Closeable {
var attempt = 1
while (true) {
try {
return GroupsV2ProcessingLock.acquireGroupProcessingLock()
} catch (e: GroupChangeBusyException) {
if (terminated || attempt >= MAX_GROUP_LOCK_ATTEMPTS) {
throw e
}
Log.w(TAG, "Group processing lock is busy, waiting again. Attempt $attempt of $MAX_GROUP_LOCK_ATTEMPTS. ${e.message}")
attempt++
}
}
}
/**
* Attempts to process the entire batch in a single transaction for performance.
*