feat(quartz): replace items + ready with sealed ProjectionState
EventStoreProjection now exposes `state: StateFlow<ProjectionState<T>>`
instead of separate `items` + `ready` signals. ProjectionState is a
sealed interface with two cases:
data object Loading // initial state
data class Loaded(items: List<MutableStateFlow<T>>) // post-seed view
This lets the UI distinguish "still seeding" from "seeded but empty"
— the previous `emptyList()` initial value conflated the two. Future
extension to `Failed(throwable)` is a single case addition.
The `ready: CompletableDeferred<Unit>` signal is gone; callers use
`state.first { it is Loaded }` (or the test `awaitReady()` helper).
Test fixtures gain two private extensions on EventStoreProjection<T>:
- val items: List<MutableStateFlow<T>> — terse access to the loaded
list (returns empty if still seeding).
- suspend awaitReady() — replaces ready.await().
awaitItems(predicate) now matches against the Loaded state. 17/17
projection + 240/240 store + interner tests pass.
https://claude.ai/code/session_01Jny85MTu1ynKgFBgysfWu5
This commit is contained in:
+30
-10
@@ -33,7 +33,6 @@ import com.vitorpamplona.quartz.nip40Expiration.isExpirationBefore
|
||||
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.flow.MutableStateFlow
|
||||
@@ -43,6 +42,23 @@ import kotlinx.coroutines.flow.onSubscription
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlinx.coroutines.yield
|
||||
|
||||
/**
|
||||
* Lifecycle state of an [EventStoreProjection]'s [state][EventStoreProjection.state].
|
||||
*
|
||||
* - [Loading] is the initial state before the seed query completes.
|
||||
* The UI should show a spinner / skeleton here.
|
||||
* - [Loaded] holds the current deduped, sorted list of slots. Every
|
||||
* membership change publishes a fresh [Loaded] instance; in-place
|
||||
* addressable updates do *not* publish (the slots' own flows do).
|
||||
*/
|
||||
sealed interface ProjectionState<out T : Event> {
|
||||
data object Loading : ProjectionState<Nothing>
|
||||
|
||||
data class Loaded<T : Event>(
|
||||
val items: List<MutableStateFlow<T>>,
|
||||
) : ProjectionState<T>
|
||||
}
|
||||
|
||||
/**
|
||||
* A reactive projection over an [ObservableEventStore] for a fixed
|
||||
* set of [filters]. Each visible event is wrapped in a
|
||||
@@ -119,8 +135,16 @@ class EventStoreProjection<T : Event>(
|
||||
scope: CoroutineScope,
|
||||
private val nowProvider: () -> Long = TimeUtils::now,
|
||||
) : AutoCloseable {
|
||||
private val _items = MutableStateFlow<List<MutableStateFlow<T>>>(emptyList())
|
||||
val items: StateFlow<List<MutableStateFlow<T>>> = _items.asStateFlow()
|
||||
private val _state = MutableStateFlow<ProjectionState<T>>(ProjectionState.Loading)
|
||||
|
||||
/**
|
||||
* Current projection state. Starts as [ProjectionState.Loading]
|
||||
* until the seed query completes; thereafter every membership
|
||||
* change publishes a fresh [ProjectionState.Loaded] with the new
|
||||
* list. Distinguishing `Loading` from `Loaded(emptyList())` lets
|
||||
* the UI tell "still seeding" apart from "seeded but no matches".
|
||||
*/
|
||||
val state: StateFlow<ProjectionState<T>> = _state.asStateFlow()
|
||||
|
||||
/** Slots keyed by the *current* event id. Re-keyed when a replaceable / addressable handle takes a new version. */
|
||||
private val byId = HashMap<HexKey, Slot<T>>()
|
||||
@@ -137,9 +161,6 @@ class EventStoreProjection<T : Event>(
|
||||
private val perFilter: Map<Filter, java.util.SortedSet<Slot<T>>> =
|
||||
filters.associateWith { sortedSetOf(slotComparator()) }
|
||||
|
||||
/** Set when the seed has been written to [items], so callers can suspend until the projection is hot. */
|
||||
val ready: CompletableDeferred<Unit> = CompletableDeferred()
|
||||
|
||||
private val collectorJob: Job =
|
||||
scope.launch {
|
||||
// `onSubscription` runs after the SharedFlow subscription is
|
||||
@@ -150,7 +171,6 @@ class EventStoreProjection<T : Event>(
|
||||
store.changes
|
||||
.onSubscription {
|
||||
seed()
|
||||
ready.complete(Unit)
|
||||
}.collect { storeEvent -> apply(storeEvent) }
|
||||
}
|
||||
|
||||
@@ -326,12 +346,12 @@ class EventStoreProjection<T : Event>(
|
||||
// sets. Cheaper than maintaining a separate `ordered` field
|
||||
// alongside every insert / remove.
|
||||
if (byId.isEmpty()) {
|
||||
_items.value = emptyList()
|
||||
_state.value = ProjectionState.Loaded(emptyList())
|
||||
return
|
||||
}
|
||||
val union = sortedSetOf(slotComparator<T>())
|
||||
for (set in perFilter.values) union.addAll(set)
|
||||
_items.value = union.map { it.flow }
|
||||
_state.value = ProjectionState.Loaded(union.map { it.flow })
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -344,7 +364,7 @@ class EventStoreProjection<T : Event>(
|
||||
byId.clear()
|
||||
byAddress.clear()
|
||||
for (set in perFilter.values) set.clear()
|
||||
_items.value = emptyList()
|
||||
_state.value = ProjectionState.Loaded(emptyList())
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
+54
-44
@@ -71,12 +71,22 @@ class EventStoreProjectionTest {
|
||||
store.close()
|
||||
}
|
||||
|
||||
/** Snapshot of the currently-loaded items, or empty if still seeding. */
|
||||
private val <T : Event> EventStoreProjection<T>.items: List<MutableStateFlow<T>>
|
||||
get() = (state.value as? ProjectionState.Loaded)?.items.orEmpty()
|
||||
|
||||
/** Suspends until the seed completes; returns the loaded list. */
|
||||
private suspend fun <T : Event> EventStoreProjection<T>.awaitReady(timeoutMs: Long = 5_000): List<MutableStateFlow<T>> =
|
||||
withTimeout(timeoutMs) {
|
||||
(state.first { it is ProjectionState.Loaded } as ProjectionState.Loaded).items
|
||||
}
|
||||
|
||||
private suspend fun <T : Event> EventStoreProjection<T>.awaitItems(
|
||||
timeoutMs: Long = 5_000,
|
||||
predicate: (List<MutableStateFlow<T>>) -> Boolean,
|
||||
): List<MutableStateFlow<T>> =
|
||||
withTimeout(timeoutMs) {
|
||||
items.first { predicate(it) }
|
||||
(state.first { it is ProjectionState.Loaded && predicate(it.items) } as ProjectionState.Loaded).items
|
||||
}
|
||||
|
||||
private suspend fun <T : Event> awaitFlow(
|
||||
@@ -97,9 +107,9 @@ class EventStoreProjectionTest {
|
||||
observable.insert(b)
|
||||
|
||||
val projection = observable.observe<TextNoteEvent>(Filter(kinds = listOf(TextNoteEvent.KIND)), scope)
|
||||
projection.ready.await()
|
||||
projection.awaitReady()
|
||||
|
||||
val items = projection.items.value
|
||||
val items = projection.items
|
||||
assertEquals(2, items.size)
|
||||
assertEquals(b.id, items[0].value.id)
|
||||
assertEquals(a.id, items[1].value.id)
|
||||
@@ -113,8 +123,8 @@ class EventStoreProjectionTest {
|
||||
observable.insert(a)
|
||||
|
||||
val projection = observable.observe<TextNoteEvent>(Filter(kinds = listOf(TextNoteEvent.KIND)), scope)
|
||||
projection.ready.await()
|
||||
val before = projection.items.value
|
||||
projection.awaitReady()
|
||||
val before = projection.items
|
||||
assertEquals(1, before.size)
|
||||
|
||||
val b = signer.sign(TextNoteEvent.build("b", createdAt = 200))
|
||||
@@ -133,14 +143,14 @@ class EventStoreProjectionTest {
|
||||
observable.insert(text)
|
||||
|
||||
val projection = observable.observe<TextNoteEvent>(Filter(kinds = listOf(TextNoteEvent.KIND)), scope)
|
||||
projection.ready.await()
|
||||
val seed = projection.items.value
|
||||
projection.awaitReady()
|
||||
val seed = projection.items
|
||||
|
||||
val meta = signer.sign(MetadataEvent.createNew("Vitor", createdAt = 200))
|
||||
observable.insert(meta)
|
||||
|
||||
delay(150)
|
||||
assertSame(seed, projection.items.value)
|
||||
assertSame(seed, projection.items)
|
||||
projection.close()
|
||||
}
|
||||
|
||||
@@ -156,8 +166,8 @@ class EventStoreProjectionTest {
|
||||
Filter(kinds = listOf(MetadataEvent.KIND), authors = listOf(v1.pubKey)),
|
||||
scope,
|
||||
)
|
||||
projection.ready.await()
|
||||
val seedList = projection.items.value
|
||||
projection.awaitReady()
|
||||
val seedList = projection.items
|
||||
assertEquals(1, seedList.size)
|
||||
val slot = seedList[0]
|
||||
assertEquals(v1.id, slot.value.id)
|
||||
@@ -166,8 +176,8 @@ class EventStoreProjectionTest {
|
||||
observable.insert(v2)
|
||||
|
||||
awaitFlow(slot) { it.id == v2.id }
|
||||
assertSame(seedList, projection.items.value, "replaceable update must not change list reference")
|
||||
assertSame(slot, projection.items.value[0])
|
||||
assertSame(seedList, projection.items, "replaceable update must not change list reference")
|
||||
assertSame(slot, projection.items[0])
|
||||
projection.close()
|
||||
}
|
||||
|
||||
@@ -187,15 +197,15 @@ class EventStoreProjectionTest {
|
||||
),
|
||||
scope,
|
||||
)
|
||||
projection.ready.await()
|
||||
val seedList = projection.items.value
|
||||
projection.awaitReady()
|
||||
val seedList = projection.items
|
||||
val slot = seedList[0]
|
||||
|
||||
val v2 = signer.sign(LongTextNoteEvent.build("blog v2", "title", dTag = "blog", createdAt = time + 1))
|
||||
observable.insert(v2)
|
||||
|
||||
awaitFlow(slot) { it.id == v2.id }
|
||||
assertSame(seedList, projection.items.value, "addressable update must not change list reference")
|
||||
assertSame(seedList, projection.items, "addressable update must not change list reference")
|
||||
projection.close()
|
||||
}
|
||||
|
||||
@@ -222,8 +232,8 @@ class EventStoreProjectionTest {
|
||||
Filter(kinds = listOf(MetadataEvent.KIND), authors = listOf(v1.pubKey)),
|
||||
scope,
|
||||
)
|
||||
projection.ready.await()
|
||||
val slot = projection.items.value[0]
|
||||
projection.awaitReady()
|
||||
val slot = projection.items[0]
|
||||
assertEquals(v2.id, slot.value.id)
|
||||
|
||||
// The store rejects v1 because v2 already won; the
|
||||
@@ -249,8 +259,8 @@ class EventStoreProjectionTest {
|
||||
observable.insert(b)
|
||||
|
||||
val projection = observable.observe<Event>(Filter(kinds = listOf(TextNoteEvent.KIND)), scope)
|
||||
projection.ready.await()
|
||||
assertEquals(2, projection.items.value.size)
|
||||
projection.awaitReady()
|
||||
assertEquals(2, projection.items.size)
|
||||
|
||||
val deletion = signer.sign(DeletionEvent.build(listOf(a)))
|
||||
observable.insert(deletion)
|
||||
@@ -271,8 +281,8 @@ class EventStoreProjectionTest {
|
||||
observable.insert(a)
|
||||
|
||||
val projection = observable.observe<Event>(Filter(kinds = listOf(TextNoteEvent.KIND)), scope)
|
||||
projection.ready.await()
|
||||
val seed = projection.items.value
|
||||
projection.awaitReady()
|
||||
val seed = projection.items
|
||||
assertEquals(1, seed.size)
|
||||
|
||||
val foreignDeletion = otherSigner.sign(DeletionEvent.build(listOf(a)))
|
||||
@@ -280,10 +290,10 @@ class EventStoreProjectionTest {
|
||||
|
||||
// Give the projection time to process the event.
|
||||
delay(150)
|
||||
assertSame(seed, projection.items.value)
|
||||
assertSame(seed, projection.items)
|
||||
assertEquals(
|
||||
a.id,
|
||||
projection.items.value[0]
|
||||
projection.items[0]
|
||||
.value.id,
|
||||
)
|
||||
projection.close()
|
||||
@@ -299,8 +309,8 @@ class EventStoreProjectionTest {
|
||||
observable.insert(b)
|
||||
|
||||
val projection = observable.observe<Event>(Filter(kinds = listOf(TextNoteEvent.KIND)), scope)
|
||||
projection.ready.await()
|
||||
assertEquals(2, projection.items.value.size)
|
||||
projection.awaitReady()
|
||||
assertEquals(2, projection.items.size)
|
||||
|
||||
val vanish =
|
||||
signer.sign(
|
||||
@@ -328,8 +338,8 @@ class EventStoreProjectionTest {
|
||||
observable.insert(a)
|
||||
|
||||
val projection = observable.observe<Event>(Filter(kinds = listOf(TextNoteEvent.KIND)), scope)
|
||||
projection.ready.await()
|
||||
val seed = projection.items.value
|
||||
projection.awaitReady()
|
||||
val seed = projection.items
|
||||
|
||||
val foreignVanish =
|
||||
otherSigner.sign(
|
||||
@@ -341,7 +351,7 @@ class EventStoreProjectionTest {
|
||||
observable.insert(foreignVanish)
|
||||
|
||||
delay(150)
|
||||
assertSame(seed, projection.items.value)
|
||||
assertSame(seed, projection.items)
|
||||
projection.close()
|
||||
}
|
||||
|
||||
@@ -364,8 +374,8 @@ class EventStoreProjectionTest {
|
||||
Filter(kinds = listOf(TextNoteEvent.KIND)),
|
||||
scope,
|
||||
)
|
||||
projection.ready.await()
|
||||
assertEquals(2, projection.items.value.size)
|
||||
projection.awaitReady()
|
||||
assertEquals(2, projection.items.size)
|
||||
|
||||
// Let the short expiration lapse, then ask the store to
|
||||
// sweep — the projection drops the expired slot in
|
||||
@@ -398,8 +408,8 @@ class EventStoreProjectionTest {
|
||||
Filter(kinds = listOf(TextNoteEvent.KIND)),
|
||||
scope,
|
||||
)
|
||||
projection.ready.await()
|
||||
assertEquals(3, projection.items.value.size)
|
||||
projection.awaitReady()
|
||||
assertEquals(3, projection.items.size)
|
||||
|
||||
// Drop everything authored by `signer` — should leave
|
||||
// only the foreign event.
|
||||
@@ -423,8 +433,8 @@ class EventStoreProjectionTest {
|
||||
Filter(kinds = listOf(TextNoteEvent.KIND), limit = 2),
|
||||
scope,
|
||||
)
|
||||
projection.ready.await()
|
||||
assertEquals(2, projection.items.value.size)
|
||||
projection.awaitReady()
|
||||
assertEquals(2, projection.items.size)
|
||||
|
||||
val c = signer.sign(TextNoteEvent.build("c", createdAt = 300))
|
||||
observable.insert(c)
|
||||
@@ -462,8 +472,8 @@ class EventStoreProjectionTest {
|
||||
listOf(filterA, filterB),
|
||||
scope,
|
||||
)
|
||||
projection.ready.await()
|
||||
assertEquals(4, projection.items.value.size, "per-filter caps don't dedupe union")
|
||||
projection.awaitReady()
|
||||
assertEquals(4, projection.items.size, "per-filter caps don't dedupe union")
|
||||
projection.close()
|
||||
}
|
||||
|
||||
@@ -485,8 +495,8 @@ class EventStoreProjectionTest {
|
||||
Filter(kinds = listOf(TextNoteEvent.KIND), limit = 2),
|
||||
scope,
|
||||
)
|
||||
projection.ready.await()
|
||||
assertEquals(2, projection.items.value.size)
|
||||
projection.awaitReady()
|
||||
assertEquals(2, projection.items.size)
|
||||
|
||||
observable.insert(a3)
|
||||
val after = projection.awaitItems { it[0].value.id == a3.id }
|
||||
@@ -503,12 +513,12 @@ class EventStoreProjectionTest {
|
||||
observable.insert(a)
|
||||
|
||||
val projection = observable.observe<TextNoteEvent>(Filter(kinds = listOf(TextNoteEvent.KIND)), scope)
|
||||
projection.ready.await()
|
||||
projection.awaitReady()
|
||||
projection.close()
|
||||
|
||||
observable.insert(signer.sign(TextNoteEvent.build("b", createdAt = 200)))
|
||||
delay(150)
|
||||
assertTrue(projection.items.value.isEmpty())
|
||||
assertTrue(projection.items.isEmpty())
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -525,8 +535,8 @@ class EventStoreProjectionTest {
|
||||
val ephemeralKind = 22_000
|
||||
val projection =
|
||||
observable.observe<Event>(Filter(kinds = listOf(ephemeralKind)), scope)
|
||||
projection.ready.await()
|
||||
assertTrue(projection.items.value.isEmpty())
|
||||
projection.awaitReady()
|
||||
assertTrue(projection.items.isEmpty())
|
||||
|
||||
val ephemeral: Event =
|
||||
signer.sign(
|
||||
@@ -548,8 +558,8 @@ class EventStoreProjectionTest {
|
||||
// event was only ever live, not durable.
|
||||
val freshProjection =
|
||||
observable.observe<Event>(Filter(kinds = listOf(ephemeralKind)), scope)
|
||||
freshProjection.ready.await()
|
||||
assertTrue(freshProjection.items.value.isEmpty())
|
||||
freshProjection.awaitReady()
|
||||
assertTrue(freshProjection.items.isEmpty())
|
||||
|
||||
projection.close()
|
||||
freshProjection.close()
|
||||
|
||||
Reference in New Issue
Block a user