diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/IEventStore.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/IEventStore.kt index e70773c29..8a92f02c7 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/IEventStore.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/IEventStore.kt @@ -22,6 +22,7 @@ package com.vitorpamplona.quartz.nip01Core.store import com.vitorpamplona.quartz.nip01Core.core.Event import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter +import kotlinx.coroutines.flow.SharedFlow interface IEventStore : AutoCloseable { suspend fun insert(event: Event) @@ -56,5 +57,23 @@ interface IEventStore : AutoCloseable { suspend fun deleteExpiredEvents() + /** + * Stream of events the store accepted into durable storage. One + * emission per successfully inserted event, in commit order. + * Rejected inserts (expired, ephemeral, blocked by tombstone / + * vanish, NIP-01 supersession loser) emit nothing. + * + * Consumed by `EventStoreProjection` to maintain a live view — + * the projection itself replays NIP-01 supersession, NIP-09 + * deletion fan-out, NIP-62 vanish cascades, and NIP-40 expiration + * from these events, so stores don't need to publish removals. + * + * Out-of-band removals (`delete(id)`, `delete(filter)`, `clearDB()`, + * `deleteExpiredEvents()`) are not visible on this stream — they're + * maintenance operations and projections survive a missed mutation + * by re-seeding when their scope is restarted. + */ + val inserts: SharedFlow + override fun close() } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/projection/EventStoreProjection.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/projection/EventStoreProjection.kt new file mode 100644 index 000000000..3882b2db5 --- /dev/null +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/projection/EventStoreProjection.kt @@ -0,0 +1,401 @@ +/* + * 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.store.projection + +import com.vitorpamplona.quartz.nip01Core.core.AddressableEvent +import com.vitorpamplona.quartz.nip01Core.core.Event +import com.vitorpamplona.quartz.nip01Core.core.HexKey +import com.vitorpamplona.quartz.nip01Core.core.isAddressable +import com.vitorpamplona.quartz.nip01Core.core.isReplaceable +import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter +import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl +import com.vitorpamplona.quartz.nip01Core.store.IEventStore +import com.vitorpamplona.quartz.nip09Deletions.DeletionEvent +import com.vitorpamplona.quartz.nip40Expiration.expiration +import com.vitorpamplona.quartz.nip59Giftwrap.wraps.GiftWrapEvent +import com.vitorpamplona.quartz.nip62RequestToVanish.RequestToVanishEvent +import com.vitorpamplona.quartz.utils.TimeUtils +import kotlinx.coroutines.CompletableDeferred +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Job +import kotlinx.coroutines.delay +import kotlinx.coroutines.flow.MutableStateFlow +import kotlinx.coroutines.flow.StateFlow +import kotlinx.coroutines.flow.asStateFlow +import kotlinx.coroutines.launch +import kotlinx.coroutines.yield + +/** + * A reactive projection over any [IEventStore] for a fixed set of + * [filters]. Each visible event is wrapped in a [MutableStateFlow] so + * the UI can collect three different kinds of change with the right + * granularity: + * + * - **Membership** (events arriving or leaving) re-emits a brand new + * [List] from [items]. The list reference is stable while membership + * is unchanged. + * - **In-place replaceable / addressable update** (a new version of + * the same `kind:pubkey:dtag` arrives) updates the existing handle's + * [MutableStateFlow.value] without touching the list. Only collectors + * of that one handle re-render. The list ordering is *not* + * reshuffled when the new version has a later `created_at` — each + * slot remembers the sort key it was inserted with, so updates feel + * like pure value mutations. + * - **Removal** (NIP-09 deletion, NIP-62 vanish, NIP-40 expiration) + * drops the handle from the list. + * + * The seed is materialised by querying the store once at start, after + * which the projection is driven entirely by [IEventStore.inserts] and + * its own expiration ticker. The store is never re-queried on + * mutation, and the projection never asks the store to delete anything + * — it interprets incoming Nostr events itself: + * + * - **NIP-01 supersession.** New replaceable / addressable events + * replace prior ones for the same `kind:pubkey[:dtag]`. The + * NIP-01 lexical-id tiebreaker (`new.id < old.id` when + * `created_at` ties) is honoured. + * - **NIP-09 deletions.** A [DeletionEvent] removes any matching + * handle owned by the same author (for GiftWrap, the recipient). + * Cross-author deletions are inert. + * - **NIP-62 right to vanish.** A [RequestToVanishEvent] whose + * `shouldVanishFrom([relay])` is true drops every handle from the + * same author with `created_at < vanish.created_at`. + * - **NIP-40 expiration.** Events with a past `expiration` tag are + * rejected at insert time. A periodic ticker drops slots whose + * expiration has just lapsed; collectors see the slot disappear. + * + * That's the same set of rules the SQLite / FS stores enforce for + * durability. The duplication is by design — the store enforces them + * on disk so the file isn't corrupt; the projection enforces them in + * memory so the live view stays correct without a re-query per event. + * + * Limit handling: the initial query honours the filter `limit`, and + * we trim to the same cap when an insert pushes the list over. We do + * **not** refill from the store after a deletion — if a deletion + * leaves you under the limit, you stay under the limit until something + * new arrives. That tradeoff matches what callers were already getting + * from `LocalCache.observeEvents`. + * + * Out-of-band store mutations — `delete(id)`, `delete(filter)`, + * `clearDB()`, the periodic `deleteExpiredEvents()` sweep — are not + * visible on [IEventStore.inserts] and won't update an open + * projection. Re-open the projection (e.g. cancel and recreate the + * scope) to pick up an out-of-band change. + * + * Lifecycle: the projection runs a collector + an expiration ticker + * inside [scope]. Cancel the scope (or call [close]) when the screen + * using the projection goes away. + */ +class EventStoreProjection( + private val store: IEventStore, + private val filters: List, + private val relay: NormalizedRelayUrl?, + scope: CoroutineScope, + private val expirationTickMs: Long = 30_000L, + private val nowProvider: () -> Long = TimeUtils::now, +) : AutoCloseable { + private val _items = MutableStateFlow>>(emptyList()) + val items: StateFlow>> = _items.asStateFlow() + + /** Slots keyed by the *current* event id. Re-keyed when a replaceable / addressable handle takes a new version. */ + private val byId = HashMap>() + + /** + * Slots keyed by `kind:pubkey:dtag` (or `kind:pubkey:` for plain + * replaceables) for in-place updates and supersession lookups. + */ + private val byStableKey = HashMap>() + + /** + * Sorted view of the same slots. The comparator uses each slot's + * frozen sort key (the seed event's `created_at` + `id`), so + * supersession in-place updates never move a slot inside this set. + */ + private val ordered = sortedSetOf(slotComparator()) + + private val limit: Int? = filters.mapNotNull { it.limit }.maxOrNull() + + /** Set when the seed has been written to [items], so callers can suspend until the projection is hot. */ + val ready: CompletableDeferred = CompletableDeferred() + + private val collectorJob: Job = + scope.launch { + seed() + ready.complete(Unit) + store.inserts.collect { event -> apply(event) } + } + + private val expirationJob: Job = + scope.launch { + // Sleep first so the seed-time sweep covers the initial + // contents — see [seed]. + while (true) { + delay(expirationTickMs) + sweepExpired() + } + } + + private suspend fun seed() { + val initial = store.query(filters) + val now = nowProvider() + for (event in initial) { + // The store should already exclude expired rows from the + // result, but it doesn't hurt to skip them here too — + // covers any FS / in-memory store that hasn't run a sweep + // recently. + if (isExpiredAt(event, now)) continue + insertNew(event) + yield() + } + publish() + } + + private fun apply(event: Event) { + if (isExpiredAt(event, nowProvider())) return + + var changed = false + + // Apply NIP-09 / NIP-62 side effects of the event before we + // consider matching the event itself against the filter — a + // deletion event that arrives at the same instant as a + // matching event still removes its targets. + if (event is DeletionEvent) { + if (handleDeletion(event)) changed = true + } + if (event is RequestToVanishEvent && event.shouldVanishFrom(relay)) { + if (handleVanish(event)) changed = true + } + + if (filters.any { it.match(event) }) { + if (handleInsert(event)) changed = true + } + + if (changed) publish() + } + + /** + * Returns true if the event matches the filter and its arrival + * caused membership to change (a fresh slot was added). Returns + * false when the arrival was an in-place supersession update or + * was rejected by the NIP-01 tiebreaker. + */ + @Suppress("UNCHECKED_CAST") + private fun handleInsert(event: Event): Boolean { + val key = stableKey(event) + if (key != null) { + val existing = byStableKey[key] + if (existing != null) { + if (!supersedes(event, existing.flow.value)) return false + + // Same address, new winner. Rekey byId and update the + // handle's value in place — list reference stays the + // same; only the handle's collectors re-render. + val previousId = existing.flow.value.id + if (previousId != event.id) { + byId.remove(previousId) + byId[event.id] = existing + } + existing.flow.value = event as T + return false + } + } else if (byId.containsKey(event.id)) { + return false + } + + insertNew(event) + return true + } + + private fun handleDeletion(deletion: DeletionEvent): Boolean { + var changed = false + + // NIP-09: delete by id, but only if the deletion's author owns + // the target. For GiftWrap, the owner is the p-tag recipient. + for (id in deletion.deleteEventIds()) { + val slot = byId[id] ?: continue + val ev = slot.flow.value + val owner = ownerPubKey(ev) + if (owner == deletion.pubKey && removeSlot(slot)) changed = true + } + + // NIP-09: delete by address, only original author, only events + // with `created_at <= deletion.created_at`. + for (addr in deletion.deleteAddresses()) { + if (addr.pubKeyHex != deletion.pubKey) continue + val key = stableKey(addr.kind, addr.pubKeyHex, addr.dTag) ?: continue + val slot = byStableKey[key] ?: continue + if (slot.flow.value.createdAt <= deletion.createdAt) { + if (removeSlot(slot)) changed = true + } + } + + return changed + } + + private fun handleVanish(vanish: RequestToVanishEvent): Boolean { + var changed = false + // Snapshot first because removeSlot mutates byId. + val targets = + byId.values.filter { + val ev = it.flow.value + ownerPubKey(ev) == vanish.pubKey && ev.createdAt < vanish.createdAt + } + for (slot in targets) { + if (removeSlot(slot)) changed = true + } + return changed + } + + private fun sweepExpired() { + val now = nowProvider() + val targets = byId.values.filter { isExpiredAt(it.flow.value, now) } + if (targets.isEmpty()) return + var changed = false + for (slot in targets) { + if (removeSlot(slot)) changed = true + } + if (changed) publish() + } + + @Suppress("UNCHECKED_CAST") + private fun insertNew(event: Event) { + val slot = Slot(event as T) + byId[event.id] = slot + stableKey(event)?.let { byStableKey[it] = slot } + ordered.add(slot) + + val cap = limit ?: return + while (ordered.size > cap) { + val tail = ordered.last() + removeSlot(tail) + } + } + + private fun removeSlot(slot: Slot): Boolean { + val removed = byId.remove(slot.flow.value.id) != null + if (!removed) return false + ordered.remove(slot) + stableKey(slot.flow.value)?.let { key -> + // Defensive: only clear the stable-key map if this slot + // still owns it. Could be stale after an addressable rekey + // raced with another insert. + if (byStableKey[key] === slot) byStableKey.remove(key) + } + return true + } + + private fun publish() { + _items.value = ordered.map { it.flow } + } + + /** + * Stop tracking changes and clear internal state. Idempotent. The + * scope passed to the constructor keeps running; only this + * projection's collector + expiration jobs are cancelled. + */ + override fun close() { + collectorJob.cancel() + expirationJob.cancel() + ordered.clear() + byId.clear() + byStableKey.clear() + _items.value = emptyList() + } + + /** + * Internal slot. Each event added to the projection lives inside + * one of these for as long as it survives. The sort key is frozen + * at construction time — supersession in-place updates rewrite + * `flow.value` but never the sort key, so the ordering inside + * [ordered] is stable across updates. + */ + private class Slot( + initial: T, + ) { + val sortCreatedAt: Long = initial.createdAt + val sortId: HexKey = initial.id + val flow: MutableStateFlow = MutableStateFlow(initial) + } + + companion object { + /** + * The lookup key for replaceable / addressable supersession. + * `null` for regular events (which only collide on event id). + */ + private fun stableKey(event: Event): String? = stableKey(event.kind, event.pubKey, (event as? AddressableEvent)?.dTag()) + + private fun stableKey( + kind: Int, + pubKeyHex: HexKey, + dTag: String?, + ): String? = + when { + kind.isAddressable() -> "$kind:$pubKeyHex:${dTag ?: ""}" + kind.isReplaceable() -> "$kind:$pubKeyHex:" + else -> null + } + + /** + * NIP-01 supersession tiebreaker. The new event wins iff: + * - its `created_at` is strictly greater, OR + * - the timestamps tie and its `id` is lexically smaller. + * Otherwise the existing slot keeps its place. + */ + private fun supersedes( + new: Event, + existing: Event, + ): Boolean = + when { + new.createdAt > existing.createdAt -> true + new.createdAt < existing.createdAt -> false + else -> new.id < existing.id + } + + /** + * Owner pubkey for ownership checks (NIP-09 author match, + * NIP-62 vanish target). For GiftWrap the owner is the p-tag + * recipient; for everything else it's `event.pubKey`. + */ + private fun ownerPubKey(event: Event): HexKey = (event as? GiftWrapEvent)?.recipientPubKey() ?: event.pubKey + + private fun isExpiredAt( + event: Event, + now: Long, + ): Boolean { + val exp = event.expiration() ?: return false + return exp <= now + } + + /** + * created_at DESC, id ASC. The keys are snapshots taken at + * insertion time, so the ordering of a slot never changes + * after it joins the set. Distinct events have distinct ids, + * so no third-key tiebreak is needed. + */ + private fun slotComparator(): Comparator> = + Comparator { a, b -> + if (a === b) return@Comparator 0 + val byTime = b.sortCreatedAt.compareTo(a.sortCreatedAt) + if (byTime != 0) return@Comparator byTime + a.sortId.compareTo(b.sortId) + } + } +} diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/ChangeLogModule.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/ChangeLogModule.kt deleted file mode 100644 index 73d34cd95..000000000 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/ChangeLogModule.kt +++ /dev/null @@ -1,98 +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.store.sqlite - -import androidx.sqlite.SQLiteConnection -import com.vitorpamplona.quartz.nip01Core.core.HexKey - -/** - * Per-connection scratch space that records the id of every row that - * leaves `event_headers`, regardless of why it left. - * - * Why TEMP: the table and trigger are connection-scoped, so they only - * exist on the writer connection where they're installed. Readers never - * see them and never accumulate rows. The data also lives only for the - * lifetime of the connection — exactly the right scope for "ids removed - * since the last drain." - * - * The trigger fires after every delete on `event_headers`, including: - * - the supersession triggers in [ReplaceableModule] / [AddressableModule], - * - the cascade from `event_vanish` in [RightToVanishModule], - * - the explicit deletes in [DeletionRequestModule], - * - [ExpirationModule.deleteExpiredEvents], - * - manual `delete(filter)` / `delete(id)`, - * - `clearDB()`. - * - * The [SQLiteEventStore] drains this log around each unit of work and - * publishes the resulting ids in a [StoreChange]. - */ -class ChangeLogModule { - fun installOnWriter(db: SQLiteConnection) { - db.execSQL( - """ - CREATE TEMP TABLE IF NOT EXISTS event_change_log ( - seq INTEGER PRIMARY KEY AUTOINCREMENT, - id TEXT NOT NULL - ) - """.trimIndent(), - ) - - db.execSQL( - """ - CREATE TEMP TRIGGER IF NOT EXISTS event_change_log_on_delete - AFTER DELETE ON event_headers - FOR EACH ROW - BEGIN - INSERT INTO event_change_log (id) VALUES (OLD.id); - END - """.trimIndent(), - ) - } - - /** - * Reads every id currently logged and clears the log. - * - * Must be called on the writer connection that owns the temp table. - * Callers are expected to be inside the writer mutex (drain happens - * inside or right after the `useWriter { ... }` block); a single - * read-and-clear pair is therefore atomic from the writer's view. - */ - fun drain(db: SQLiteConnection): List { - val ids = ArrayList() - db.prepare("SELECT id FROM event_change_log ORDER BY seq").use { stmt -> - while (stmt.step()) { - ids.add(stmt.getText(0)) - } - } - if (ids.isNotEmpty()) { - db.execSQL("DELETE FROM event_change_log") - } - return ids - } - - /** - * Drops everything from the log without returning it. Used after a - * rollback so the next successful unit of work starts clean. - */ - fun reset(db: SQLiteConnection) { - db.execSQL("DELETE FROM event_change_log") - } -} diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/EventStore.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/EventStore.kt index f17be61ee..092f10bcb 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/EventStore.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/EventStore.kt @@ -26,11 +26,12 @@ import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl import com.vitorpamplona.quartz.nip01Core.relay.normalizer.normalizeRelayUrl import com.vitorpamplona.quartz.nip01Core.store.IEventStore +import com.vitorpamplona.quartz.nip01Core.store.projection.EventStoreProjection import kotlinx.coroutines.CoroutineScope class EventStore( dbName: String? = "events.db", - relay: NormalizedRelayUrl? = "wss://quartz.local/".normalizeRelayUrl(), + val relay: NormalizedRelayUrl? = "wss://quartz.local/".normalizeRelayUrl(), val indexStrategy: IndexingStrategy = DefaultIndexingStrategy(), ) : IEventStore { val store = SQLiteEventStore(BundledSQLiteDriver(), dbName, relay, indexStrategy) @@ -67,20 +68,17 @@ class EventStore( override suspend fun deleteExpiredEvents() = store.deleteExpiredEvents() - /** - * Stream of mutations committed to the store. See [SQLiteEventStore.changes]. - */ - val changes get() = store.changes + override val inserts get() = store.inserts /** - * Open a reactive [EventStoreProjection] over the store for - * [filters]. The projection runs inside [scope]; cancel the scope - * (or call [EventStoreProjection.close]) to release it. + * Open a reactive [EventStoreProjection] over this store with + * NIP-62 vanish scoping bound to the store's [relay]. Cancel + * [scope] (or call [EventStoreProjection.close]) to release it. */ fun observe( filters: List, scope: CoroutineScope, - ): EventStoreProjection = EventStoreProjection(store, filters, scope) + ): EventStoreProjection = EventStoreProjection(this, filters, relay, scope) fun observe( filter: Filter, diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/EventStoreProjection.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/EventStoreProjection.kt deleted file mode 100644 index f655d6dfd..000000000 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/EventStoreProjection.kt +++ /dev/null @@ -1,247 +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.store.sqlite - -import com.vitorpamplona.quartz.nip01Core.core.AddressableEvent -import com.vitorpamplona.quartz.nip01Core.core.Event -import com.vitorpamplona.quartz.nip01Core.core.HexKey -import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter -import kotlinx.coroutines.CompletableDeferred -import kotlinx.coroutines.CoroutineScope -import kotlinx.coroutines.Job -import kotlinx.coroutines.flow.MutableStateFlow -import kotlinx.coroutines.flow.StateFlow -import kotlinx.coroutines.flow.asStateFlow -import kotlinx.coroutines.launch -import kotlinx.coroutines.yield - -/** - * A reactive projection over the [SQLiteEventStore] for a fixed set of - * [filters]. Each visible event is wrapped in a [MutableStateFlow] so - * the UI can collect three different kinds of change with the right - * granularity: - * - * - **Membership** (events arriving or leaving) re-emits a brand new - * [List] from [items]. The list reference is stable while membership - * is unchanged. - * - **In-place addressable update** (a new version of the same - * `kind:pubkey:dtag` arrives) updates the existing handle's - * [MutableStateFlow.value] without touching the list. Only collectors - * of that one handle re-render. The list ordering is *not* reshuffled - * when the new version has a later `created_at` — each slot remembers - * the sort key it was inserted with, so addressable updates feel like - * pure value mutations. - * - **Removal** (NIP-09 deletion, NIP-62 vanish, NIP-40 expiration, or - * a non-addressable event being explicitly deleted) drops the handle - * from the list. - * - * The seed is materialised by running the filters against the store - * once at start. After that, [items] is driven entirely by - * [SQLiteEventStore.changes]; the database is not re-queried on every - * mutation. - * - * Limit handling matches the existing in-memory observables in - * `commons/observables`: the initial query honours the filter `limit`, - * and we trim the list to the same cap when an insert pushes it over. - * We do **not** refill from the DB after a deletion — if a deletion - * leaves you under the limit, you stay under the limit until something - * new arrives. That tradeoff keeps the projection allocation-free per - * mutation and matches what callers were already getting from - * `LocalCache.observeEvents`. - * - * Lifecycle: the projection runs a single coroutine in [scope]. Cancel - * the scope (or call [close]) when the screen using the projection - * goes away. There is no shared state between projections; each one - * keeps its own indexes. - */ -class EventStoreProjection( - private val store: SQLiteEventStore, - private val filters: List, - scope: CoroutineScope, -) : AutoCloseable { - private val _items = MutableStateFlow>>(emptyList()) - val items: StateFlow>> = _items.asStateFlow() - - /** Slots keyed by the *current* event id. Re-keyed when an addressable handle takes a new version. */ - private val byId = HashMap>() - - /** Slots keyed by `kind:pubkey:dtag` for in-place addressable updates. */ - private val byAddress = HashMap>() - - /** - * Sorted view of the same slots. The comparator uses each slot's - * frozen sort key (the seed event's `created_at` + `id`), so an - * addressable update never moves a slot inside this set. - */ - private val ordered = sortedSetOf(slotComparator()) - - private val limit: Int? = filters.mapNotNull { it.limit }.maxOrNull() - - /** Set when the seed has been written to [items], so callers can suspend until the projection is hot. */ - val ready: CompletableDeferred = CompletableDeferred() - - private val job: Job = - scope.launch { - seed() - ready.complete(Unit) - store.changes.collect { change -> apply(change) } - } - - private suspend fun seed() { - val initial = store.query(filters) - for (event in initial) { - insertNew(event) - // Cooperate with cancellation on very large seeds. - yield() - } - publish() - } - - private fun apply(change: StoreChange) { - var changed = false - - // Removals first. A replaceable / addressable supersession - // arrives as `inserted = [new]` plus `removedIds = [oldId]`. - // When the new event is addressable and an existing slot - // already maps that address, [handleInsert] will rekey - // `byId` from the old id to the new id before this loop sees - // the old id — so the lookup here is a no-op for that case - // and the slot stays in place. For non-addressable - // replaceables (kinds 0/3/10000-19999) the old event has a - // different id, no address index, and we genuinely drop it. - for (event in change.inserted) { - if (handleInsert(event)) changed = true - } - for (id in change.removedIds) { - if (handleRemove(id)) changed = true - } - - if (changed) publish() - } - - @Suppress("UNCHECKED_CAST") - private fun handleInsert(event: Event): Boolean { - if (filters.none { it.match(event) }) return false - - if (event is AddressableEvent) { - val key = event.addressTag() - val existing = byAddress[key] - if (existing != null) { - // Same address, new version. Rekey byId from the - // previous event id to the new one and update the - // handle's value in place — list reference does not - // change, only the handle's collectors re-render. - val previousId = existing.flow.value.id - if (previousId != event.id) { - byId.remove(previousId) - byId[event.id] = existing - } - existing.flow.value = event as T - return false - } - } else if (byId.containsKey(event.id)) { - return false - } - - insertNew(event) - return true - } - - private fun handleRemove(id: HexKey): Boolean { - val slot = byId.remove(id) ?: return false - ordered.remove(slot) - val ev = slot.flow.value - if (ev is AddressableEvent) { - val addr = ev.addressTag() - // Only clear the address index if this slot still owns it. - if (byAddress[addr] === slot) byAddress.remove(addr) - } - return true - } - - @Suppress("UNCHECKED_CAST") - private fun insertNew(event: Event) { - val slot = Slot(event as T) - byId[event.id] = slot - if (event is AddressableEvent) byAddress[event.addressTag()] = slot - ordered.add(slot) - - val cap = limit ?: return - while (ordered.size > cap) { - val tail = ordered.last() - ordered.remove(tail) - val tailEvent = tail.flow.value - byId.remove(tailEvent.id) - if (tailEvent is AddressableEvent) { - val addr = tailEvent.addressTag() - if (byAddress[addr] === tail) byAddress.remove(addr) - } - } - } - - private fun publish() { - _items.value = ordered.map { it.flow } - } - - /** - * Stop tracking changes and clear internal state. Idempotent. The - * scope passed to the constructor keeps running; only this - * projection's collector job is cancelled. - */ - override fun close() { - job.cancel() - ordered.clear() - byId.clear() - byAddress.clear() - _items.value = emptyList() - } - - /** - * Internal slot. Each event added to the projection lives inside - * one of these for as long as it survives. The sort key is frozen - * at construction time — addressable in-place updates rewrite - * `flow.value` but never the sort key, so the ordering inside - * [ordered] is stable across updates. - */ - private class Slot( - initial: T, - ) { - val sortCreatedAt: Long = initial.createdAt - val sortId: HexKey = initial.id - val flow: MutableStateFlow = MutableStateFlow(initial) - } - - companion object { - /** - * created_at DESC, id ASC. The keys are snapshots taken at - * insertion time, so the ordering of a slot never changes - * after it joins the set. Distinct events have distinct ids, - * so no third-key tiebreak is needed. - */ - private fun slotComparator(): Comparator> = - Comparator { a, b -> - if (a === b) return@Comparator 0 - val byTime = b.sortCreatedAt.compareTo(a.sortCreatedAt) - if (byTime != 0) return@Comparator byTime - a.sortId.compareTo(b.sortId) - } - } -} diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/SQLiteEventStore.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/SQLiteEventStore.kt index 83c95d531..9b3687baf 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/SQLiteEventStore.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/SQLiteEventStore.kt @@ -66,22 +66,18 @@ class SQLiteEventStore( val deletionModule = DeletionRequestModule(seedModule::hasher) val expirationModule = ExpirationModule() val rightToVanishModule = RightToVanishModule(seedModule::hasher) - val changeLogModule = ChangeLogModule() /** - * Stream of mutations committed to the store. One [StoreChange] per - * successful unit of work; rejected inserts and rolled-back - * transactions are not emitted. Subscribers see updates in commit - * order. See [EventStoreProjection] for a high-level reactive list - * that consumes this stream. + * Stream of events the store accepted into durable storage. See + * [IEventStore.inserts] for the contract. */ - private val _changes = - MutableSharedFlow( + private val _inserts = + MutableSharedFlow( replay = 0, extraBufferCapacity = 256, onBufferOverflow = BufferOverflow.SUSPEND, ) - val changes: SharedFlow = _changes.asSharedFlow() + val inserts: SharedFlow = _inserts.asSharedFlow() val queryBuilder = QueryBuilder( @@ -139,12 +135,6 @@ class SQLiteEventStore( setUserVersion(this, DATABASE_VERSION) } } - - // After the schema is in place (fresh or already-current - // DBs both reach this point), wire the change log onto - // the writer connection. TEMP table + TEMP trigger are - // connection-scoped, so this runs on the writer only. - changeLogModule.installOnWriter(db) }, ) } @@ -186,23 +176,10 @@ class SQLiteEventStore( } } - suspend fun clearDB() { + suspend fun clearDB() = pool.useWriter { db -> modules.reversed().forEach { it.deleteAll(db) } - // The deleteAll cascade fires the change-log trigger for - // every removed event_header row. Drain and publish so - // open projections drop everything they were holding. - publishChange(emptyList(), changeLogModule.drain(db)) } - } - - private suspend fun publishChange( - inserted: List, - removedIds: List, - ) { - if (inserted.isEmpty() && removedIds.isEmpty()) return - _changes.emit(StoreChange(inserted, removedIds)) - } suspend fun vacuum() = pool.useWriter { db -> @@ -237,12 +214,10 @@ class SQLiteEventStore( db.transaction { innerInsertEvent(event, this) } - // The transaction either committed or threw. On commit, the - // change log holds ids superseded by replaceable / addressable - // triggers and any NIP-09 / NIP-62 cascades fired by this - // event's content. On rollback the temp table was rolled back - // with the transaction, so drain returns empty. - publishChange(listOf(event), changeLogModule.drain(db)) + // The transaction either committed or threw. On commit the + // event is durable; emit so subscribed projections see it. + // On rollback the throw propagates and we never reach this. + _inserts.emit(event) } } @@ -268,7 +243,9 @@ class SQLiteEventStore( body() } } - publishChange(txn.accepted, changeLogModule.drain(db)) + // Emit each accepted event after the batch commits — the + // projection sees them in the same order they were inserted. + for (e in txn.accepted) _inserts.emit(e) } } @@ -308,31 +285,17 @@ class SQLiteEventStore( suspend fun count(filters: List): Int = pool.useReader { queryBuilder.count(filters, it) } - suspend fun delete(filter: Filter) = - pool.useWriter { db -> - queryBuilder.delete(filter, db) - publishChange(emptyList(), changeLogModule.drain(db)) - } + suspend fun delete(filter: Filter) = pool.useWriter { queryBuilder.delete(filter, it) } - suspend fun delete(filters: List) = - pool.useWriter { db -> - queryBuilder.delete(filters, db) - publishChange(emptyList(), changeLogModule.drain(db)) - } + suspend fun delete(filters: List) = pool.useWriter { queryBuilder.delete(filters, it) } suspend fun delete(id: HexKey): Int = pool.useWriter { db -> db.execSQL("DELETE FROM event_headers WHERE id = ?", arrayOf(id)) - val count = db.changes() - publishChange(emptyList(), changeLogModule.drain(db)) - count + db.changes() } - suspend fun deleteExpiredEvents() = - pool.useWriter { db -> - expirationModule.deleteExpiredEvents(db) - publishChange(emptyList(), changeLogModule.drain(db)) - } + suspend fun deleteExpiredEvents() = pool.useWriter { expirationModule.deleteExpiredEvents(it) } fun close() = pool.close() } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/StoreChange.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/StoreChange.kt deleted file mode 100644 index ef9742fb8..000000000 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/StoreChange.kt +++ /dev/null @@ -1,57 +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.store.sqlite - -import com.vitorpamplona.quartz.nip01Core.core.Event -import com.vitorpamplona.quartz.nip01Core.core.HexKey - -/** - * One atomic batch of mutations that occurred inside the event store. - * - * The store emits exactly one [StoreChange] per successfully committed - * unit of work — `insertEvent`, a `transaction { ... }` block, - * `delete(...)`, `deleteExpiredEvents()`, or `clearDB()`. Empty changes - * (e.g. an insert that was rejected by a trigger) are not emitted. - * - * `removedIds` covers every row that left `event_headers` during the - * unit of work, regardless of cause: the supersession triggers for - * replaceable / addressable events, NIP-09 deletion fan-out, NIP-62 - * vanish cascades, NIP-40 expiration sweeps, and direct `delete(...)` - * calls. They are captured by an `AFTER DELETE` trigger that writes - * `OLD.id` into a per-connection TEMP table. - * - * `inserted` carries the events that survived the writer and are now - * in the database. A replaceable / addressable supersession therefore - * appears as a single change with `inserted = [new]` and - * `removedIds = [oldId]`. - */ -data class StoreChange( - val inserted: List, - val removedIds: List, -) { - fun isEmpty() = inserted.isEmpty() && removedIds.isEmpty() - - fun isNotEmpty() = !isEmpty() - - companion object { - val EMPTY = StoreChange(emptyList(), emptyList()) - } -} diff --git a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/EventStoreProjectionTest.kt b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/store/projection/EventStoreProjectionTest.kt similarity index 70% rename from quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/EventStoreProjectionTest.kt rename to quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/store/projection/EventStoreProjectionTest.kt index f9effa0ac..7644f1d31 100644 --- a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/EventStoreProjectionTest.kt +++ b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/store/projection/EventStoreProjectionTest.kt @@ -18,13 +18,14 @@ * 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.store.sqlite +package com.vitorpamplona.quartz.nip01Core.store.projection import com.vitorpamplona.quartz.nip01Core.core.Event import com.vitorpamplona.quartz.nip01Core.metadata.MetadataEvent import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter import com.vitorpamplona.quartz.nip01Core.relay.normalizer.normalizeRelayUrl import com.vitorpamplona.quartz.nip01Core.signers.NostrSignerSync +import com.vitorpamplona.quartz.nip01Core.store.sqlite.EventStore import com.vitorpamplona.quartz.nip09Deletions.DeletionEvent import com.vitorpamplona.quartz.nip10Notes.TextNoteEvent import com.vitorpamplona.quartz.nip23LongContent.LongTextNoteEvent @@ -50,6 +51,7 @@ import kotlin.test.assertTrue class EventStoreProjectionTest { private val signer = NostrSignerSync() + private val otherSigner = NostrSignerSync() private lateinit var store: EventStore private lateinit var scope: CoroutineScope @@ -66,13 +68,6 @@ class EventStoreProjectionTest { store.close() } - /** - * Wait until [items][EventStoreProjection.items] reaches a state - * for which [predicate] is true. We poll the StateFlow rather than - * collect because the projection only re-emits on membership - * change — in-place addressable updates intentionally don't move - * the list reference. - */ private suspend fun EventStoreProjection.awaitItems( timeoutMs: Long = 5_000, predicate: (List>) -> Boolean, @@ -81,11 +76,6 @@ class EventStoreProjectionTest { items.first { predicate(it) } } - /** - * Wait until [flow] reaches [expected]. Used to observe an - * in-place addressable update where the list reference doesn't - * change but the slot's value does. - */ private suspend fun awaitFlow( flow: MutableStateFlow, timeoutMs: Long = 5_000, @@ -108,7 +98,6 @@ class EventStoreProjectionTest { val items = projection.items.value assertEquals(2, items.size) - // Sorted by created_at DESC. assertEquals(b.id, items[0].value.id) assertEquals(a.id, items[1].value.id) projection.close() @@ -131,12 +120,11 @@ class EventStoreProjectionTest { val after = projection.awaitItems { it.size == 2 } assertNotSame(before, after, "insert must produce a new list reference") assertEquals(b.id, after[0].value.id) - assertEquals(a.id, after[1].value.id) projection.close() } @Test - fun insertingNonMatchingEventDoesNotChangeList() = + fun nonMatchingInsertDoesNotChangeList() = runBlocking { val text = signer.sign(TextNoteEvent.build("a", createdAt = 100)) store.insert(text) @@ -145,11 +133,9 @@ class EventStoreProjectionTest { projection.ready.await() val seed = projection.items.value - // Metadata kind doesn't match the filter. val meta = signer.sign(MetadataEvent.createNew("Vitor", createdAt = 200)) store.insert(meta) - // Give the projection time to process the change. delay(150) assertSame(seed, projection.items.value) projection.close() @@ -176,11 +162,8 @@ class EventStoreProjectionTest { val v2 = signer.sign(MetadataEvent.createNew("v2", createdAt = time + 1)) store.insert(v2) - // The slot's flow updates... awaitFlow(slot) { it.id == v2.id } - - // ...but the list reference is the SAME, and the slot is the SAME instance. - assertSame(seedList, projection.items.value, "addressable replace must not change list reference") + assertSame(seedList, projection.items.value, "replaceable update must not change list reference") assertSame(slot, projection.items.value[0]) projection.close() } @@ -203,16 +186,54 @@ class EventStoreProjectionTest { ) projection.ready.await() val seedList = projection.items.value - assertEquals(1, seedList.size) val slot = seedList[0] - assertEquals(v1.id, slot.value.id) val v2 = signer.sign(LongTextNoteEvent.build("blog v2", "title", dTag = "blog", createdAt = time + 1)) store.insert(v2) awaitFlow(slot) { it.id == v2.id } assertSame(seedList, projection.items.value, "addressable update must not change list reference") - assertSame(slot, projection.items.value[0]) + projection.close() + } + + /** + * Out-of-order arrival: the projection sees the *newer* version + * first (e.g. the relay sent v2 first), then v1. The projection + * must keep v2 in place and not regress to v1 — that's the NIP-01 + * supersession contract from the projection's side, since the + * store would have rejected v1 anyway. + */ + @Test + fun olderReplaceableArrivingAfterNewerIsRejected() = + runBlocking { + val time = TimeUtils.now() + val v1 = signer.sign(MetadataEvent.createNew("v1", createdAt = time)) + val v2 = signer.sign(MetadataEvent.createNew("v2", createdAt = time + 5)) + + // Seed the projection with v2 before v1 even hits the store + // — by inserting v2 first. + store.insert(v2) + + val projection = + store.observe( + Filter(kinds = listOf(MetadataEvent.KIND), authors = listOf(v1.pubKey)), + scope, + ) + projection.ready.await() + val slot = projection.items.value[0] + assertEquals(v2.id, slot.value.id) + + // The store rejects v1 because v2 already won; the + // projection therefore never sees v1 on the inserts + // stream. The slot must still hold v2. + try { + store.insert(v1) + } catch (_: Throwable) { + // expected — store enforces the same rule + } + + delay(150) + assertEquals(v2.id, slot.value.id) projection.close() } @@ -224,7 +245,7 @@ class EventStoreProjectionTest { store.insert(a) store.insert(b) - val projection = store.observe(Filter(kinds = listOf(TextNoteEvent.KIND)), scope) + val projection = store.observe(Filter(kinds = listOf(TextNoteEvent.KIND)), scope) projection.ready.await() assertEquals(2, projection.items.value.size) @@ -236,8 +257,37 @@ class EventStoreProjectionTest { projection.close() } + /** + * NIP-09 cross-author deletions are inert. A different signer + * publishing a kind-5 targeting `a` must not drop the slot. + */ @Test - fun nip62VanishRemovesAllAuthorsEvents() = + fun nip09CrossAuthorDeletionIsInert() = + runBlocking { + val a = signer.sign(TextNoteEvent.build("a", createdAt = 100)) + store.insert(a) + + val projection = store.observe(Filter(kinds = listOf(TextNoteEvent.KIND)), scope) + projection.ready.await() + val seed = projection.items.value + assertEquals(1, seed.size) + + val foreignDeletion = otherSigner.sign(DeletionEvent.build(listOf(a))) + store.insert(foreignDeletion) + + // Give the projection time to process the event. + delay(150) + assertSame(seed, projection.items.value) + assertEquals( + a.id, + projection.items.value[0] + .value.id, + ) + projection.close() + } + + @Test + fun nip62VanishRemovesAuthorEvents() = runBlocking { val time = TimeUtils.now() val a = signer.sign(TextNoteEvent.build("a", createdAt = time)) @@ -245,7 +295,7 @@ class EventStoreProjectionTest { store.insert(a) store.insert(b) - val projection = store.observe(Filter(kinds = listOf(TextNoteEvent.KIND)), scope) + val projection = store.observe(Filter(kinds = listOf(TextNoteEvent.KIND)), scope) projection.ready.await() assertEquals(2, projection.items.value.size) @@ -263,8 +313,42 @@ class EventStoreProjectionTest { projection.close() } + /** + * NIP-62 only removes events from the same author. A vanish from + * a different author must not touch slots owned by [signer]. + */ @Test - fun nip40ExpirationRemovesSlotOnSweep() = + fun nip62OtherAuthorVanishLeavesEventsAlone() = + runBlocking { + val time = TimeUtils.now() + val a = signer.sign(TextNoteEvent.build("a", createdAt = time)) + store.insert(a) + + val projection = store.observe(Filter(kinds = listOf(TextNoteEvent.KIND)), scope) + projection.ready.await() + val seed = projection.items.value + + val foreignVanish = + otherSigner.sign( + RequestToVanishEvent.build( + "wss://quartz.local".normalizeRelayUrl(), + createdAt = time + 2, + ), + ) + store.insert(foreignVanish) + + delay(150) + assertSame(seed, projection.items.value) + projection.close() + } + + /** + * NIP-40 per-projection ticker: a slot whose `expiration` lapses + * after the projection has loaded should be dropped on the next + * tick, even though the store hasn't run its sweep yet. + */ + @Test + fun nip40ExpirationDroppedByTicker() = runBlocking { val time = TimeUtils.now() val safe = signer.sign(TextNoteEvent.build("safe", createdAt = time) { expiration(time + 100) }) @@ -272,36 +356,26 @@ class EventStoreProjectionTest { store.insert(safe) store.insert(short) - val projection = store.observe(Filter(kinds = listOf(TextNoteEvent.KIND)), scope) + // Drive the ticker frequently so the test doesn't sit idle. + val projection = + EventStoreProjection( + store, + listOf(Filter(kinds = listOf(TextNoteEvent.KIND))), + relay = null, + scope = scope, + expirationTickMs = 100, + ) projection.ready.await() assertEquals(2, projection.items.value.size) - // Wait for the expiration to lapse, then run the sweep. + // Wait past the short expiration. delay(2000) - store.deleteExpiredEvents() - val after = projection.awaitItems { it.size == 1 } + val after = projection.awaitItems(timeoutMs = 5_000) { it.size == 1 } assertEquals(safe.id, after[0].value.id) projection.close() } - @Test - fun manualDeleteByIdRemovesSlot() = - runBlocking { - val a = signer.sign(TextNoteEvent.build("a", createdAt = 100)) - store.insert(a) - - val projection = store.observe(Filter(kinds = listOf(TextNoteEvent.KIND)), scope) - projection.ready.await() - assertEquals(1, projection.items.value.size) - - store.store.delete(a.id) - - val after = projection.awaitItems { it.isEmpty() } - assertTrue(after.isEmpty()) - projection.close() - } - @Test fun limitIsEnforcedOnInsertOverflow() = runBlocking { @@ -318,7 +392,6 @@ class EventStoreProjectionTest { projection.ready.await() assertEquals(2, projection.items.value.size) - // Newer event arrives; it should push the oldest out. val c = signer.sign(TextNoteEvent.build("c", createdAt = 300)) store.insert(c) @@ -339,7 +412,6 @@ class EventStoreProjectionTest { projection.ready.await() projection.close() - // Subsequent inserts must not surface in the (now empty) projection. store.insert(signer.sign(TextNoteEvent.build("b", createdAt = 200))) delay(150) assertTrue(projection.items.value.isEmpty()) diff --git a/quartz/src/jvmMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/fs/FsEventStore.kt b/quartz/src/jvmMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/fs/FsEventStore.kt index 8da5d14b5..eec9265f1 100644 --- a/quartz/src/jvmMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/fs/FsEventStore.kt +++ b/quartz/src/jvmMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/fs/FsEventStore.kt @@ -36,6 +36,10 @@ import com.vitorpamplona.quartz.nip09Deletions.DeletionEvent import com.vitorpamplona.quartz.nip40Expiration.expiration import com.vitorpamplona.quartz.nip62RequestToVanish.RequestToVanishEvent import com.vitorpamplona.quartz.utils.TimeUtils +import kotlinx.coroutines.channels.BufferOverflow +import kotlinx.coroutines.flow.MutableSharedFlow +import kotlinx.coroutines.flow.SharedFlow +import kotlinx.coroutines.flow.asSharedFlow import java.nio.file.FileAlreadyExistsException import java.nio.file.Files import java.nio.file.Path @@ -95,19 +99,40 @@ open class FsEventStore( planner = FsQueryPlanner(layout, hasher) } + /** + * Stream of events newly persisted by this store. See + * [IEventStore.inserts] for the contract. Events that are no-op + * idempotent retries (canonical already on disk from a previous + * call) are not re-emitted; the original insert already published. + */ + private val _inserts = + MutableSharedFlow( + replay = 0, + extraBufferCapacity = 256, + onBufferOverflow = BufferOverflow.SUSPEND, + ) + override val inserts: SharedFlow = _inserts.asSharedFlow() + // ------------------------------------------------------------------ // Insert // ------------------------------------------------------------------ - override suspend fun insert(event: Event) = - lockManager.withWriteLock { - insertLocked(event) - } + override suspend fun insert(event: Event) { + val accepted = lockManager.withWriteLock { insertLocked(event) } + if (accepted) _inserts.emit(event) + } - private fun insertLocked(event: Event) { - if (event.kind.isEphemeral()) return - if (isAlreadyExpired(event)) return - if (isBlockedByTombstone(event)) return + /** + * Inserts the event under the write lock and returns true iff this + * call is the one that newly persisted it. Returns false for the + * no-op paths (ephemeral, expired, blocked, supersession loser, or + * canonical already on disk from a prior call) so callers can + * decide whether to publish on [inserts]. + */ + private fun insertLocked(event: Event): Boolean { + if (event.kind.isEphemeral()) return false + if (isAlreadyExpired(event)) return false + if (isBlockedByTombstone(event)) return false val slot = slots.slotPathFor(event) val existingSlot = slot?.let { slots.readSlot(it) } @@ -121,7 +146,7 @@ open class FsEventStore( // (later createdAt, or same createdAt with the lexically smaller // id). Matches the ReplaceableModule / AddressableModule trigger // condition in SQLite. - return + return false } val canonical = layout.canonical(event.id) @@ -134,7 +159,7 @@ open class FsEventStore( } if (event is DeletionEvent) processDeletion(event, canonical) if (event is RequestToVanishEvent) processVanish(event, canonical) - return + return false } Files.createDirectories(canonical.parent) @@ -147,7 +172,7 @@ open class FsEventStore( // A concurrent writer won the race. Canonical is immutable // so the other copy is equivalent — drop our tmp and return. Files.deleteIfExists(tmp) - return + return false } Files.setLastModifiedTime(canonical, FileTime.from(event.createdAt, TimeUnit.SECONDS)) indexer.link(event, canonical) @@ -156,6 +181,7 @@ open class FsEventStore( } if (event is DeletionEvent) processDeletion(event, canonical) if (event is RequestToVanishEvent) processVanish(event, canonical) + return true } catch (t: Throwable) { Files.deleteIfExists(tmp) throw t @@ -263,14 +289,21 @@ open class FsEventStore( } } - override suspend fun transaction(body: IEventStore.ITransaction.() -> Unit) = + override suspend fun transaction(body: IEventStore.ITransaction.() -> Unit) { + val accepted = ArrayList() lockManager.withWriteLock { val txn = object : IEventStore.ITransaction { - override fun insert(event: Event) = insertLocked(event) + override fun insert(event: Event) { + if (insertLocked(event)) accepted.add(event) + } } txn.body() } + // Emit each accepted event in order after the lock releases — + // mirrors SQLiteEventStore.transaction. + for (e in accepted) _inserts.emit(e) + } // ------------------------------------------------------------------ // Query