mirror of
https://github.com/greenart7c3/Amber.git
synced 2026-10-05 19:08:23 +00:00
Adapt to Quartz 1.14.0 suspend listener changes
Quartz 1.14.0 made WebSocketListener.onMessage and RelayConnectionListener.onIncomingMessage suspend functions, which broke compilation. Bridge the suspend boundary in OkHttpWebSocket the same way Amethyst and Quartz's own BasicOkHttpWebSocket do: pump incoming frames through an unbounded Channel consumed by a Dispatchers.IO coroutine, reusing BasicOkHttpWebSocket's exceptionHandler. Kept Amber's needsReconnect() proxy/timeout diffing that the Quartz builtin lacks. Mark all onIncomingMessage overrides suspend and wrap direct test invocations in runBlocking. Also bumps AGP to 9.3.2.
This commit is contained in:
@@ -24,6 +24,13 @@ import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.WebSocket
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.WebSocketListener
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.WebsocketBuilder
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.BasicOkHttpWebSocket.Companion.exceptionHandler
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.cancel
|
||||
import kotlinx.coroutines.channels.Channel
|
||||
import kotlinx.coroutines.channels.trySendBlocking
|
||||
import kotlinx.coroutines.launch
|
||||
import okhttp3.OkHttpClient
|
||||
import okhttp3.Request
|
||||
import okhttp3.Response
|
||||
@@ -33,7 +40,6 @@ class OkHttpWebSocket(
|
||||
val httpClient: (url: NormalizedRelayUrl) -> OkHttpClient,
|
||||
val out: WebSocketListener,
|
||||
) : WebSocket {
|
||||
private val listener = OkHttpWebsocketListener()
|
||||
private var usingOkHttp: OkHttpClient? = null
|
||||
private var socket: okhttp3.WebSocket? = null
|
||||
|
||||
@@ -62,10 +68,28 @@ class OkHttpWebSocket(
|
||||
|
||||
override fun connect() {
|
||||
usingOkHttp = httpClient(url)
|
||||
socket = usingOkHttp?.newWebSocket(buildRequest(), listener)
|
||||
socket = usingOkHttp?.newWebSocket(buildRequest(), OkHttpWebsocketListener(out))
|
||||
}
|
||||
|
||||
inner class OkHttpWebsocketListener : okhttp3.WebSocketListener() {
|
||||
inner class OkHttpWebsocketListener(
|
||||
val out: WebSocketListener,
|
||||
) : okhttp3.WebSocketListener() {
|
||||
val scope = CoroutineScope(Dispatchers.IO + exceptionHandler)
|
||||
|
||||
// UNLIMITED on purpose — do NOT bound this channel. The app holds
|
||||
// many relay connections; a bounded buffer under a slow consumer
|
||||
// would block OkHttp reader threads and park the backlog on the
|
||||
// relay's outbound buffers — infrastructure that isn't ours. Drain
|
||||
// the remote as fast as it can send; consumer speed is handled
|
||||
// downstream.
|
||||
val incomingMessages: Channel<String> = Channel(Channel.UNLIMITED)
|
||||
val job = // Launch a coroutine to process messages from the channel.
|
||||
scope.launch {
|
||||
for (message in incomingMessages) {
|
||||
out.onMessage(message)
|
||||
}
|
||||
}
|
||||
|
||||
override fun onOpen(
|
||||
webSocket: okhttp3.WebSocket,
|
||||
response: Response,
|
||||
@@ -77,13 +101,23 @@ class OkHttpWebSocket(
|
||||
override fun onMessage(
|
||||
webSocket: okhttp3.WebSocket,
|
||||
text: String,
|
||||
) = out.onMessage(text)
|
||||
) {
|
||||
// Asynchronously send the received message to the channel.
|
||||
// `trySendBlocking` is used here for simplicity within the callback,
|
||||
// but it's important to understand potential thread blocking if the buffer is full.
|
||||
incomingMessages.trySendBlocking(text)
|
||||
}
|
||||
|
||||
override fun onClosed(
|
||||
webSocket: okhttp3.WebSocket,
|
||||
code: Int,
|
||||
reason: String,
|
||||
) {
|
||||
// Close the channel when the WebSocket connection is closed.
|
||||
incomingMessages.close()
|
||||
job.cancel()
|
||||
scope.cancel()
|
||||
|
||||
socket = null
|
||||
out.onClosed(code, reason)
|
||||
}
|
||||
@@ -93,6 +127,11 @@ class OkHttpWebSocket(
|
||||
t: Throwable,
|
||||
response: Response?,
|
||||
) {
|
||||
// Close the channel on failure, and propagate the error.
|
||||
incomingMessages.close()
|
||||
job.cancel()
|
||||
scope.cancel()
|
||||
|
||||
socket = null
|
||||
out.onFailure(t, response?.code, response?.message)
|
||||
}
|
||||
|
||||
@@ -160,7 +160,7 @@ class NostrClientLoggerListener(
|
||||
super.onSent(relay, cmdStr, cmd, success)
|
||||
}
|
||||
|
||||
override fun onIncomingMessage(relay: IRelayClient, msgStr: String, msg: Message) {
|
||||
override suspend fun onIncomingMessage(relay: IRelayClient, msgStr: String, msg: Message) {
|
||||
// Defense-in-depth (GHSA-8844-q5vh-9j8f, L2): log only the message
|
||||
// type and wire size, not the raw frame (which may carry NIP-46
|
||||
// envelopes, DM/gift-wrap ciphertexts and event content).
|
||||
|
||||
@@ -164,7 +164,7 @@ object ApplicationBackup {
|
||||
val latest = mutableMapOf<Long, Event>()
|
||||
|
||||
val listener = object : RelayConnectionListener {
|
||||
override fun onIncomingMessage(relay: IRelayClient, msgStr: String, msg: Message) {
|
||||
override suspend fun onIncomingMessage(relay: IRelayClient, msgStr: String, msg: Message) {
|
||||
if (msg is EventMessage && msg.subId == subId && msg.event.kind == INBOX_KIND && msg.event.pubKey == account.hexKey) {
|
||||
if (msg.event.verify()) {
|
||||
latest[msg.event.createdAt] = msg.event
|
||||
@@ -250,7 +250,7 @@ object ApplicationBackup {
|
||||
val received = mutableMapOf<Long, Event>()
|
||||
|
||||
val listener = object : RelayConnectionListener {
|
||||
override fun onIncomingMessage(relay: IRelayClient, msgStr: String, msg: Message) {
|
||||
override suspend fun onIncomingMessage(relay: IRelayClient, msgStr: String, msg: Message) {
|
||||
if (msg is EventMessage &&
|
||||
msg.subId == subId &&
|
||||
msg.event.kind == BACKUP_KIND &&
|
||||
|
||||
@@ -54,7 +54,7 @@ class NotificationSubscription(
|
||||
client.addConnectionListener(this)
|
||||
}
|
||||
|
||||
override fun onIncomingMessage(relay: IRelayClient, msgStr: String, msg: Message) {
|
||||
override suspend fun onIncomingMessage(relay: IRelayClient, msgStr: String, msg: Message) {
|
||||
if (msg is EventMessage) {
|
||||
if (subIds.containsValue(msg.subId)) {
|
||||
Amber.instance.applicationIOScope.launch {
|
||||
|
||||
@@ -79,7 +79,7 @@ class ProfileSubscription(
|
||||
client.addConnectionListener(this)
|
||||
}
|
||||
|
||||
override fun onIncomingMessage(relay: IRelayClient, msgStr: String, msg: Message) {
|
||||
override suspend fun onIncomingMessage(relay: IRelayClient, msgStr: String, msg: Message) {
|
||||
if (msg is EoseMessage) {
|
||||
val subId = msg.subId
|
||||
val relays = relaysPerSubId[subId]
|
||||
|
||||
@@ -91,7 +91,7 @@ class ZapstoreUpdater(
|
||||
}
|
||||
}
|
||||
|
||||
override fun onIncomingMessage(relay: IRelayClient, msgStr: String, msg: Message) {
|
||||
override suspend fun onIncomingMessage(relay: IRelayClient, msgStr: String, msg: Message) {
|
||||
when (msg) {
|
||||
is EventMessage if (msg.subId == releaseSubId || msg.subId == fileSubId) -> messages.trySend(msg)
|
||||
is EoseMessage if (msg.subId == releaseSubId || msg.subId == fileSubId) -> messages.trySend(msg)
|
||||
|
||||
@@ -434,7 +434,7 @@ fun onAddRelay(
|
||||
super.onCannotConnect(relay, errorMessage)
|
||||
}
|
||||
|
||||
override fun onIncomingMessage(relay: IRelayClient, msgStr: String, msg: Message) {
|
||||
override suspend fun onIncomingMessage(relay: IRelayClient, msgStr: String, msg: Message) {
|
||||
if (msg is EventMessage) {
|
||||
if (ncSub == msg.subId && msg.event.kind == NostrConnectEvent.KIND && msg.event.id == signedEvent.id) {
|
||||
filterResult = true
|
||||
|
||||
+1
-1
@@ -156,7 +156,7 @@ class NotificationSubscriptionTest {
|
||||
threads += Thread {
|
||||
repeat(300) { i ->
|
||||
try {
|
||||
subscription.onIncomingMessage(relay, "", EventMessage("unknown-sub-$i", event))
|
||||
runBlocking { subscription.onIncomingMessage(relay, "", EventMessage("unknown-sub-$i", event)) }
|
||||
} catch (e: Throwable) {
|
||||
errors.add(e)
|
||||
}
|
||||
|
||||
@@ -185,8 +185,10 @@ class ProfileSubscriptionTest {
|
||||
repeat(200) {
|
||||
try {
|
||||
val subId = sentSubIds.toList().randomOrNull() ?: return@repeat
|
||||
subscription.onIncomingMessage(relay, "", EoseMessage(subId))
|
||||
subscription.onIncomingMessage(relay, "", EventMessage(subId, event))
|
||||
runBlocking {
|
||||
subscription.onIncomingMessage(relay, "", EoseMessage(subId))
|
||||
subscription.onIncomingMessage(relay, "", EventMessage(subId, event))
|
||||
}
|
||||
} catch (e: Throwable) {
|
||||
errors.add(e)
|
||||
}
|
||||
|
||||
@@ -11,7 +11,7 @@ junitVersion = "1.3.0"
|
||||
lifecycle_version = "2.11.0"
|
||||
material3 = "1.4.0"
|
||||
nav_version = "2.9.8"
|
||||
quartz = "1.13.1"
|
||||
quartz = "1.14.0"
|
||||
compose_ui = "1.12.0"
|
||||
roomKtx = "2.8.4"
|
||||
securityCryptoKtx = "1.1.0"
|
||||
@@ -19,7 +19,7 @@ zxingAndroidEmbedded = "4.3.0"
|
||||
okhttp = "5.5.0"
|
||||
kotlin = "2.4.10"
|
||||
workRuntimeKtx = "2.11.2"
|
||||
agp = "9.3.1"
|
||||
agp = "9.3.2"
|
||||
ktlint = "14.2.0"
|
||||
ksp = "2.3.11"
|
||||
coil = "3.5.0"
|
||||
|
||||
Reference in New Issue
Block a user