mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-05 11:18:24 +00:00
perf: replace synchronous OkHttp execute() with executeAsync()
Every coroutine dispatched to Dispatchers.IO is stamped BlockingContext at
dispatch time, so its worker releases its CPU permit and the shared kotlinx
scheduler grows past ncpu. A blocking execute() holds one of those threads for
the whole request; executeAsync() suspends until the response headers arrive
and releases the thread across the network round-trip.
Converts 27 call sites across quartz, commons, amethyst, nestsClient,
desktopApp and cli. executeAsync() was already the house pattern (63 existing
uses); these were the stragglers.
The withContext(Dispatchers.IO) wrappers are kept on purpose: executeAsync()
only suspends until headers, and reading the body (string()/bytes()) is still
a blocking read. Moving those onto Dispatchers.Default would hold its
core-sized CPU permits and starve the pool.
Four private helpers become suspend (decodeGifFrames, fetchFromNetwork,
downloadFirstChunk, extractFirstFrame). Each had exactly one caller, already
inside a withContext(Dispatchers.IO) in a suspend function, so nothing
propagates further.
desktopApp gains the okhttp-coroutines dependency (Apache-2.0, same version as
the okhttp it already ships).
Not converted: NipCommand.fetchText and RelayCommands.info in the CLI. amy is a
short-lived process whose main is runBlocking { dispatch(argv) } — there is no
long-lived dispatcher pool or UI thread to protect there, so blocking is by
design and making them suspend would only add risk.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 5
parent
c94bb81df5
commit
8020102364
@@ -30,6 +30,7 @@ import okhttp3.MediaType.Companion.toMediaType
|
||||
import okhttp3.OkHttpClient
|
||||
import okhttp3.Request
|
||||
import okhttp3.RequestBody.Companion.toRequestBody
|
||||
import okhttp3.coroutines.executeAsync
|
||||
|
||||
/**
|
||||
* Mints a `block/buzz` workspace invite link by POSTing to the relay's **Buzz-specific**
|
||||
@@ -86,7 +87,7 @@ object BuzzInviteMinter {
|
||||
.post(bodyBytes.toRequestBody("application/json".toMediaType()))
|
||||
.build()
|
||||
|
||||
okHttpClient(url).newCall(request).execute().use { response ->
|
||||
okHttpClient(url).newCall(request).executeAsync().use { response ->
|
||||
val payload = response.body.string()
|
||||
val tree = runCatching { json.readTree(payload) }.getOrNull()
|
||||
|
||||
|
||||
+2
-1
@@ -30,6 +30,7 @@ import com.vitorpamplona.quartz.utils.Log
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.withContext
|
||||
import okhttp3.Request
|
||||
import okhttp3.coroutines.executeAsync
|
||||
import java.io.File
|
||||
import java.util.concurrent.ConcurrentHashMap
|
||||
|
||||
@@ -70,7 +71,7 @@ fun rememberConcordImageModel(
|
||||
if (!cacheFile.exists()) {
|
||||
val client = accountViewModel.httpClientBuilder.okHttpClientForImage(pointer.url)
|
||||
val ciphertext =
|
||||
client.newCall(Request.Builder().url(pointer.url).build()).execute().use { resp ->
|
||||
client.newCall(Request.Builder().url(pointer.url).build()).executeAsync().use { resp ->
|
||||
if (!resp.isSuccessful) return@runCatching null
|
||||
resp.body.bytes()
|
||||
}
|
||||
|
||||
+3
-2
@@ -29,6 +29,7 @@ import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.withContext
|
||||
import okhttp3.OkHttpClient
|
||||
import okhttp3.Request
|
||||
import okhttp3.coroutines.executeAsync
|
||||
import java.io.File
|
||||
import java.security.MessageDigest
|
||||
|
||||
@@ -63,7 +64,7 @@ internal object RoomFontLoader {
|
||||
}.getOrNull()
|
||||
}
|
||||
|
||||
private fun ensureCached(
|
||||
private suspend fun ensureCached(
|
||||
url: String,
|
||||
context: Context,
|
||||
clientFor: (String) -> OkHttpClient,
|
||||
@@ -74,7 +75,7 @@ internal object RoomFontLoader {
|
||||
|
||||
val client = clientFor(url)
|
||||
val request = Request.Builder().url(url).build()
|
||||
client.newCall(request).execute().use { resp ->
|
||||
client.newCall(request).executeAsync().use { resp ->
|
||||
if (!resp.isSuccessful) return null
|
||||
file.outputStream().use { out -> resp.body.byteStream().copyTo(out) }
|
||||
}
|
||||
|
||||
@@ -74,6 +74,7 @@ import kotlinx.coroutines.withTimeoutOrNull
|
||||
import okhttp3.Dispatcher
|
||||
import okhttp3.OkHttpClient
|
||||
import okhttp3.Request
|
||||
import okhttp3.coroutines.executeAsync
|
||||
import java.lang.management.ManagementFactory
|
||||
import java.util.concurrent.TimeUnit
|
||||
|
||||
@@ -304,7 +305,7 @@ class Context(
|
||||
.url(relay.toHttp())
|
||||
.header("Accept", "application/nostr+json")
|
||||
.build()
|
||||
okhttp.newCall(request).execute().use { resp ->
|
||||
okhttp.newCall(request).executeAsync().use { resp ->
|
||||
resp.body.string().let { Nip11RelayInformation.fromJson(it) }
|
||||
}
|
||||
}.getOrNull()
|
||||
|
||||
@@ -57,6 +57,7 @@ import okhttp3.MediaType.Companion.toMediaType
|
||||
import okhttp3.OkHttpClient
|
||||
import okhttp3.Request
|
||||
import okhttp3.RequestBody.Companion.toRequestBody
|
||||
import okhttp3.coroutines.executeAsync
|
||||
|
||||
/**
|
||||
* `amy buzz …` — first-class access to the `block/buzz` workspace protocol, driving the
|
||||
@@ -372,7 +373,7 @@ object BuzzCommands {
|
||||
.url(url)
|
||||
.get()
|
||||
.build(),
|
||||
).execute()
|
||||
).executeAsync()
|
||||
.use { it.code to it.body.string() }
|
||||
}
|
||||
|
||||
@@ -385,7 +386,7 @@ object BuzzCommands {
|
||||
withContext(Dispatchers.IO) {
|
||||
val builder = Request.Builder().url(url).post(body.toRequestBody(jsonMedia))
|
||||
if (auth != null) builder.header("Authorization", auth)
|
||||
http.newCall(builder.build()).execute().use { it.code to it.body.string() }
|
||||
http.newCall(builder.build()).executeAsync().use { it.code to it.body.string() }
|
||||
}
|
||||
|
||||
/** `buzz post RELAY GID <text>` → publishes a kind-40002 stream message with an `h` tag. */
|
||||
|
||||
+2
-1
@@ -24,6 +24,7 @@ import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.withContext
|
||||
import okhttp3.OkHttpClient
|
||||
import okhttp3.Request
|
||||
import okhttp3.coroutines.executeAsync
|
||||
|
||||
/**
|
||||
* Fetches and parses the live georelays CSV so the [GeoRelayDirectory] can route
|
||||
@@ -44,7 +45,7 @@ class GeoRelayCsvLoader(
|
||||
.url(url)
|
||||
.get()
|
||||
.build()
|
||||
okHttpClient(url).newCall(request).execute().use { response ->
|
||||
okHttpClient(url).newCall(request).executeAsync().use { response ->
|
||||
if (response.isSuccessful) {
|
||||
GeoRelayDirectory.parseCsv(response.body.string())
|
||||
} else {
|
||||
|
||||
+3
-2
@@ -30,6 +30,7 @@ import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.withContext
|
||||
import okhttp3.OkHttpClient
|
||||
import okhttp3.Request
|
||||
import okhttp3.coroutines.executeAsync
|
||||
import java.math.BigDecimal
|
||||
import java.math.RoundingMode
|
||||
import java.net.URLEncoder
|
||||
@@ -191,7 +192,7 @@ class LightningAddressResolver(
|
||||
withContext(Dispatchers.IO) {
|
||||
try {
|
||||
val request = Request.Builder().url(url).build()
|
||||
httpClient.newCall(request).execute().use { response ->
|
||||
httpClient.newCall(request).executeAsync().use { response ->
|
||||
if (response.isSuccessful) {
|
||||
response.body.string()
|
||||
} else {
|
||||
@@ -222,7 +223,7 @@ class LightningAddressResolver(
|
||||
}
|
||||
|
||||
val request = Request.Builder().url(url).build()
|
||||
httpClient.newCall(request).execute().use { response ->
|
||||
httpClient.newCall(request).executeAsync().use { response ->
|
||||
// Return body even on error — caller extracts "reason" or "message" from JSON
|
||||
response.body.string()
|
||||
}
|
||||
|
||||
+2
-1
@@ -28,6 +28,7 @@ import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.withContext
|
||||
import okhttp3.OkHttpClient
|
||||
import okhttp3.Request
|
||||
import okhttp3.coroutines.executeAsync
|
||||
import kotlin.coroutines.cancellation.CancellationException
|
||||
|
||||
/**
|
||||
@@ -58,7 +59,7 @@ class OkHttpLnurlEndpointResolver(
|
||||
try {
|
||||
val client = okHttpClient(url)
|
||||
val request = Request.Builder().url(url).build()
|
||||
client.newCall(request).execute().use { response ->
|
||||
client.newCall(request).executeAsync().use { response ->
|
||||
if (!response.isSuccessful) return@use null
|
||||
val body = response.body.string()
|
||||
val root = mapper.readTree(body) ?: return@use null
|
||||
|
||||
+9
-8
@@ -36,6 +36,7 @@ import okhttp3.Request
|
||||
import okhttp3.RequestBody
|
||||
import okhttp3.RequestBody.Companion.toRequestBody
|
||||
import okhttp3.Response
|
||||
import okhttp3.coroutines.executeAsync
|
||||
import okio.BufferedSink
|
||||
import okio.source
|
||||
import java.io.File
|
||||
@@ -149,7 +150,7 @@ open class BlossomClient(
|
||||
paymentProof?.headers()?.forEach { (name, value) -> addHeader(name, value) }
|
||||
}.put(body)
|
||||
.build()
|
||||
okHttpClient.newCall(request).execute().use { response ->
|
||||
okHttpClient.newCall(request).executeAsync().use { response ->
|
||||
if (response.code in MIRROR_UNSUPPORTED_CODES) {
|
||||
throw BlossomMirrorUnsupportedException(serverBaseUrl, response.code)
|
||||
}
|
||||
@@ -208,7 +209,7 @@ open class BlossomClient(
|
||||
.apply { authHeader?.let { addHeader("Authorization", it) } }
|
||||
.get()
|
||||
.build()
|
||||
okHttpClient.newCall(request).execute().use { response ->
|
||||
okHttpClient.newCall(request).executeAsync().use { response ->
|
||||
check402(response, serverBaseUrl)
|
||||
if (!response.isSuccessful) {
|
||||
val reason = response.headers[BlossomServerUrl.REASON_HEADER] ?: response.code.toString()
|
||||
@@ -237,7 +238,7 @@ open class BlossomClient(
|
||||
.apply { authHeader?.let { addHeader("Authorization", it) } }
|
||||
.delete()
|
||||
.build()
|
||||
okHttpClient.newCall(request).execute().use { it.isSuccessful }
|
||||
okHttpClient.newCall(request).executeAsync().use { it.isSuccessful }
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -256,7 +257,7 @@ open class BlossomClient(
|
||||
.head()
|
||||
.build()
|
||||
try {
|
||||
okHttpClient.newCall(request).execute().use { it.isSuccessful }
|
||||
okHttpClient.newCall(request).executeAsync().use { it.isSuccessful }
|
||||
} catch (e: CancellationException) {
|
||||
throw e
|
||||
} catch (_: Exception) {
|
||||
@@ -289,7 +290,7 @@ open class BlossomClient(
|
||||
.addHeader(BlossomServerUrl.X_CONTENT_TYPE_HEADER, contentType)
|
||||
.apply { authHeader?.let { addHeader("Authorization", it) } }
|
||||
.build()
|
||||
okHttpClient.newCall(request).execute().use { response ->
|
||||
okHttpClient.newCall(request).executeAsync().use { response ->
|
||||
BlossomPreflightResult(
|
||||
accepted = response.isSuccessful,
|
||||
status = response.code,
|
||||
@@ -313,7 +314,7 @@ open class BlossomClient(
|
||||
.url(BlossomServerUrl.report(serverBaseUrl))
|
||||
.put(reportEventJson.toRequestBody("application/json".toMediaType()))
|
||||
.build()
|
||||
okHttpClient.newCall(request).execute().use { it.isSuccessful }
|
||||
okHttpClient.newCall(request).executeAsync().use { it.isSuccessful }
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -335,7 +336,7 @@ open class BlossomClient(
|
||||
.url(url)
|
||||
.get()
|
||||
.build()
|
||||
okHttpClient.newCall(request).execute().use { response ->
|
||||
okHttpClient.newCall(request).executeAsync().use { response ->
|
||||
if (response.isSuccessful) response.body.bytes() else null
|
||||
}
|
||||
}
|
||||
@@ -368,7 +369,7 @@ open class BlossomClient(
|
||||
.apply { authHeader?.let { addHeader("Authorization", it) } }
|
||||
.put(body)
|
||||
.build()
|
||||
okHttpClient.newCall(request).execute().use { parseDescriptor(it, serverBaseUrl) }
|
||||
okHttpClient.newCall(request).executeAsync().use { parseDescriptor(it, serverBaseUrl) }
|
||||
}
|
||||
|
||||
private fun parseDescriptor(
|
||||
|
||||
@@ -47,6 +47,7 @@ dependencies {
|
||||
|
||||
// Networking
|
||||
implementation(libs.okhttp)
|
||||
implementation(libs.okhttpCoroutines)
|
||||
|
||||
// JSON
|
||||
implementation(libs.jackson.module.kotlin)
|
||||
|
||||
+3
-2
@@ -30,6 +30,7 @@ import kotlinx.coroutines.sync.withLock
|
||||
import kotlinx.coroutines.sync.withPermit
|
||||
import kotlinx.coroutines.withContext
|
||||
import okhttp3.Request
|
||||
import okhttp3.coroutines.executeAsync
|
||||
import java.util.concurrent.ConcurrentHashMap
|
||||
|
||||
/**
|
||||
@@ -67,7 +68,7 @@ class Nip11Fetcher {
|
||||
}
|
||||
}
|
||||
|
||||
private fun fetchFromNetwork(url: NormalizedRelayUrl): Nip11RelayInformation? {
|
||||
private suspend fun fetchFromNetwork(url: NormalizedRelayUrl): Nip11RelayInformation? {
|
||||
// FAIL-CLOSED: use currentClient() not getHttpClient()
|
||||
val client = DesktopHttpClient.currentClient()
|
||||
val httpUrl = url.toHttp()
|
||||
@@ -78,7 +79,7 @@ class Nip11Fetcher {
|
||||
.header("Accept", "application/nostr+json")
|
||||
.build()
|
||||
return try {
|
||||
client.newCall(request).execute().use { response ->
|
||||
client.newCall(request).executeAsync().use { response ->
|
||||
if (response.isSuccessful) {
|
||||
val source = response.body.source()
|
||||
source.request(MAX_RESPONSE_BYTES) // buffer up to limit
|
||||
|
||||
+2
-1
@@ -25,6 +25,7 @@ import com.vitorpamplona.quartz.utils.ciphers.AESGCM
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.withContext
|
||||
import okhttp3.Request
|
||||
import okhttp3.coroutines.executeAsync
|
||||
import java.util.concurrent.ConcurrentHashMap
|
||||
|
||||
/**
|
||||
@@ -51,7 +52,7 @@ object EncryptedMediaService {
|
||||
|
||||
return withContext(Dispatchers.IO) {
|
||||
val request = Request.Builder().url(url).build()
|
||||
val response = httpClient.newCall(request).execute()
|
||||
val response = httpClient.newCall(request).executeAsync()
|
||||
val encryptedBytes =
|
||||
response.use {
|
||||
if (!it.isSuccessful) throw RuntimeException("Download failed: ${it.code}")
|
||||
|
||||
+2
-1
@@ -24,6 +24,7 @@ import com.vitorpamplona.amethyst.desktop.network.DesktopHttpClient
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.withContext
|
||||
import okhttp3.Request
|
||||
import okhttp3.coroutines.executeAsync
|
||||
|
||||
object ServerHealthCheck {
|
||||
enum class ServerStatus {
|
||||
@@ -45,7 +46,7 @@ object ServerHealthCheck {
|
||||
.url(url)
|
||||
.head()
|
||||
.build()
|
||||
val response = DesktopHttpClient.currentClient().newCall(request).execute()
|
||||
val response = DesktopHttpClient.currentClient().newCall(request).executeAsync()
|
||||
response.use {
|
||||
if (it.isSuccessful || it.code == 405) ServerStatus.ONLINE else ServerStatus.OFFLINE
|
||||
}
|
||||
|
||||
+4
-3
@@ -27,6 +27,7 @@ import com.vitorpamplona.amethyst.desktop.network.DesktopHttpClient
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.withContext
|
||||
import okhttp3.Request
|
||||
import okhttp3.coroutines.executeAsync
|
||||
import org.jcodec.api.FrameGrab
|
||||
import org.jcodec.common.io.NIOUtils
|
||||
import org.jcodec.common.model.ColorSpace
|
||||
@@ -122,7 +123,7 @@ object VideoThumbnailCache {
|
||||
}
|
||||
}
|
||||
|
||||
private fun extractFirstFrame(url: String): ImageBitmap? {
|
||||
private suspend fun extractFirstFrame(url: String): ImageBitmap? {
|
||||
// For HLS we skip straight to ffmpeg — JCodec can't read m3u8.
|
||||
val isHls = url.contains(".m3u8", ignoreCase = true) || url.contains("/hls/", ignoreCase = true)
|
||||
|
||||
@@ -163,7 +164,7 @@ object VideoThumbnailCache {
|
||||
*
|
||||
* Cleans up zero-byte cache files on failure so a transient empty response isn't sticky.
|
||||
*/
|
||||
private fun downloadFirstChunk(url: String): Download? {
|
||||
private suspend fun downloadFirstChunk(url: String): Download? {
|
||||
val hash = sha1Hex(url)
|
||||
val cached = File(downloadCacheDir, "$hash.mp4")
|
||||
if (cached.length() > 0L) return Download(cached, persistable = true)
|
||||
@@ -171,7 +172,7 @@ object VideoThumbnailCache {
|
||||
|
||||
var wrote = false
|
||||
var rangeHonored = false
|
||||
DesktopHttpClient.currentClient().newCall(buildRangeRequest(url)).execute().use { resp ->
|
||||
DesktopHttpClient.currentClient().newCall(buildRangeRequest(url)).executeAsync().use { resp ->
|
||||
if (!resp.isSuccessful && resp.code != 206) return null
|
||||
val contentType = resp.header("Content-Type")?.lowercase().orEmpty()
|
||||
if (contentType.startsWith("text/") || "html" in contentType) return null
|
||||
|
||||
+2
-1
@@ -82,6 +82,7 @@ import kotlinx.coroutines.launch
|
||||
import kotlinx.coroutines.withContext
|
||||
import okhttp3.OkHttpClient
|
||||
import okhttp3.Request
|
||||
import okhttp3.coroutines.executeAsync
|
||||
import java.net.URLEncoder
|
||||
import java.util.concurrent.TimeUnit
|
||||
|
||||
@@ -141,7 +142,7 @@ private suspend fun resolveNip05Http(identifier: String): String? {
|
||||
.readTimeout(10, TimeUnit.SECONDS)
|
||||
.build()
|
||||
val request = Request.Builder().url(url).build()
|
||||
val response = client.newCall(request).execute()
|
||||
val response = client.newCall(request).executeAsync()
|
||||
response.use { resp ->
|
||||
if (!resp.isSuccessful) return@withContext null
|
||||
val body = resp.body.string()
|
||||
|
||||
+3
-2
@@ -41,6 +41,7 @@ import kotlinx.coroutines.delay
|
||||
import kotlinx.coroutines.isActive
|
||||
import kotlinx.coroutines.withContext
|
||||
import okhttp3.Request
|
||||
import okhttp3.coroutines.executeAsync
|
||||
import org.jetbrains.skia.Bitmap
|
||||
import org.jetbrains.skia.Codec
|
||||
import org.jetbrains.skia.Data
|
||||
@@ -131,10 +132,10 @@ fun AnimatedGifImage(
|
||||
}
|
||||
}
|
||||
|
||||
private fun decodeGifFrames(url: String): GifFrames? =
|
||||
private suspend fun decodeGifFrames(url: String): GifFrames? =
|
||||
try {
|
||||
val request = Request.Builder().url(url).build()
|
||||
val response = gifHttpClient.newCall(request).execute()
|
||||
val response = gifHttpClient.newCall(request).executeAsync()
|
||||
val bytes = response.body.bytes()
|
||||
|
||||
val skData = Data.makeFromBytes(bytes)
|
||||
|
||||
+2
-1
@@ -24,6 +24,7 @@ import com.vitorpamplona.amethyst.desktop.network.DesktopHttpClient
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.withContext
|
||||
import okhttp3.Request
|
||||
import okhttp3.coroutines.executeAsync
|
||||
import java.awt.FileDialog
|
||||
import java.awt.Frame
|
||||
import java.io.File
|
||||
@@ -59,7 +60,7 @@ object SaveMediaAction {
|
||||
return withContext(Dispatchers.IO) {
|
||||
try {
|
||||
val request = Request.Builder().url(url).build()
|
||||
val response = httpClient.newCall(request).execute()
|
||||
val response = httpClient.newCall(request).executeAsync()
|
||||
response.use { resp ->
|
||||
if (!resp.isSuccessful) return@withContext null
|
||||
val total = resp.body.contentLength()
|
||||
|
||||
+2
-1
@@ -31,6 +31,7 @@ import okhttp3.OkHttpClient
|
||||
import okhttp3.Request
|
||||
import okhttp3.RequestBody.Companion.toRequestBody
|
||||
import okhttp3.Response
|
||||
import okhttp3.coroutines.executeAsync
|
||||
import java.io.EOFException
|
||||
import java.io.IOException
|
||||
import java.net.SocketException
|
||||
@@ -167,7 +168,7 @@ class OkHttpNestsClient(
|
||||
val request = buildRequest()
|
||||
val response: Response =
|
||||
try {
|
||||
httpClient(url).newCall(request).execute()
|
||||
httpClient(url).newCall(request).executeAsync()
|
||||
} catch (e: SocketException) {
|
||||
transportError = e
|
||||
if (++transportAttempts >= MAX_TRANSPORT_RETRIES) throw NestsException("Failed to reach $url", e)
|
||||
|
||||
+1
-1
@@ -58,7 +58,7 @@ class OkHttpBitcoinExplorer(
|
||||
.get()
|
||||
.build()
|
||||
|
||||
return client.newCall(request).execute().use {
|
||||
return client.newCall(request).executeAsync().use {
|
||||
if (it.isSuccessful) {
|
||||
Log.d("OkHttpBlockstreamExplorer") { "$baseAPI/block/$hash" }
|
||||
|
||||
|
||||
Reference in New Issue
Block a user