Makes sure subscriptions are closed on the NostrClient utility methods
This commit is contained in:
+9
-6
@@ -65,15 +65,18 @@ suspend fun INostrClient.queryCountSuspend(
|
|||||||
|
|
||||||
subscribe(listener)
|
subscribe(listener)
|
||||||
|
|
||||||
queryCount(subId = subId, filters = mapOf(relay to listOf(filter)))
|
|
||||||
|
|
||||||
val result =
|
val result =
|
||||||
withTimeoutOrNull(timeoutMs) {
|
try {
|
||||||
resultChannel.receive()
|
queryCount(subId = subId, filters = mapOf(relay to listOf(filter)))
|
||||||
|
|
||||||
|
withTimeoutOrNull(timeoutMs) {
|
||||||
|
resultChannel.receive()
|
||||||
|
}
|
||||||
|
} finally {
|
||||||
|
close(subId)
|
||||||
|
unsubscribe(listener)
|
||||||
}
|
}
|
||||||
|
|
||||||
close(subId)
|
|
||||||
unsubscribe(listener)
|
|
||||||
resultChannel.close()
|
resultChannel.close()
|
||||||
|
|
||||||
return result
|
return result
|
||||||
|
|||||||
+27
-24
@@ -98,38 +98,41 @@ suspend fun INostrClient.sendAndWaitForResponseDetailed(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
subscribe(subscription)
|
val receivedResults =
|
||||||
|
try {
|
||||||
|
subscribe(subscription)
|
||||||
|
|
||||||
// subscribe before sending the result.
|
// subscribe before sending the result.
|
||||||
val resultSubscription =
|
val resultSubscription =
|
||||||
coroutineScope {
|
coroutineScope {
|
||||||
val result =
|
val result =
|
||||||
async {
|
async {
|
||||||
val receivedResults = mutableMapOf<NormalizedRelayUrl, Boolean>()
|
val receivedResults = mutableMapOf<NormalizedRelayUrl, Boolean>()
|
||||||
// The withTimeout block will cancel the coroutine if the loop takes too long
|
// The withTimeout block will cancel the coroutine if the loop takes too long
|
||||||
withTimeoutOrNull(timeoutInSeconds * 1000) {
|
withTimeoutOrNull(timeoutInSeconds * 1000) {
|
||||||
while (receivedResults.size < relayList.size) {
|
while (receivedResults.size < relayList.size) {
|
||||||
val result = resultChannel.receive()
|
val result = resultChannel.receive()
|
||||||
|
|
||||||
val currentResult = receivedResults[result.relay]
|
val currentResult = receivedResults[result.relay]
|
||||||
// do not override a successful result.
|
// do not override a successful result.
|
||||||
if (currentResult == null || !currentResult) {
|
if (currentResult == null || !currentResult) {
|
||||||
receivedResults[result.relay] = result.success
|
receivedResults[result.relay] = result.success
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
receivedResults
|
||||||
}
|
}
|
||||||
}
|
|
||||||
receivedResults
|
send(event, relayList)
|
||||||
|
|
||||||
|
result
|
||||||
}
|
}
|
||||||
|
|
||||||
send(event, relayList)
|
resultSubscription.await()
|
||||||
|
} finally {
|
||||||
result
|
unsubscribe(subscription)
|
||||||
}
|
}
|
||||||
|
|
||||||
val receivedResults = resultSubscription.await()
|
|
||||||
|
|
||||||
unsubscribe(subscription)
|
|
||||||
|
|
||||||
// Clean up the channel
|
// Clean up the channel
|
||||||
resultChannel.close()
|
resultChannel.close()
|
||||||
|
|
||||||
|
|||||||
+8
-6
@@ -104,14 +104,16 @@ suspend fun INostrClient.downloadFirstEvent(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
openReqSubscription(subscriptionId, filters, listener)
|
|
||||||
|
|
||||||
val result =
|
val result =
|
||||||
withTimeoutOrNull(30000) {
|
try {
|
||||||
resultChannel.receive()
|
openReqSubscription(subscriptionId, filters, listener)
|
||||||
}
|
|
||||||
|
|
||||||
close(subscriptionId)
|
withTimeoutOrNull(30000) {
|
||||||
|
resultChannel.receive()
|
||||||
|
}
|
||||||
|
} finally {
|
||||||
|
close(subscriptionId)
|
||||||
|
}
|
||||||
|
|
||||||
resultChannel.close()
|
resultChannel.close()
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user