diff --git a/client/src/commonMain/kotlin/pw/binom/agentik/client/ReconnectingOutbox.kt b/client/src/commonMain/kotlin/pw/binom/agentik/client/ReconnectingOutbox.kt index bb88273..f9da95d 100644 --- a/client/src/commonMain/kotlin/pw/binom/agentik/client/ReconnectingOutbox.kt +++ b/client/src/commonMain/kotlin/pw/binom/agentik/client/ReconnectingOutbox.kt @@ -13,8 +13,10 @@ import kotlinx.coroutines.launch import pw.binom.agentik.outbox.CommonEvent import pw.binom.agentik.outbox.OutboxStore import kotlin.concurrent.atomics.AtomicBoolean +import kotlin.concurrent.atomics.AtomicReference import kotlin.concurrent.atomics.ExperimentalAtomicApi import kotlin.math.min +import kotlin.math.pow import kotlin.random.Random import kotlin.time.Duration import kotlin.time.Duration.Companion.seconds @@ -151,8 +153,8 @@ class ReconnectingOutbox( private val started = AtomicBoolean(false) private var job: Job? = null - @Volatile - private var lastSeen: Instant? = null + @OptIn(ExperimentalAtomicApi::class) + private val lastSeen: AtomicReference = AtomicReference(null) /** * Live-события из [outbox] с авто-reconnect. [after] — начальный курсор; @@ -183,10 +185,11 @@ class ReconnectingOutbox( @OptIn(ExperimentalAtomicApi::class) private fun ensureStarted(initialCursor: Instant?) { if (!started.compareAndSet(false, true)) return - lastSeen = initialCursor + lastSeen.store(initialCursor) job = scope.launch { runLoop() } } + @OptIn(ExperimentalAtomicApi::class) private suspend fun runLoop() { var attempt = 0 var connected = false @@ -194,12 +197,12 @@ class ReconnectingOutbox( attempt++ _status.emit(ConnectionStatus.Connecting(attempt)) val error: Throwable? = try { - outbox.events(after = lastSeen).collect { event -> - lastSeen = event.date + outbox.events(after = lastSeen.load()).collect { event -> + lastSeen.store(event.date) _events.emit(event) if (!connected) { connected = true - _status.emit(ConnectionStatus.Connected(lastSeen!!)) + _status.emit(ConnectionStatus.Connected(event.date)) } } null @@ -224,7 +227,7 @@ class ReconnectingOutbox( private fun computeBackoff(attempt: Int): Duration { // attempt 1 → initial, 2 → initial * m, 3 → initial * m^2, ... val base = (policy.initial.inWholeMilliseconds.toDouble() * - Math.pow(policy.multiplier, (attempt - 1).toDouble())) + policy.multiplier.pow((attempt - 1).toDouble())) .toLong() val capped = min(base, policy.max.inWholeMilliseconds) val jitterMs = (capped * policy.jitter * random.nextDouble()).toLong()