diff --git a/app/src/main/java/org/session/libsession/messaging/sending_receiving/pollers/ConfigTtlExtensionThrottle.kt b/app/src/main/java/org/session/libsession/messaging/sending_receiving/pollers/ConfigTtlExtensionThrottle.kt new file mode 100644 index 0000000000..6ceba41e61 --- /dev/null +++ b/app/src/main/java/org/session/libsession/messaging/sending_receiving/pollers/ConfigTtlExtensionThrottle.kt @@ -0,0 +1,51 @@ +package org.session.libsession.messaging.sending_receiving.pollers + +import java.util.concurrent.ConcurrentHashMap +import javax.inject.Inject +import javax.inject.Singleton +import kotlin.time.ComparableTimeMark +import kotlin.time.Duration.Companion.hours +import kotlin.time.TimeSource + +/** + * Limits how often a poller asks the swarm to extend the TTL of our config messages. + * + * Every extend is a write on every storage node holding those messages, and pollers run every few + * seconds, so extending on each poll puts real disk I/O load on service nodes for no benefit: the + * extension is to weeks from now, so doing it once an hour loses nothing. + * + * Tracked per swarm, so one group's renewal never suppresses another group's or the user's own. + */ +@Singleton +class ConfigTtlExtensionThrottle( + private val timeSource: TimeSource.WithComparableMarks, +) { + @Inject + constructor() : this(TimeSource.Monotonic) + + private val lastSuccessfulExtension = ConcurrentHashMap() + + /** + * Runs [extend] unless an extension for [swarmPubKeyHex] succeeded within [COOLDOWN]. + * + * The cooldown starts only when [extend] returns normally. A failed extension that started it would + * leave the configs un-renewed for an hour while looking handled, and repeated failures would let + * them age out of the swarm — so any exception propagates and the next poll tries again. + * + * @return whether [extend] ran. + */ + suspend fun extendIfDue(swarmPubKeyHex: String, extend: suspend () -> Unit): Boolean { + val last = lastSuccessfulExtension[swarmPubKeyHex] + if (last != null && last.elapsedNow() < COOLDOWN) { + return false + } + + extend() + lastSuccessfulExtension[swarmPubKeyHex] = timeSource.markNow() + return true + } + + companion object { + val COOLDOWN = 1.hours + } +} diff --git a/app/src/main/java/org/session/libsession/messaging/sending_receiving/pollers/Poller.kt b/app/src/main/java/org/session/libsession/messaging/sending_receiving/pollers/Poller.kt index 55210422ed..dbf1275f55 100644 --- a/app/src/main/java/org/session/libsession/messaging/sending_receiving/pollers/Poller.kt +++ b/app/src/main/java/org/session/libsession/messaging/sending_receiving/pollers/Poller.kt @@ -56,6 +56,7 @@ class Poller @Inject constructor( private val swarmSnodeSelector: SwarmSnodeSelector, private val swarmDirectory: SwarmDirectory, private val snodeApiExecutor: SnodeApiExecutor, + private val configTtlExtensionThrottle: ConfigTtlExtensionThrottle, appVisibilityManager: AppVisibilityManager, ) : BasePoller( debugLabel = "MainPoller", @@ -248,18 +249,20 @@ class Poller @Inject constructor( if (hashesToExtend.isNotEmpty()) { launch { try { - swarmApiExecutor.execute( - SwarmApiRequest( - swarmPubKeyHex = userAuth.accountId.hexString, - api = alterTtlApiFactory.create( - messageHashes = hashesToExtend, - auth = userAuth, - alterType = AlterTtlApi.AlterType.Extend, - newExpiry = snodeClock.currentTimeMillis() + 14.days.inWholeMilliseconds - ), - swarmNodeOverride = snode, + configTtlExtensionThrottle.extendIfDue(userAuth.accountId.hexString) { + swarmApiExecutor.execute( + SwarmApiRequest( + swarmPubKeyHex = userAuth.accountId.hexString, + api = alterTtlApiFactory.create( + messageHashes = hashesToExtend, + auth = userAuth, + alterType = AlterTtlApi.AlterType.Extend, + newExpiry = snodeClock.currentTimeMillis() + 14.days.inWholeMilliseconds + ), + swarmNodeOverride = snode, + ) ) - ) + } } catch (e: Exception) { if (e is CancellationException) throw e diff --git a/app/src/main/java/org/thoughtcrime/securesms/groups/GroupPoller.kt b/app/src/main/java/org/thoughtcrime/securesms/groups/GroupPoller.kt index f5fd906d50..4196267a8c 100644 --- a/app/src/main/java/org/thoughtcrime/securesms/groups/GroupPoller.kt +++ b/app/src/main/java/org/thoughtcrime/securesms/groups/GroupPoller.kt @@ -13,6 +13,7 @@ import network.loki.messenger.libsession_util.Namespace import org.session.libsession.messaging.sending_receiving.MessageParser import org.session.libsession.messaging.sending_receiving.ReceivedMessageProcessor import org.session.libsession.messaging.sending_receiving.pollers.BasePoller +import org.session.libsession.messaging.sending_receiving.pollers.ConfigTtlExtensionThrottle import org.session.libsession.network.SnodeClock import org.session.libsession.snode.model.RetrieveMessageResponse import org.session.libsession.utilities.Address @@ -52,6 +53,7 @@ class GroupPoller @AssistedInject constructor( private val alterTtlApiApiFactory: AlterTtlApi.Factory, private val swarmApiExecutor: SwarmApiExecutor, private val swarmSnodeSelector: SwarmSnodeSelector, + private val configTtlExtensionThrottle: ConfigTtlExtensionThrottle, networkConnectivity: NetworkConnectivity, appVisibilityManager: AppVisibilityManager, ): BasePoller( @@ -126,18 +128,20 @@ class GroupPoller @AssistedInject constructor( if (configHashesToExtends.isNotEmpty() && adminKey != null) { pollingTasks += "extending group config TTL" to async { - swarmApiExecutor.execute( - SwarmApiRequest( - swarmNodeOverride = snode, - swarmPubKeyHex = groupId.hexString, - api = alterTtlApiApiFactory.create( - messageHashes = configHashesToExtends, - auth = groupAuth, - alterType = AlterTtlApi.AlterType.Extend, - newExpiry = clock.currentTimeMillis() + 14.days.inWholeMilliseconds, + configTtlExtensionThrottle.extendIfDue(groupId.hexString) { + swarmApiExecutor.execute( + SwarmApiRequest( + swarmNodeOverride = snode, + swarmPubKeyHex = groupId.hexString, + api = alterTtlApiApiFactory.create( + messageHashes = configHashesToExtends, + auth = groupAuth, + alterType = AlterTtlApi.AlterType.Extend, + newExpiry = clock.currentTimeMillis() + 14.days.inWholeMilliseconds, + ) ) ) - ) + } } } diff --git a/app/src/test/java/org/session/libsession/messaging/sending_receiving/pollers/ConfigTtlExtensionThrottleTest.kt b/app/src/test/java/org/session/libsession/messaging/sending_receiving/pollers/ConfigTtlExtensionThrottleTest.kt new file mode 100644 index 0000000000..52215acfbc --- /dev/null +++ b/app/src/test/java/org/session/libsession/messaging/sending_receiving/pollers/ConfigTtlExtensionThrottleTest.kt @@ -0,0 +1,112 @@ +package org.session.libsession.messaging.sending_receiving.pollers + +import kotlinx.coroutines.CancellationException +import kotlinx.coroutines.test.runTest +import org.junit.Test +import kotlin.test.assertEquals +import kotlin.test.assertFailsWith +import kotlin.test.assertFalse +import kotlin.test.assertTrue +import kotlin.time.Duration.Companion.minutes +import kotlin.time.Duration.Companion.seconds +import kotlin.time.TestTimeSource + +class ConfigTtlExtensionThrottleTest { + + private val userSwarm = "05" + "a".repeat(64) + private val groupSwarm = "03" + "b".repeat(64) + private val otherGroupSwarm = "03" + "c".repeat(64) + + private val time = TestTimeSource() + private val throttle = ConfigTtlExtensionThrottle(time) + + private var extensionsSent = 0 + private val succeed: suspend () -> Unit = { extensionsSent++ } + private val fail: suspend () -> Unit = { + extensionsSent++ + throw RuntimeException("storage server rejected the extension") + } + + @Test + fun `two polls inside the window send one extension`() = runTest { + assertTrue(throttle.extendIfDue(userSwarm, succeed)) + time += 30.minutes + assertFalse(throttle.extendIfDue(userSwarm, succeed)) + + assertEquals(1, extensionsSent) + } + + @Test + fun `the window is still closed just before an hour`() = runTest { + throttle.extendIfDue(userSwarm, succeed) + time += ConfigTtlExtensionThrottle.COOLDOWN - 1.seconds + + assertFalse(throttle.extendIfDue(userSwarm, succeed)) + assertEquals(1, extensionsSent) + } + + @Test + fun `a poll after the window sends another extension`() = runTest { + throttle.extendIfDue(userSwarm, succeed) + time += ConfigTtlExtensionThrottle.COOLDOWN + + assertTrue(throttle.extendIfDue(userSwarm, succeed)) + assertEquals(2, extensionsSent) + } + + @Test + fun `a failed extension does not start the cooldown`() = runTest { + assertFailsWith { throttle.extendIfDue(userSwarm, fail) } + time += 1.seconds + + assertTrue(throttle.extendIfDue(userSwarm, succeed)) + assertEquals(2, extensionsSent) + + // Positive control: the success above did start the cooldown, so the retry is observable as a + // retry rather than as a throttle that never engages. + assertFalse(throttle.extendIfDue(userSwarm, succeed)) + assertEquals(2, extensionsSent) + } + + @Test + fun `repeated failures keep retrying on every poll`() = runTest { + repeat(3) { + assertFailsWith { throttle.extendIfDue(userSwarm, fail) } + time += 10.seconds + } + + assertEquals(3, extensionsSent) + } + + @Test + fun `a failure after the window re-opens it does not extend the old cooldown`() = runTest { + throttle.extendIfDue(userSwarm, succeed) + time += ConfigTtlExtensionThrottle.COOLDOWN + + assertFailsWith { throttle.extendIfDue(userSwarm, fail) } + time += 1.seconds + + assertTrue(throttle.extendIfDue(userSwarm, succeed)) + assertEquals(3, extensionsSent) + } + + @Test + fun `a cancelled extension does not start the cooldown`() = runTest { + assertFailsWith { + throttle.extendIfDue(userSwarm) { throw CancellationException("poller stopped") } + } + + assertTrue(throttle.extendIfDue(userSwarm, succeed)) + assertEquals(1, extensionsSent) + } + + @Test + fun `each swarm has its own cooldown`() = runTest { + throttle.extendIfDue(groupSwarm, succeed) + + assertTrue(throttle.extendIfDue(otherGroupSwarm, succeed)) + assertTrue(throttle.extendIfDue(userSwarm, succeed)) + assertFalse(throttle.extendIfDue(groupSwarm, succeed)) + assertEquals(3, extensionsSent) + } +} diff --git a/app/src/test/java/org/thoughtcrime/securesms/groups/GroupPollerCursorResetTest.kt b/app/src/test/java/org/thoughtcrime/securesms/groups/GroupPollerCursorResetTest.kt index 4e0086f4f3..6d21dc3f09 100644 --- a/app/src/test/java/org/thoughtcrime/securesms/groups/GroupPollerCursorResetTest.kt +++ b/app/src/test/java/org/thoughtcrime/securesms/groups/GroupPollerCursorResetTest.kt @@ -18,6 +18,7 @@ import org.junit.After import org.junit.Before import org.junit.Rule import org.junit.Test +import org.session.libsession.messaging.sending_receiving.pollers.ConfigTtlExtensionThrottle import org.session.libsession.snode.SwarmAuth import org.session.libsession.snode.model.RetrieveMessageResponse import org.session.libsession.utilities.ConfigFactoryProtocol @@ -130,6 +131,7 @@ class GroupPollerCursorResetTest { alterTtlApiApiFactory = mockk(relaxed = true), swarmApiExecutor = swarmApiExecutor, swarmSnodeSelector = snodeSelector, + configTtlExtensionThrottle = ConfigTtlExtensionThrottle(), networkConnectivity = networkConnectivity, appVisibilityManager = appVisibilityManager, ).also { it.namespaces = Namespaces }