Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -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<String, ComparableTimeMark>()

/**
* 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
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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<Unit>(
debugLabel = "MainPoller",
Expand Down Expand Up @@ -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

Expand Down
24 changes: 14 additions & 10 deletions app/src/main/java/org/thoughtcrime/securesms/groups/GroupPoller.kt
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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<GroupPoller.GroupPollResult>(
Expand Down Expand Up @@ -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,
)
)
)
)
}
}
}

Expand Down
Original file line number Diff line number Diff line change
@@ -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<RuntimeException> { 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<RuntimeException> { 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<RuntimeException> { 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<CancellationException> {
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)
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -130,6 +131,7 @@ class GroupPollerCursorResetTest {
alterTtlApiApiFactory = mockk(relaxed = true),
swarmApiExecutor = swarmApiExecutor,
swarmSnodeSelector = snodeSelector,
configTtlExtensionThrottle = ConfigTtlExtensionThrottle(),
networkConnectivity = networkConnectivity,
appVisibilityManager = appVisibilityManager,
).also { it.namespaces = Namespaces }
Expand Down
Loading