From 958a460ea97b1f56563d975d9cf4b50d45311550 Mon Sep 17 00:00:00 2001 From: Vitor Pamplona Date: Thu, 26 Mar 2026 11:10:45 -0400 Subject: [PATCH] Adds a try/catch to make sure this procedure unsubscribes from the req --- .../service/broadcast/BroadcastTracker.kt | 71 ++++++++++--------- 1 file changed, 37 insertions(+), 34 deletions(-) diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/service/broadcast/BroadcastTracker.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/service/broadcast/BroadcastTracker.kt index 1c3d62232c..c65d10d6fd 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/service/broadcast/BroadcastTracker.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/service/broadcast/BroadcastTracker.kt @@ -142,55 +142,58 @@ class BroadcastTracker { } } - client.subscribe(subscription) + try { + client.subscribe(subscription) - val finalBroadcast = - coroutineScope { - val resultCollector = - async { - val receivedRelays = mutableSetOf() - var currentBroadcast = broadcast + val finalBroadcast = + coroutineScope { + val resultCollector = + async { + val receivedRelays = mutableSetOf() + var currentBroadcast = broadcast - withTimeoutOrNull(TIMEOUT_SECONDS * 1000) { - while (receivedRelays.size < relays.size) { - val response = resultChannel.receive() + withTimeoutOrNull(TIMEOUT_SECONDS * 1000) { + while (receivedRelays.size < relays.size) { + val response = resultChannel.receive() - // Skip if already received (don't override success) - if (response.relay in receivedRelays) continue + // Skip if already received (don't override success) + if (response.relay in receivedRelays) continue - receivedRelays.add(response.relay) - currentBroadcast = currentBroadcast.withResult(response.relay, response.result) + receivedRelays.add(response.relay) + currentBroadcast = currentBroadcast.withResult(response.relay, response.result) - // Update active broadcasts with new progress - _activeBroadcasts.update { list -> - list.map { if (it.id == trackingId) currentBroadcast else it }.toImmutableList() + // Update active broadcasts with new progress + _activeBroadcasts.update { list -> + list.map { if (it.id == trackingId) currentBroadcast else it }.toImmutableList() + } } } + + // Mark remaining relays as timeout + relays.filter { it !in receivedRelays }.forEach { relay -> + currentBroadcast = currentBroadcast.withResult(relay, RelayResult.Timeout) + } + + currentBroadcast } - // Mark remaining relays as timeout - relays.filter { it !in receivedRelays }.forEach { relay -> - currentBroadcast = currentBroadcast.withResult(relay, RelayResult.Timeout) - } + // Send after setting up listener + client.send(event, relays) - currentBroadcast - } + resultCollector.await() + } - // Send after setting up listener - client.send(event, relays) + resultChannel.close() - resultCollector.await() + // Remove from active, emit to completed + _activeBroadcasts.update { list -> + list.map { if (it.id == trackingId) finalBroadcast else it }.toImmutableList() } - client.unsubscribe(subscription) - resultChannel.close() - - // Remove from active, emit to completed - _activeBroadcasts.update { list -> - list.map { if (it.id == trackingId) finalBroadcast else it }.toImmutableList() + Log.d(TAG, "Broadcast $trackingId complete: ${finalBroadcast.successCount}/${finalBroadcast.totalRelays} success") + } finally { + client.unsubscribe(subscription) } - - Log.d(TAG, "Broadcast $trackingId complete: ${finalBroadcast.successCount}/${finalBroadcast.totalRelays} success") } /**