mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-05 19:28:25 +00:00
revert(quartz): inline small-REQ fast path — no wire-level win, floor is transport-side
Reverts the queryRawInline fast path (fb29d655,1b786f31) per the keep-only-winners rule. Three relayBench runs at 50k (baseline, cap-256 where the path never engaged, cap-512 where author-archive/by-ids/ 500-limit feeds genuinely took it) showed no movement outside the container drift band — strfry's own numbers drifted ±30% between runs and inline-eligible scenarios moved the same as ineligible ones. The in-process win was real but small (~17%, 0.60 -> 0.50 ms per ~21-row REQ); the wire-level p50 is 1.2-1.7 ms, so the missing ~1 ms per REQ sits in the transport (Ktor frame send path + client round trip) — backlog item 6 territory, not dispatch. Findings, numbers, and the do-not-retry note live in quartz/plans/2026-07-04-small-req-floor.md; SmallReqFloorBenchmark stays as the measurement tool. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TtDNpayEYvJH7QuPswND3A
This commit is contained in:
@@ -0,0 +1,58 @@
|
||||
# Small-REQ dispatch floor — investigated, inline fast path reverted
|
||||
|
||||
**Status: closed (negative result recorded).** Backlog item 2 of the
|
||||
relay performance campaign.
|
||||
|
||||
## The gap
|
||||
|
||||
relayBench at 50k events: geode WINS most 500-event query scenarios
|
||||
(hashtag 5.4 vs 8.4 ms, recent-window 4.1 vs 8.5) but loses ~2.5× on
|
||||
small results — author-archive (19 events) 1.2–1.7 ms vs strfry's
|
||||
~0.5–0.6, thread (23) likewise — and the @8conn throughput inverts
|
||||
(strfry 2–3× geode). With ~20-row responses, throughput ≈ 1/latency:
|
||||
there is a fixed per-REQ floor.
|
||||
|
||||
## Decomposition (SmallReqFloorBenchmark, kept in jvmTest)
|
||||
|
||||
In-process at 50k events, ~21 rows/REQ, medians of 400:
|
||||
|
||||
| stage | ms |
|
||||
|---|---:|
|
||||
| A raw store query (SQL + row decode) | 0.18 |
|
||||
| B + live machinery (FilterIndex reg/unreg, dedupe set) | 0.36 |
|
||||
| C + session dispatch (parse, launch, frames) | 0.60 |
|
||||
|
||||
Note: in-memory DBs have no reader pool (`useReader` falls back to the
|
||||
writer mutex), so absolute numbers are conservative vs the file-DB
|
||||
bench setup.
|
||||
|
||||
## What was tried and why it was reverted
|
||||
|
||||
An inline fast path (`SessionBackend.queryRawInline`): REQs with
|
||||
provably bounded replays (limit or ids-count summing ≤ 512) ran their
|
||||
stored replay on the receive coroutine and kept only a live-tail
|
||||
handle — no per-REQ `launch`, no Job, no dispatcher handoffs. It cut
|
||||
in-process time-to-EOSE ~17% (0.60 → 0.50 ms) with full wire-behavior
|
||||
parity (stored→EOSE order, live tail, CLOSE, same-subId replacement).
|
||||
|
||||
Three relayBench runs (baseline, cap-256 [path not engaged — bench
|
||||
filters carry limit=500 or no limit], cap-512 [engaged for
|
||||
author-archive/by-ids/500-limit feeds]) showed **no movement outside
|
||||
the container drift band** — strfry's own numbers drifted ±30% run to
|
||||
run, and inline-eligible scenarios moved the same as ineligible ones.
|
||||
Reverted per the keep-only-winners rule.
|
||||
|
||||
## Where the floor actually is
|
||||
|
||||
In-process REQ→EOSE is ~0.5–0.6 ms, but the wire-level p50 is
|
||||
1.2–1.7 ms: the missing ~1 ms per REQ is transport-side — the Ktor CIO
|
||||
frame write path, per-frame sends with no batching, and the client
|
||||
round trip — which is backlog item 6 (websocket send path, 13–25% of
|
||||
ingest CPU in the JFR profile) territory. strfry completes the whole
|
||||
round trip in under 0.6 ms on a single event loop with uWebSockets.
|
||||
|
||||
**Do not retry** coroutine-dispatch shaving for this gap without first
|
||||
measuring the transport side: instrument the time between
|
||||
`RelaySession`'s `onSend` invocation and the frame actually leaving
|
||||
the socket, and compare permessage-deflate/frame-batching settings
|
||||
against strfry's uWS configuration.
|
||||
@@ -1,6 +1,6 @@
|
||||
# quartz plans
|
||||
|
||||
_Audited 2026-06-30. 10 plans: 7 shipped (archived), 0 in-progress, 3 queued, 0 abandoned._
|
||||
_Audited 2026-06-30. 11 plans: 7 shipped (archived), 0 in-progress, 3 queued, 1 closed (negative result)._
|
||||
|
||||
## Queued
|
||||
| Plan | Summary |
|
||||
@@ -8,6 +8,7 @@ _Audited 2026-06-30. 10 plans: 7 shipped (archived), 0 in-progress, 3 queued, 0
|
||||
| [2026-05-08-local-headers-explorer.md](2026-05-08-local-headers-explorer.md) | Headers-only Bitcoin P2P client to verify NIP-03 OTS attestations without a trusted block explorer. |
|
||||
| [2026-06-12-giftwrap-deletion-requests.md](2026-06-12-giftwrap-deletion-requests.md) | Let a recipient-authored kind-5 delete/block a gift wrap (kind 1059) addressed to them. |
|
||||
| [2026-07-03-incremental-negentropy-storage.md](2026-07-03-incremental-negentropy-storage.md) | Always-current (created_at, id) index so cold NEG-OPENs stop paying a full scan + seal (~340 ms at 50k vs strfry's ~21 ms). |
|
||||
| [2026-07-04-small-req-floor.md](2026-07-04-small-req-floor.md) | Small-REQ dispatch floor: decomposed, inline fast path tried and reverted (no wire-level win); floor is transport-side. |
|
||||
|
||||
## Archived (shipped)
|
||||
| Plan | Summary |
|
||||
|
||||
+24
-100
@@ -35,7 +35,6 @@ import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.CloseCmd
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.CountCmd
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.EventCmd
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.ReqCmd
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.server.backend.RequestContext
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.server.backend.SessionBackend
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.server.policies.IRelayPolicy
|
||||
@@ -50,6 +49,7 @@ import com.vitorpamplona.quartz.utils.Log
|
||||
import com.vitorpamplona.quartz.utils.cache.LargeCache
|
||||
import kotlinx.coroutines.CancellationException
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.Job
|
||||
import kotlinx.coroutines.channels.ClosedSendChannelException
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlin.concurrent.atomics.AtomicLong
|
||||
@@ -75,16 +75,7 @@ class RelaySession(
|
||||
*/
|
||||
val id: Long = nextConnectionId(),
|
||||
) : AutoCloseable {
|
||||
/**
|
||||
* One active REQ. Launched subscriptions cancel their coroutine;
|
||||
* inline subscriptions (the bounded small-REQ fast path, which never
|
||||
* had a coroutine) close their live-tail registration directly.
|
||||
*/
|
||||
private fun interface ActiveSubscription {
|
||||
fun cancel()
|
||||
}
|
||||
|
||||
private val subscriptions = LargeCache<String, ActiveSubscription>()
|
||||
private val subscriptions = LargeCache<String, Job>()
|
||||
|
||||
/**
|
||||
* The authenticated-identity store for this connection. The engine is the
|
||||
@@ -112,8 +103,8 @@ class RelaySession(
|
||||
|
||||
private fun addSubscription(
|
||||
subId: String,
|
||||
sub: ActiveSubscription,
|
||||
) = subscriptions.put(subId, sub)
|
||||
job: Job,
|
||||
) = subscriptions.put(subId, job)
|
||||
|
||||
private fun cancelSubscription(subId: String): Boolean =
|
||||
subscriptions.remove(subId)?.let {
|
||||
@@ -122,7 +113,7 @@ class RelaySession(
|
||||
} ?: false
|
||||
|
||||
fun cancelAllSubscriptions() {
|
||||
subscriptions.forEach { _, sub -> sub.cancel() }
|
||||
subscriptions.forEach { _, job -> job.cancel() }
|
||||
subscriptions.clear()
|
||||
negentropy.clear()
|
||||
}
|
||||
@@ -272,30 +263,7 @@ class RelaySession(
|
||||
}
|
||||
|
||||
// -- NIP-01: REQ ----------------------------------------------------------
|
||||
|
||||
/**
|
||||
* A REQ qualifies for the inline fast path when its replay is
|
||||
* provably bounded: every filter carries a `limit`, or an `ids` list
|
||||
* (ids are unique keys, so the result can't exceed the list), and
|
||||
* the bounds sum to at most [INLINE_REPLAY_MAX_ROWS]. Inline replays
|
||||
* run on this connection's receive coroutine — no per-REQ launch, no
|
||||
* dispatcher handoffs (profiled at ~2/3 of a ~20-row REQ's
|
||||
* time-to-EOSE) — at the price of delaying the connection's NEXT
|
||||
* command until EOSE, which the row cap keeps in the
|
||||
* single-digit-milliseconds range. Unbounded REQs keep the launched
|
||||
* path so a CLOSE can always interrupt a genuinely giant replay.
|
||||
*/
|
||||
private fun isInlineEligible(filters: List<Filter>): Boolean {
|
||||
var total = 0
|
||||
for (f in filters) {
|
||||
val bound = f.limit ?: f.ids?.size ?: return false
|
||||
total += bound
|
||||
if (total > INLINE_REPLAY_MAX_ROWS) return false
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
private suspend fun handleReq(cmd: ReqCmd) {
|
||||
private fun handleReq(cmd: ReqCmd) {
|
||||
// Ask the policy whether a *new* subscription may open (e.g. a
|
||||
// max_subscriptions cap). A re-REQ on an existing id replaces it
|
||||
// 1-for-1 and doesn't grow the count, so it's exempt. Checked before
|
||||
@@ -320,54 +288,6 @@ class RelaySession(
|
||||
// Policy may rewrite filters to match the user's access level.
|
||||
val filters = (result as PolicyResult.Accepted).cmd.filters
|
||||
|
||||
// The `["EVENT","<subId>",` prefix is built once per
|
||||
// subscription, not per row (zero-decode paths only).
|
||||
fun framePrefix() =
|
||||
buildString {
|
||||
append("[\"EVENT\",")
|
||||
RawEvent.appendJsonQuoted(this, cmd.subId)
|
||||
append(',')
|
||||
}
|
||||
|
||||
fun sendStored(
|
||||
framePrefix: String,
|
||||
raw: RawEvent,
|
||||
) = sendRaw(
|
||||
buildString(framePrefix.length + raw.jsonTags.length + raw.content.length + 256) {
|
||||
append(framePrefix)
|
||||
raw.appendJsonObjectTo(this)
|
||||
append(']')
|
||||
},
|
||||
)
|
||||
|
||||
// Small-REQ fast path: a provably bounded replay runs inline on
|
||||
// this receive coroutine — no launch, no dispatcher handoffs —
|
||||
// and only the live-tail handle is retained. See
|
||||
// [SessionBackend.queryRawInline] for the contract.
|
||||
if (!policy.filtersOutgoingEvents && isInlineEligible(filters)) {
|
||||
val handle =
|
||||
try {
|
||||
val prefix = framePrefix()
|
||||
store.queryRawInline(
|
||||
ctx = requestContext,
|
||||
filters = filters,
|
||||
onEachStored = { raw -> sendStored(prefix, raw) },
|
||||
onEachLive = { event -> send(EventMessage(cmd.subId, event)) },
|
||||
onEose = { send(EoseMessage(cmd.subId)) },
|
||||
)
|
||||
} catch (e: CancellationException) {
|
||||
throw e
|
||||
} catch (e: Exception) {
|
||||
send(ClosedMessage.of(cmd.subId, MachineReadablePrefix.ERROR, e.message ?: "query failed"))
|
||||
return
|
||||
}
|
||||
if (handle != null) {
|
||||
addSubscription(cmd.subId) { handle.close() }
|
||||
return
|
||||
}
|
||||
// Backend without an inline path: fall through to launch.
|
||||
}
|
||||
|
||||
val job =
|
||||
scope.launch {
|
||||
try {
|
||||
@@ -388,11 +308,26 @@ class RelaySession(
|
||||
// Zero-decode path: the stored replay splices raw
|
||||
// storage strings straight into wire frames — no tags
|
||||
// parse, no Event materialization, no re-serialize.
|
||||
val prefix = framePrefix()
|
||||
// The `["EVENT","<subId>",` prefix is built once per
|
||||
// subscription, not per row.
|
||||
val framePrefix =
|
||||
buildString {
|
||||
append("[\"EVENT\",")
|
||||
RawEvent.appendJsonQuoted(this, cmd.subId)
|
||||
append(',')
|
||||
}
|
||||
store.queryRaw(
|
||||
ctx = requestContext,
|
||||
filters = filters,
|
||||
onEachStored = { raw -> sendStored(prefix, raw) },
|
||||
onEachStored = { raw ->
|
||||
sendRaw(
|
||||
buildString(framePrefix.length + raw.jsonTags.length + raw.content.length + 256) {
|
||||
append(framePrefix)
|
||||
raw.appendJsonObjectTo(this)
|
||||
append(']')
|
||||
},
|
||||
)
|
||||
},
|
||||
onEachLive = { event -> send(EventMessage(cmd.subId, event)) },
|
||||
onEose = { send(EoseMessage(cmd.subId)) },
|
||||
)
|
||||
@@ -408,7 +343,7 @@ class RelaySession(
|
||||
}
|
||||
}
|
||||
|
||||
addSubscription(cmd.subId) { job.cancel() }
|
||||
addSubscription(cmd.subId, job)
|
||||
}
|
||||
|
||||
// -- NIP-01: CLOSE --------------------------------------------------------
|
||||
@@ -428,16 +363,5 @@ class RelaySession(
|
||||
|
||||
/** Allocates the next process-unique connection id. */
|
||||
fun nextConnectionId(): Long = connectionIdSeq.fetchAndAdd(1L)
|
||||
|
||||
/**
|
||||
* Ceiling on the summed per-filter bounds (limit or ids count)
|
||||
* for the inline REQ fast path. Sized so the worst-case inline
|
||||
* replay stays in single-digit milliseconds (500-row replays
|
||||
* measure ~5–8 ms of storage+frame work) — small enough that
|
||||
* delaying the connection's next command is unnoticeable, large
|
||||
* enough to cover the limit≤500 feed/notifications/archive
|
||||
* shapes real clients hammer relays with.
|
||||
*/
|
||||
const val INLINE_REPLAY_MAX_ROWS = 512
|
||||
}
|
||||
}
|
||||
|
||||
+2
-30
@@ -258,31 +258,6 @@ class LiveEventStore(
|
||||
onEachLive: (Event) -> Unit,
|
||||
onEose: () -> Unit,
|
||||
) {
|
||||
val handle = queryRawInline(ctx, filters, onEachStored, onEachLive, onEose)
|
||||
try {
|
||||
// Suspend until the caller's coroutine is cancelled (NIP-01
|
||||
// CLOSE or connection drop); the live tail runs meanwhile.
|
||||
awaitCancellation()
|
||||
} finally {
|
||||
handle.close()
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* The registration + replay + EOSE core of [queryRaw], run entirely
|
||||
* on the calling coroutine. Small bounded REQs take this directly
|
||||
* (see [SessionBackend.queryRawInline]) and skip the per-REQ
|
||||
* coroutine; [queryRaw] wraps it with `awaitCancellation` for the
|
||||
* launched path. Never returns `null` — this backend always
|
||||
* supports the inline path.
|
||||
*/
|
||||
override suspend fun queryRawInline(
|
||||
ctx: RequestContext,
|
||||
filters: List<Filter>,
|
||||
onEachStored: (RawEvent) -> Unit,
|
||||
onEachLive: (Event) -> Unit,
|
||||
onEose: () -> Unit,
|
||||
): SessionBackend.LiveSubscriptionHandle {
|
||||
drainFtsIfSearching(filters)
|
||||
val seenLock = AtomicBoolean(false)
|
||||
var seenIds: HashSet<String>? = HashSet(1024)
|
||||
@@ -315,14 +290,11 @@ class LiveEventStore(
|
||||
onEachStored(raw)
|
||||
}
|
||||
onEose()
|
||||
// Drop the dedupe set so the live path stops paying for it.
|
||||
seenLocked { seenIds = null }
|
||||
} catch (e: Throwable) {
|
||||
// A failed replay must not leak the registration.
|
||||
awaitCancellation()
|
||||
} finally {
|
||||
index.unregister(sub)
|
||||
throw e
|
||||
}
|
||||
return SessionBackend.LiveSubscriptionHandle { index.unregister(sub) }
|
||||
}
|
||||
|
||||
override suspend fun count(
|
||||
|
||||
-36
@@ -82,42 +82,6 @@ interface SessionBackend {
|
||||
onEose: () -> Unit,
|
||||
): Unit = query(ctx, filters, onEachLive, onEose)
|
||||
|
||||
/**
|
||||
* Keeps a live registration alive after [queryRawInline] returned.
|
||||
* [close] detaches it — idempotent, callable from any thread.
|
||||
*/
|
||||
fun interface LiveSubscriptionHandle {
|
||||
fun close()
|
||||
}
|
||||
|
||||
/**
|
||||
* Inline variant of [queryRaw] for BOUNDED replays: runs the stored
|
||||
* replay and [onEose] on the *calling* coroutine — no per-REQ
|
||||
* `launch`, no dispatcher handoffs, no job to keep alive — and
|
||||
* returns a [LiveSubscriptionHandle] whose `close()` ends the live
|
||||
* tail (live events keep arriving via [onEachLive] from the ingest
|
||||
* path until then). This is the small-REQ fast path: profiling put
|
||||
* the per-REQ coroutine machinery at ~2/3 of a ~20-row REQ's
|
||||
* time-to-EOSE.
|
||||
*
|
||||
* Callers MUST only use it when the replay is bounded (e.g. every
|
||||
* filter carries a small `limit`): the replay occupies the session's
|
||||
* receive coroutine, so a giant replay here would delay that
|
||||
* connection's subsequent commands — including the CLOSE that could
|
||||
* have cancelled it. Unbounded REQs belong on [queryRaw] in a
|
||||
* separate coroutine.
|
||||
*
|
||||
* Returns `null` when the backend has no inline path (the default) —
|
||||
* callers fall back to [queryRaw].
|
||||
*/
|
||||
suspend fun queryRawInline(
|
||||
ctx: RequestContext,
|
||||
filters: List<Filter>,
|
||||
onEachStored: (RawEvent) -> Unit,
|
||||
onEachLive: (Event) -> Unit,
|
||||
onEose: () -> Unit,
|
||||
): LiveSubscriptionHandle? = null
|
||||
|
||||
/** Answers a NIP-45 COUNT with an exact cardinality for the caller in [ctx]. */
|
||||
suspend fun count(
|
||||
ctx: RequestContext,
|
||||
|
||||
-169
@@ -1,169 +0,0 @@
|
||||
/*
|
||||
* Copyright (c) 2025 Vitor Pamplona
|
||||
*
|
||||
* Permission is hereby granted, free of charge, to any person obtaining a copy of
|
||||
* this software and associated documentation files (the "Software"), to deal in
|
||||
* the Software without restriction, including without limitation the rights to use,
|
||||
* copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the
|
||||
* Software, and to permit persons to whom the Software is furnished to do so,
|
||||
* subject to the following conditions:
|
||||
*
|
||||
* The above copyright notice and this permission notice shall be included in all
|
||||
* copies or substantial portions of the Software.
|
||||
*
|
||||
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
|
||||
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS
|
||||
* FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR
|
||||
* COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN
|
||||
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
|
||||
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
|
||||
*/
|
||||
package com.vitorpamplona.quartz.nip01Core.relay.server
|
||||
|
||||
import com.vitorpamplona.quartz.nip01Core.core.Event
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.server.policies.EmptyPolicy
|
||||
import com.vitorpamplona.quartz.nip01Core.store.sqlite.EventStore
|
||||
import kotlinx.coroutines.ExperimentalCoroutinesApi
|
||||
import kotlinx.coroutines.test.UnconfinedTestDispatcher
|
||||
import kotlinx.coroutines.test.runTest
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertTrue
|
||||
|
||||
/**
|
||||
* NIP-01 semantics of the inline small-REQ fast path (a REQ whose
|
||||
* filters all carry a small `limit` runs its replay on the receive
|
||||
* coroutine — see [RelaySession.INLINE_REPLAY_MAX_ROWS]). The wire
|
||||
* behavior must be indistinguishable from the launched path: stored
|
||||
* replay, EOSE, live tail, CLOSE handling, and same-subId replacement.
|
||||
*/
|
||||
@OptIn(ExperimentalCoroutinesApi::class)
|
||||
class InlineReqFastPathTest {
|
||||
private fun hexId(n: Int): String = n.toString().padStart(64, '0')
|
||||
|
||||
private val pubkey = "46fcbe3065eaf1ae7811465924e48923363ff3f526bd6f73d7c184b16bd8ce4d"
|
||||
private val sig = "0".repeat(128)
|
||||
|
||||
private fun testEvent(
|
||||
idSeed: Int,
|
||||
kind: Int = 1,
|
||||
createdAt: Long = idSeed.toLong(),
|
||||
) = Event(hexId(idSeed), pubkey, createdAt, kind, emptyArray(), "note $idSeed", sig)
|
||||
|
||||
private suspend fun serverWith(
|
||||
dispatcher: kotlinx.coroutines.CoroutineDispatcher,
|
||||
vararg events: Event,
|
||||
): NostrServer {
|
||||
val store = EventStore(null)
|
||||
events.forEach { store.insert(it) }
|
||||
return NostrServer(store = store, policyBuilder = { EmptyPolicy }, parentContext = dispatcher)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun boundedReqRepliesStoredThenEoseThenLiveTail() =
|
||||
runTest {
|
||||
val dispatcher = UnconfinedTestDispatcher(testScheduler)
|
||||
val server = serverWith(dispatcher, testEvent(1), testEvent(2, kind = 7))
|
||||
|
||||
val frames = mutableListOf<String>()
|
||||
val session = server.connect { frames.add(it) }
|
||||
|
||||
// limit=10 → inline path.
|
||||
session.receive("""["REQ","small",{"kinds":[1],"limit":10}]""")
|
||||
|
||||
assertEquals(1, frames.count { it.startsWith("[\"EVENT\",\"small\"") })
|
||||
assertTrue(frames.first { it.startsWith("[\"EVENT\"") }.contains(hexId(1)))
|
||||
assertTrue(frames.last().startsWith("[\"EOSE\",\"small\""))
|
||||
|
||||
// Live tail: a matching publish after EOSE must reach the sub.
|
||||
session.receive("""["EVENT",${testEvent(3).toJson()}]""")
|
||||
assertTrue(frames.any { it.startsWith("[\"EVENT\",\"small\"") && it.contains(hexId(3)) })
|
||||
// Non-matching kind stays out.
|
||||
session.receive("""["EVENT",${testEvent(4, kind = 7).toJson()}]""")
|
||||
assertTrue(frames.none { it.startsWith("[\"EVENT\",\"small\"") && it.contains(hexId(4)) })
|
||||
|
||||
server.close()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun closeDetachesTheInlineLiveTail() =
|
||||
runTest {
|
||||
val dispatcher = UnconfinedTestDispatcher(testScheduler)
|
||||
val server = serverWith(dispatcher, testEvent(1))
|
||||
|
||||
val frames = mutableListOf<String>()
|
||||
val session = server.connect { frames.add(it) }
|
||||
|
||||
session.receive("""["REQ","small",{"kinds":[1],"limit":10}]""")
|
||||
session.receive("""["CLOSE","small"]""")
|
||||
|
||||
session.receive("""["EVENT",${testEvent(2).toJson()}]""")
|
||||
assertTrue(frames.none { it.startsWith("[\"EVENT\",\"small\"") && it.contains(hexId(2)) })
|
||||
|
||||
server.close()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun sameSubIdReplacesInlineSubscription() =
|
||||
runTest {
|
||||
val dispatcher = UnconfinedTestDispatcher(testScheduler)
|
||||
val server = serverWith(dispatcher, testEvent(1), testEvent(2, kind = 7))
|
||||
|
||||
val frames = mutableListOf<String>()
|
||||
val session = server.connect { frames.add(it) }
|
||||
|
||||
session.receive("""["REQ","sub",{"kinds":[1],"limit":10}]""")
|
||||
// Replace with a kind-7 filter (still inline). The old live
|
||||
// tail must be gone: a new kind-1 publish stays silent, a
|
||||
// kind-7 one is delivered.
|
||||
session.receive("""["REQ","sub",{"kinds":[7],"limit":10}]""")
|
||||
|
||||
session.receive("""["EVENT",${testEvent(3, kind = 1).toJson()}]""")
|
||||
assertTrue(frames.none { it.startsWith("[\"EVENT\",\"sub\"") && it.contains(hexId(3)) })
|
||||
session.receive("""["EVENT",${testEvent(4, kind = 7).toJson()}]""")
|
||||
assertTrue(frames.any { it.startsWith("[\"EVENT\",\"sub\"") && it.contains(hexId(4)) })
|
||||
|
||||
server.close()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun unboundedReqStillWorksViaLaunchedPath() =
|
||||
runTest {
|
||||
val dispatcher = UnconfinedTestDispatcher(testScheduler)
|
||||
val server = serverWith(dispatcher, testEvent(1))
|
||||
|
||||
val frames = mutableListOf<String>()
|
||||
val session = server.connect { frames.add(it) }
|
||||
|
||||
// No limit → not inline-eligible; must behave identically.
|
||||
session.receive("""["REQ","big",{"kinds":[1]}]""")
|
||||
|
||||
assertTrue(frames.any { it.startsWith("[\"EVENT\",\"big\"") && it.contains(hexId(1)) })
|
||||
assertTrue(frames.any { it.startsWith("[\"EOSE\",\"big\"") })
|
||||
|
||||
session.receive("""["EVENT",${testEvent(2).toJson()}]""")
|
||||
assertTrue(frames.any { it.startsWith("[\"EVENT\",\"big\"") && it.contains(hexId(2)) })
|
||||
|
||||
server.close()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun oversizedLimitSumIsNotInlineEligible() =
|
||||
runTest {
|
||||
// Behavior parity either way — this documents the boundary:
|
||||
// summed limits over the cap route to the launched path and
|
||||
// still answer correctly.
|
||||
val dispatcher = UnconfinedTestDispatcher(testScheduler)
|
||||
val server = serverWith(dispatcher, testEvent(1))
|
||||
|
||||
val frames = mutableListOf<String>()
|
||||
val session = server.connect { frames.add(it) }
|
||||
|
||||
session.receive("""["REQ","big",{"kinds":[1],"limit":${RelaySession.INLINE_REPLAY_MAX_ROWS + 1}}]""")
|
||||
|
||||
assertTrue(frames.any { it.startsWith("[\"EVENT\",\"big\"") && it.contains(hexId(1)) })
|
||||
assertTrue(frames.any { it.startsWith("[\"EOSE\",\"big\"") })
|
||||
|
||||
server.close()
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user