A production ANR on a Pixel 8 (Amethyst 1.13.1, anr_2026-08-03-12-55-26-256)
showed the app burning 596% CPU — 6 of 9 cores — with the main thread stuck in
WaitingForGcToComplete. 37 of 52 runnable DefaultDispatcher workers sat at ONE
program point inside PoolRequests.onIncomingMessage and 12 more at one point in
syncState, all state=R, while the single thread actually holding the lock was
itself parked in GC.
Root cause: RequestSubscriptionState.withLock was a raw busy-wait
(`while (lock.exchange(true)) { while (lock.load()) {} }`) with no yield or
backoff, and being `inline` it disappeared into its callers' frames. The lock is
per subId, but one subId spans every relay it runs on — 191 live sockets on that
device — so dozens of relay-dispatch threads piled onto a single AtomicBoolean.
Spinning is only correct when the holder cannot be descheduled; on Android it
always can.
The fix stripes the lock per (subId, relay) rather than making waiting cheaper.
All 11 withLock bodies in PoolRequests are already scoped to exactly one relay,
and every field of RequestSubscriptionState is keyed by relay, so the sharing was
purely an artifact of mutableMapOf not being thread-safe. State moves into a
ConcurrentMap<T, RelayState>; locks live in a fixed 32-entry stripe array that is
never mutated, so lock identity stays stable — if locks lived inside the map
values, a thread holding one while another dropped and re-created that entry
would leave both inside the critical section excluding nothing.
A suspending Mutex was measured and rejected: it needs 262 method overrides and
110 call sites to become suspend, and ran at 0.35-0.63x the current throughput.
Measured (LockDesignComparisonBenchmark, 191 relays / 64 dispatcher threads):
striped vs per-sub lock 1.5-2.8x throughput, bystander p50 halved
On device (SM-T220, playBenchmark, same account, n=3 per design):
DefaultDispatcher CPU -35% mean / -30% median vs the spin lock,
with non-overlapping ranges; GC -18%
Plus 10 min of driven UI stress (feed, profiles, chat, notifications,
communities): no ANRs, no crashes, thread pools stable.
Also here:
- PlatformLock: new expect/actual parking lock (ReentrantLock on jvmAndroid,
NSRecursiveLock on Apple, spin only on linuxX64 which is a CI target). quartz
cannot use commons' equivalent KmpLock because commons depends on quartz.
- LiveNegentropyIndex had the identical busy-wait with a full list SORT inside
the critical section; switched to PlatformLock.
- ConcurrentMap.remove (+ tests), with a caution that a removable value must not
own a lock callers acquire.
- SpinLockConvoyBenchmark: regression guard asserting contended waiters PARK
rather than spin (fails-before / passes-after). Pure benchmarks are gated
behind -PprodRelayBench=1, so CI cost is 0.3s rather than 51.5s.
Analysis and measurements: quartz/plans/2026-08-03-poolrequests-lock-contention.md
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Quartz Guide for Clients
Here's how to structure a new Twitter-like client.
Architecture
Set up a Context class to wire Quartz components together. Usually there is only one instance of this class.
object AppGraph {
// application-wide scope
private val scope = CoroutineScope(Dispatchers.IO + SupervisorJob())
// the local db
val sqlite = EventStore(dbName = "demo-events.db")
// the local cache that keeps only one copy of each event in memory
val interned = InterningEventStore(sqlite)
// the observable db, that you can produce flows that auto update
val db = ObservableEventStore(interned)
// the client to access relays
val client = NostrClient(websocketBuilder = KtorWebSocket.Builder())
// sends all events, regardless of the subscription, to the local db
val collector = EventCollector(client) { event, _ ->
runCatching {
db.insert(event)
}
}
// update this variable when a user logs in, starts with a guest
var signer: NostrSigner = NostrSignerInternal(KeyPair())
init {
// Periodic NIP-40 sweep — drops expired events from SQLite and
// emits StoreChange.DeleteExpired so live projections drop them
// too. Without this the on-disk store grows monotonically.
scope.launch {
while (isActive) {
delay(15.minutes)
runCatching { db.deleteExpiredEvents() }
}
}
}
}
Then use a view model to subscribe to relays and the local db at the same time, like this:
class NotesFeed(
private val db: ObservableEventStore,
private val client: NostrClient,
) {
private val subId = newSubId()
private val filter = Filter(kinds = listOf(TextNoteEvent.KIND), limit = 100)
private val relays =
setOf(
"wss://relay.damus.io".normalizeRelayUrl(),
"wss://nos.lol".normalizeRelayUrl(),
"wss://relay.nostr.band".normalizeRelayUrl(),
)
val notes: Flow<ProjectionState<TextNoteEvent>> =
db
.project<TextNoteEvent>(filter)
.filterItems { it.value.isNewThread() }
.onStart { client.subscribe(subId, relays.associateWith { listOf(filter) }) }
.onCompletion { client.unsubscribe(subId) }
}
class FeedViewModel(
private val db: ObservableEventStore,
private val client: NostrClient,
) : ViewModel() {
val notesFeed = NotesFeed(db, client)
val feed = notesFeed
.flow
.stateIn(viewModelScope, SharingStarted.WhileSubscribed(5_000), ProjectionState.Loading)
fun send(text: String, signer: NostrSigner) {
viewModelScope.launch {
val signed = signer.sign<TextNoteEvent>(TextNoteEvent.build(text))
// Hits the bus → projection picks it up alongside any inbound relay copy.
db.insert(signed)
client.publish(signed, relays)
}
}
}
Notice that the notes flow is ready for the UI and automatically subscribes
and unsubscribes to any group of relays and filters the user wants. Similarly,
the send function updates both the local db and the relay.
NostrClient connects on-demand: the first subscribe(...) or publish(...) to a relay triggers the socket. There's no need to call client.connect() at startup — it's only useful for resuming after a prior disconnect().
Building a reactive feed UI
A feed screen reads from the view model's feed flow, which only updates when new events arrive or are deleted due to kind 5 deletions, vanish requests or expirations.
fun main() {
application {
val state = rememberWindowState(size = DpSize(560.dp, 720.dp))
Window(onCloseRequest = ::exitApplication, state = state, title = "Nostr Kind 1 Demo") {
MaterialTheme {
val viewModel = remember {
FeedViewModel(AppGraph.db, AppGraph.client, AppGraph.signer)
}
val noteState by viewModel.feed.collectAsStateWithLifecycle()
when (noteState) {
is ProjectionState.Loading -> LoadingFeed()
is ProjectionState.Loaded -> Feed(noteState.items)
}
}
}
}
}
@Composable
private fun LoadingFeed() {
Box(modifier = Modifier.fillMaxSize(), contentAlignment = Alignment.Center) {
CircularProgressIndicator()
}
}
@Composable
private fun Feed(items: List<MutableStateFlow<TextNoteEvent>>) {
LazyColumn(modifier = Modifier.fillMaxSize()) {
items(items = items, key = { it.value.id }) { handle ->
NoteRow(handle)
HorizontalDivider()
}
}
}
@Composable
private fun NoteRow(handle: MutableStateFlow<TextNoteEvent>) {
val event by handle.collectAsStateWithLifecycle()
Text(
text = event.content,
style = MaterialTheme.typography.bodyMedium,
modifier = Modifier.padding(top = 4.dp),
)
}
Notice how each how also subscribe for changes. This is important to receive updates from replaceable and addressable events.
Appendix A
Quartz doesn't offer a Ktor websocket, but you can use this one as reference.
/**
* Ktor-based [WebSocket] for talking to a Nostr relay.
*
* Quartz exposes [WebsocketBuilder] as the only seam between its relay-pool
* and the underlying transport, so all this class has to do is open a Ktor
* websocket session, forward incoming text frames to [out], and let Quartz
* drive sends.
*/
class KtorWebSocket(
private val url: NormalizedRelayUrl,
private val httpClient: HttpClient,
private val out: WebSocketListener,
) : WebSocket {
private val scope = CoroutineScope(Dispatchers.IO + SupervisorJob())
private var session: DefaultWebSocketSession? = null
private var readerJob: Job? = null
override fun needsReconnect(): Boolean = session == null
override fun connect() {
readerJob =
scope.launch {
try {
val s = httpClient.webSocketSession(urlString = url.url)
session = s
out.onOpen(0, false)
for (frame in s.incoming) {
if (frame is Frame.Text) {
out.onMessage(frame.readText())
}
}
val reason = s.closeReason.await()
out.onClosed(
code =
reason?.code?.toInt() ?: CloseReason.Codes.NORMAL.code
.toInt(),
reason = reason?.message ?: "",
)
} catch (t: Throwable) {
out.onFailure(t, null, null)
} finally {
session = null
}
}
}
override fun disconnect() {
val s = session
session = null
readerJob?.cancel()
readerJob = null
if (s != null) {
runBlocking { s.close(CloseReason(CloseReason.Codes.NORMAL, "client disconnect")) }
}
scope.cancel()
}
override fun send(msg: String): Boolean {
val s = session ?: return false
scope.launch { s.send(msg) }
return true
}
/**
* The factory Quartz hands to [com.vitorpamplona.quartz.nip01Core.relay.client.NostrClient].
* One [HttpClient] is shared by every relay in the pool.
*/
class Builder(
private val httpClient: HttpClient = defaultClient(),
) : WebsocketBuilder {
override fun build(
url: NormalizedRelayUrl,
out: WebSocketListener,
): WebSocket = KtorWebSocket(url, httpClient, out)
companion object {
fun defaultClient() =
HttpClient(CIO) {
install(WebSockets)
}
}
}
}