Files
amethyst/quartz/RELAY.md
Claude 57ffb3386d feat(geode): enable the tag+kind+pubkey index, refresh measured docs
TagAuthorIndexBenchmark at 1M events settles the flag: the DM-room
shape (kinds + authors + #p, 65 client assembler call sites) drops
14.2 ms -> 0.66 ms (~21x, growing with corpus size) while batch-insert
cost stays inside run noise (49.0 vs 47.4 us/event). Existing relay
DBs build the index on next open via ensureOptionalIndexes.

Also refreshes the docs the numbers made stale: IndexingStrategy KDoc
now records the 200k and 1M measurements instead of a TODO,
MergeQueryExecutor's tag-merge note points at the new relayBench
reactions-watch scenario, FsQueryPlanner/FsDriverSelectionBenchmark
reflect the landed cost-based pick (149 ms -> 4.0 ms at 30k events),
and RELAY.md documents that strategy flag flips materialize indexes on
the next open.

Verified: quartz jvmTest store suites, geode test (126), desktopApp
LocalRelayStore tests (5, incl. reopening a default-strategy DB with
the new pubkey-alone flag), relayBench compiles.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01RzkoN3SJHCZWRAiadXXG4w
2026-07-21 15:07:12 +00:00

20 KiB

Building a Nostr Relay with Quartz

Quartz provides a transport-agnostic relay engine. You provide a send callback per connection, and it gives you a RelaySession that accepts raw JSON strings. Plug it into Ktor or any WebSocket transport.

There are two engines on the same RelaySession core: NostrServer for storage-backed relays (an IEventStore with a live tail after EOSE) and EventSourceServer for non-storage relays (search, redirector, computed data — see Non-Storage Relays).

Both NostrServer and EventStore implement AutoCloseable.

Quick Start

1. Add Dependencies

dependencies {
    implementation("io.ktor:ktor-server-core:3.1.1")
    implementation("io.ktor:ktor-server-netty:3.1.1")
    implementation("io.ktor:ktor-server-websockets:3.1.1")
    implementation(project(":quartz"))
}

2. Create the Relay Server

fun main() {
    val store = EventStore(dbName = "relay-events.db")
    val server = NostrServer(store)

    embeddedServer(Netty, port = 7777) {
        install(WebSockets)

        routing {
            webSocket("/") {
                server.serve(
                    send = { json -> launch { send(Frame.Text(json)) } },
                ) { session ->
                    for (frame in incoming) {
                        if (frame is Frame.Text) {
                            session.receive(frame.readText())
                        }
                    }
                }
            }
        }
    }.start(wait = true)
}

That's it. You now have a NIP-01 compliant relay running on ws://localhost:7777.

Event Store

SQLite (Persistent)

val store = EventStore(dbName = "relay-events.db")

Uses androidx.sqlite with WAL journal mode and a 32 MB memory cache.

In-Memory (Testing)

val store = EventStore(null)

Indexing Strategy

Control which indexes are created:

val store = EventStore(
    dbName = "relay-events.db",
    indexStrategy = DefaultIndexingStrategy(
        indexEventsByCreatedAtAlone = false,
        indexTagsByCreatedAtAlone = false,
        indexTagsWithKindAndPubkey = false,
        useAndIndexIdOnOrderBy = false,
        indexFullTextSearch = true,
    ),
)

By default, all single-letter tags with values are indexed. Override shouldIndex(kind, tag) for custom behavior. More indexes = faster queries but larger database.

Flag flips are safe on existing databases: any flag-gated index the strategy wants but the on-disk schema lacks is created on the next open (idempotent CREATE INDEX IF NOT EXISTS, one-time build cost) — no schema version bump involved. Disabling a flag never drops an existing index.

indexFullTextSearch defaults to true and controls the NIP-50 full-text index (event_fts). Set it to false when search is served elsewhere (e.g. a Vespa backend, or a SearchEventSource as shown below): inserts skip the FTS tokenization cost, no event_fts table/trigger is created, and any filter carrying a non-empty search term returns no matches.

Non-Storage Relays (search, redirector, computed)

A relay's job is to answer REQs. When the answer doesn't come from a stored set — a NIP-50 search that forwards to an HTTP backend, a relay that emits computed/projected data — implement EventSource and serve it with EventSourceServer. You supply a Flow<Event>; the engine owns the wire protocol (challenge/auth, command parsing, policy, EVENT/EOSE/CLOSED framing, subscription lifecycle), so there's no hand-written read loop.

class SearchEventSource(private val backend: SearchApi) : EventSource {
    override fun events(ctx: RequestContext, filters: List<Filter>): Flow<Event> = flow {
        // ctx says who is asking — score/restrict results from their perspective.
        val viewer = ctx.authenticatedUsers.firstOrNull()
        filters.forEach { f ->
            f.search?.let { raw ->
                val q = SearchQuery.parse(raw)
                backend.search(q.terms, domain = q.domain, language = q.language, viewer = viewer)
                    .forEach { emit(it) }
            }
        }
    }
    // COUNT defaults to counting events(); override for a cheaper backend count.
}

fun main() {
    val server = EventSourceServer(
        source = SearchEventSource(searchApi),
        policyBuilder = { FullAuthPolicy(relay) }, // optional NIP-42 gating
    )

    embeddedServer(Netty, port = 7777) {
        install(WebSockets)
        routing {
            webSocket("/") {
                server.serve(send = { json -> launch { send(Frame.Text(json)) } }) { session ->
                    for (frame in incoming) {
                        if (frame is Frame.Text) session.receive(frame.readText())
                    }
                }
            }
        }
    }.start(wait = true)
}

EOSE = flow completion. The engine sends EOSE when events(...) completes, which is the natural shape for finite queries. Relays that need an open-ended live tail after EOSE should use the storage path (NostrServer + IEventStore) instead. EVENT publishes are rejected (OK falseblocked: this relay does not accept events) and negentropy is disabled, since there's no stored set. A failure thrown from the flow ends the subscription with CLOSED error: <message>.

EventSourceServer and NostrServer are the two concrete dispatch engines; both build on RelaySession and the shared SessionBackend seam (LiveEventStore is the storage-backed SessionBackend; EventSourceBackend adapts a EventSource). Implement SessionBackend directly only if you need custom control over the EVENT/negentropy paths as well as REQ/COUNT.

Caller-aware sources (who is asking)

Every events/count/countResult call receives a RequestContext carrying the connection's authenticatedUsers (the pubkeys that completed NIP-42 on this socket) and a stable connectionId. That is what makes NIP-42 useful on a non-storage relay: a single shared EventSource instance can serve caller-relative results (trust/relevance scored from the viewer's perspective, "for-you" feeds), restricted content (a pubkey's DMs returned only to that pubkey), or per-connection tenancy — without smuggling auth state through a side channel. ctx.authenticatedUsers is a live view of the engine-owned connection scope, so a REQ that arrives after the AUTH sees the freshly recorded pubkey(s). (The engine records them on a successful NIP-42 AUTH; the policy reads the same scope to gate.)

For per-connection state richer than the pubkey (e.g. a backend session token minted in FullAuthPolicy.authorize), downcast ctx.policy to your own policy subclass and read a typed field — the policy instance is itself per-connection.

Policies

Policies control what clients can do. They validate commands and can rewrite filters.

Built-in Policies

VerifyPolicy (default) — Verifies event signatures and IDs. Rejects malformed events.

val server = NostrServer(store) // Uses VerifyPolicy by default

EmptyPolicy — Accepts everything. Useful for testing.

val server = NostrServer(store, policyBuilder = { EmptyPolicy })

FullAuthPolicy — Requires NIP-42 authentication before accepting any command.

val server = NostrServer(
    store = store,
    policyBuilder = {
        FullAuthPolicy(relay = "wss://myrelay.example.com/".normalizeRelayUrl()!!)
    },
)

FullAuthPolicy already implements the full NIP-42 challenge/verify handshake — you should not re-implement it. To bridge auth to an external system (e.g. exchange the verified event for a backend JWT), override the suspend authorize hook. It runs after the NIP-42 checks pass (and after the rest of the policy chain approves), and the pubkey is recorded only once it returns — so it can do network/disk I/O, and throwing from it rejects the login (OK false) with the connection left unauthenticated (nothing was committed).

Composing Policies

Chain policies with + or PolicyStack. All must approve; first rejection wins.

val server = NostrServer(
    store = store,
    policyBuilder = {
        VerifyPolicy + FullAuthPolicy(relay = "wss://myrelay.example.com/".normalizeRelayUrl()!!)
    },
)

Writing a Custom Policy

Implement IRelayPolicy:

class KindWhitelistPolicy(
    private val allowedKinds: Set<Int>,
) : IRelayPolicy {

    override fun onConnect(send: (Message) -> Unit) { }

    override fun accept(cmd: EventCmd): PolicyResult<EventCmd> =
        if (cmd.event.kind in allowedKinds) {
            PolicyResult.Accepted(cmd)
        } else {
            PolicyResult.Rejected("blocked: kind ${cmd.event.kind} not allowed")
        }

    override fun accept(cmd: ReqCmd) = PolicyResult.Accepted(cmd)
    override fun accept(cmd: CountCmd) = PolicyResult.Accepted(cmd)
    override fun accept(cmd: AuthCmd) = PolicyResult.Accepted(cmd)
}

Use it:

val server = NostrServer(
    store = store,
    policyBuilder = {
        VerifyPolicy + KindWhitelistPolicy(allowedKinds = setOf(0, 1, 3, 7, 30023))
    },
)

NIP-50 Search Queries

Filter.search is the raw NIP-50 string. SearchQuery parses it into the free-text terms plus the typed key:value extensions (domain:, language:, sentiment:, nsfw:, include:spam), so a search relay or redirector doesn't have to re-parse the string. Unknown extensions are preserved and readable via extension(key); toSearchString() re-assembles a canonical query.

val q = SearchQuery.parse(filter.search) // "best apps domain:example.com nsfw:false"
q.terms        // "best apps"
q.domain       // "example.com"
q.nsfwIncluded // false  (NIP-50 default is true when the token is absent)

The SQLite store expects Filter.search to be plain FTS text: : is FTS5 column-filter syntax, so a raw include:spam reaching MATCH raises "no such column: include" instead of matching. Strip the tokens before querying with SearchQuery.stripExtensions(raw) or filter.strippingSearchExtensions() — an extensions-only query collapses to an empty search, which imposes no constraint (NIP-50: unsupported extensions are ignored, not match-nothing). The storage-backed server path (NostrServer / LiveEventStore) already does this, so relays like geode comply out of the box; EventSource backends get the raw string because a real search backend wants the extensions.

A search/redirector relay is just a custom policy (or, for computed results, a custom IEventStore whose query answers the REQ) that reads the parsed query:

override fun accept(cmd: ReqCmd): PolicyResult<ReqCmd> {
    cmd.filters.forEach { f ->
        val q = SearchQuery.parse(f.search)
        if (q.domain != null && q.domain !in allowedDomains) {
            return PolicyResult.Rejected(
                MachineReadablePrefix.RESTRICTED.format("domain not searchable"),
            )
        }
    }
    return PolicyResult.Accepted(cmd)
}

Wire Helpers

Command and Message carry symmetric JSON helpers so you don't have to reach for the mapper directly:

val cmd = Command.fromJson(text)   // ["REQ", "sub", {...}] -> ReqCmd
val json = EoseMessage("sub").toJson()
val msg = Message.fromJson(json)

Build standardized OK/CLOSED reasons with MachineReadablePrefix instead of hand-writing the NIP-01 prefixes (auth-required:, restricted:, error:, …):

session.send(OkMessage.rejected(event.id, MachineReadablePrefix.AUTH_REQUIRED, "log in first"))
session.send(ClosedMessage.of(subId, MachineReadablePrefix.RESTRICTED, "not allowed yet"))
MachineReadablePrefix.parse("rate-limited: slow down") // -> RATE_LIMITED

Testing

@OptIn(ExperimentalCoroutinesApi::class)
class MyRelayTest {

    @Test
    fun clientCanPublishAndSubscribe() = runTest {
        val dispatcher = UnconfinedTestDispatcher(testScheduler)

        NostrServer(
            store = EventStore(null),
            policyBuilder = { EmptyPolicy },
            parentContext = dispatcher,
        ).use { server ->
            val messages = mutableListOf<String>()
            val session = server.connect { messages.add(it) }

            session.receive("""["EVENT",{"id":"${"0".repeat(64)}","pubkey":"${"a".repeat(64)}","created_at":1000,"kind":1,"tags":[],"content":"hello","sig":"${"b".repeat(128)}"}]""")
            assertTrue(messages.any { it.contains("OK") })

            session.receive("""["REQ","sub1",{"kinds":[1]}]""")
            assertTrue(messages.any { it.contains("EVENT") })
            assertTrue(messages.any { it.contains("EOSE") })
        }
    }
}

Limits (NIP-11) — enforce and advertise from one source

RelayLimits is the single source of truth for the relay's operational limits. Pass it to the server and every limit is enforced; call toNip11Limitation() and the same numbers are advertised in your NIP-11 document — they can't drift.

val limits = RelayLimits(
    maxMessageLength = 65_536,
    maxSubscriptions = 20,
    maxFilters = 10,
    maxLimit = 500,
    maxContentLength = 8_196,
    maxEventTags = 2_000,
    createdAtUpperLimit = TimeUtils.now() + 900,
    authRequired = true,
)

val server = NostrServer(store, policyBuilder = { FullAuthPolicy(relay) }, limits = limits)

All of it is enforced by a single LimitsPolicy (which the server prepends to your policy when you pass limits), so limits compose through a PolicyStack like any other policy — you can also split them across several policies:

  • Per-command (accept(...)): rejects EVENTs over maxContentLength / maxEventTags / outside the createdAt bounds; rejects REQ/COUNT with too many filters or an over-long sub id; and clamps each filter's limit to maxLimit (substituting defaultLimit when none is given). Rejections use the invalid: machine-readable prefix.
  • Per-connection (the acceptMessage / acceptSubscription policy hooks): oversized frames get a NOTICE (maxMessageLength); a new subscription past maxSubscriptions is CLOSED with rate-limited:. These hooks exist because a policy otherwise can't see the raw frame or the live subscription count.
  • Advertised only: minPowDifficulty (enforce with a PoW policy), authRequired (use FullAuthPolicy), paymentRequired, restrictedWrites.

Serving NIP-11

Build a Nip11RelayInformation, fold in the same limits, and serve it at the relay root with the NIP-11 media type:

val info = Nip11RelayInformation(
    name = "My Relay",
    supported_nips = listOf("1", "11", "42", "45"),
    limitation = server.limits?.toNip11Limitation(),
)

// Ktor: answer GET / when the client asks for application/nostr+json
get("/") {
    call.respondText(info.toJson(), ContentType.parse(Nip11RelayInformation.CONTENT_TYPE))
}

Nip11RelayInformation.fromJson / toJson round-trip the document (null fields are omitted), and CONTENT_TYPE is application/nostr+json.

Approximate COUNT (NIP-45 HyperLogLog)

COUNT is answered by SessionBackend.countResult(ctx, filters) (and EventSource.countResult), which defaults to an exact count. To return a mergeable HyperLogLog estimate instead — for the six canonical NIP-45 queries (reaction/repost/quote/reply/comment/follower counts) — fold matching pubkeys into an HllBuilder and return its CountResult:

override suspend fun countResult(ctx: RequestContext, filters: List<Filter>): CountResult {
    val filter = filters.first()
    val hll = HyperLogLog.builderFor(filter) ?: return CountResult(count(ctx, filters))
    store.query(filter) { event -> hll.add(event.pubKey) }
    return hll.toCountResult() // count = estimate, approximate = true, hll = registers
}

The engine frames count/approximate/hll onto the wire (["COUNT", id, {"count":N,"hll":"<512-hex>"}]). Register arrays built by two relays over the same corpus merge with HyperLogLog.merge(...) into a deduplicated cross-relay estimate.

Observability

Both servers take an optional RelayServerListener and expose a live activeConnections gauge. Connections carry a stable, process-unique RelaySession.id (used to key the server's registry and the listener callbacks), so you can correlate the open/close of the same connection in logs and metrics:

val server = NostrServer(
    store = store,
    listener = object : RelayServerListener {
        override fun onConnect(connectionId: Long) = metrics.connections.inc()
        override fun onDisconnect(connectionId: Long) = metrics.connections.dec()
    },
)
server.activeConnections // Long, current count

onConnect/onDisconnect fire at most once per connection (a double close() is accounted once) and onDisconnect is also fired for any connections still open when the server itself is closed. Callbacks can run on any transport coroutine — keep them cheap and non-blocking.

Key Source Files

quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/
├── relay/server/                  # engine + entry points
│   ├── RelayServerBase.kt      # Shared engine (connect/serve/policy/registry)
│   ├── NostrServer.kt          # Storage-backed engine (IEventStore)
│   ├── EventSourceServer.kt    # Non-storage engine (search/redirector/computed)
│   ├── RelaySession.kt         # Per-connection handler (stable .id)
│   ├── ConnectionRegistry.kt   # Live-connection bookkeeping + gauge
│   ├── RelayServerListener.kt  # Connection open/close observability hook
│   ├── NegSessionRegistry.kt   # NIP-77 negentropy per-connection state
│   ├── backend/                # data plane (how REQ/COUNT/EVENT are answered)
│   │   ├── SessionBackend.kt   # Data-plane seam (LiveEventStore / EventSourceBackend)
│   │   ├── EventSource.kt      # Flow<Event> SPI for non-storage relays
│   │   ├── EventSourceBackend.kt # Adapts an EventSource to SessionBackend
│   │   ├── LiveEventStore.kt   # Reactive event streaming (storage SessionBackend)
│   │   └── IngestQueue.kt      # Group-commit EVENT writer
│   ├── policies/               # policy model + implementations
│   │   ├── IRelayPolicy.kt     # Policy interface + PolicyResult + onAuthenticated
│   │   ├── RelayLimits.kt      # Limits: enforced + NIP-11 limitation source of truth
│   │   ├── EmptyPolicy.kt      # Accept everything
│   │   ├── VerifyPolicy.kt     # Signature verification (default)
│   │   ├── FullAuthPolicy.kt   # NIP-42 auth required (override authorize to bridge)
│   │   ├── LimitsPolicy.kt     # Enforcement of RelayLimits (per-command + hooks)
│   │   └── PolicyStack.kt      # Chain multiple policies
│   └── inprocess/              # in-memory transport for tests
├── relay/commands/
│   ├── toRelay/Command.kt      # Command.fromJson / toJson
│   └── toClient/
│       ├── Message.kt          # Message.fromJson / toJson
│       └── MachineReadablePrefix.kt # Typed OK/CLOSED reason prefixes
├── store/
│   ├── IEventStore.kt          # Storage interface
│   └── sqlite/
│       ├── EventStore.kt       # Public SQLite store wrapper
│       ├── SQLiteEventStore.kt # Full implementation
│       └── IndexingStrategy.kt # Index configuration
├── relay/filters/
│   ├── Filter.kt               # NIP-01 subscription filters
│   └── FilterMatcher.kt        # Event-to-filter matching
└── ../nip50Search/
    └── SearchQuery.kt          # NIP-50 search-string parser