From 70420083acc6fdc0662a8b3cbab75153f60823b2 Mon Sep 17 00:00:00 2001 From: kami Date: Sun, 24 May 2026 22:23:43 +0400 Subject: [PATCH] feat(server-02): promote EventStore to global-sequence store with subscribeAll StoredEvent gains sessionSequence: Long (per-session monotonic); existing sequence field becomes the global monotonic cursor across the entire store. SqliteEventStore allocates both counters inside the same BEGIN IMMEDIATE transaction and emits to a global MutableSharedFlow after commit. InMemoryEventStore mirrors this with an AtomicLong global counter and a suspend-on-overflow SharedFlow. EventStore interface gains subscribeAll() and lastGlobalSequence(). Contract tests updated and extended to cover global-sequence and cross-session subscribeAll invariants. --- .../apps/server/logging/LoggingEventStore.kt | 4 + .../core/events/events/EventEnvelope.kt | 1 + .../correx/core/events/stores/EventStore.kt | 4 + .../persistence/InMemoryEventStore.kt | 61 +++++++---- .../persistence/SqliteEventStore.kt | 101 +++++++++++------- .../events/store/EventStoreContractTest.kt | 80 +++++++++++--- .../src/test/kotlin/RouterFacadeTest.kt | 7 +- .../correx/testing/fixtures/EventFixtures.kt | 3 +- 8 files changed, 184 insertions(+), 77 deletions(-) diff --git a/apps/server/src/main/kotlin/com/correx/apps/server/logging/LoggingEventStore.kt b/apps/server/src/main/kotlin/com/correx/apps/server/logging/LoggingEventStore.kt index 93bf7533..17920fac 100644 --- a/apps/server/src/main/kotlin/com/correx/apps/server/logging/LoggingEventStore.kt +++ b/apps/server/src/main/kotlin/com/correx/apps/server/logging/LoggingEventStore.kt @@ -61,4 +61,8 @@ class LoggingEventStore(private val delegate: EventStore) : EventStore { override fun allEvents(): Sequence = delegate.allEvents() override fun allSessionIds(): Set = delegate.allSessionIds().also { log.debug("got session ids from store: {}", it) } + + override fun subscribeAll(): Flow = delegate.subscribeAll() + + override suspend fun lastGlobalSequence(): Long = delegate.lastGlobalSequence() } diff --git a/core/events/src/main/kotlin/com/correx/core/events/events/EventEnvelope.kt b/core/events/src/main/kotlin/com/correx/core/events/events/EventEnvelope.kt index f5c3dae8..42d976d3 100644 --- a/core/events/src/main/kotlin/com/correx/core/events/events/EventEnvelope.kt +++ b/core/events/src/main/kotlin/com/correx/core/events/events/EventEnvelope.kt @@ -12,5 +12,6 @@ data class NewEvent( data class StoredEvent( val metadata: EventMetadata, val sequence: Long, + val sessionSequence: Long, val payload: EventPayload ) diff --git a/core/events/src/main/kotlin/com/correx/core/events/stores/EventStore.kt b/core/events/src/main/kotlin/com/correx/core/events/stores/EventStore.kt index 73ae179a..97a91dfd 100644 --- a/core/events/src/main/kotlin/com/correx/core/events/stores/EventStore.kt +++ b/core/events/src/main/kotlin/com/correx/core/events/stores/EventStore.kt @@ -42,6 +42,10 @@ interface EventStore { */ fun subscribe(sessionId: SessionId): Flow + fun subscribeAll(): Flow + + suspend fun lastGlobalSequence(): Long + /** * Iterate every event in the store across all sessions, in monotonic insertion order. * Intended for batch jobs (e.g. CAS liveness scans). Implementations should stream lazily. diff --git a/infrastructure/persistence/src/main/kotlin/com/correx/infrastructure/persistence/InMemoryEventStore.kt b/infrastructure/persistence/src/main/kotlin/com/correx/infrastructure/persistence/InMemoryEventStore.kt index 041c2102..efdf655a 100644 --- a/infrastructure/persistence/src/main/kotlin/com/correx/infrastructure/persistence/InMemoryEventStore.kt +++ b/infrastructure/persistence/src/main/kotlin/com/correx/infrastructure/persistence/InMemoryEventStore.kt @@ -5,59 +5,71 @@ import com.correx.core.events.events.StoredEvent import com.correx.core.events.stores.EventStore import com.correx.core.events.types.EventId import com.correx.core.events.types.SessionId +import kotlinx.coroutines.channels.BufferOverflow import kotlinx.coroutines.flow.Flow import kotlinx.coroutines.flow.MutableSharedFlow +import kotlinx.coroutines.flow.asSharedFlow import java.util.concurrent.* import java.util.concurrent.atomic.* class InMemoryEventStore : EventStore { private val streams = ConcurrentHashMap>() private val sequences = ConcurrentHashMap() + private val globalSeq = AtomicLong(0) private val seenEventIds = ConcurrentHashMap.newKeySet() private val subscriptions = ConcurrentHashMap>() + private val globalFlow = MutableSharedFlow( + replay = 0, + extraBufferCapacity = 1024, + onBufferOverflow = BufferOverflow.SUSPEND, + ) override suspend fun append(event: NewEvent): StoredEvent { val stream = streams.computeIfAbsent(event.metadata.sessionId) { mutableListOf() } - + val stored: StoredEvent synchronized(stream) { if (!seenEventIds.add(event.metadata.eventId)) { error("duplicate event_id: ${event.metadata.eventId}") } - return doAppend(event, stream) + stored = doAppend(event, stream) } + subscriptions[event.metadata.sessionId]?.tryEmit(stored) + globalFlow.emit(stored) + return stored } override suspend fun appendAll(events: List): List { if (events.isEmpty()) return emptyList() - val sessionId = events.first().metadata.sessionId val stream = streams.computeIfAbsent(sessionId) { mutableListOf() } - + val stored: List synchronized(stream) { - return events.map { doAppend(it, stream) } + stored = events.map { doAppend(it, stream) } } + val flow = subscriptions[sessionId] + if (flow != null) stored.forEach { flow.tryEmit(it) } + stored.forEach { globalFlow.emit(it) } + return stored } - override fun read(sessionId: SessionId): List { - return streams[sessionId] - ?.sortedBy { it.sequence } - ?: emptyList() - } + override fun read(sessionId: SessionId): List = + streams[sessionId]?.sortedBy { it.sessionSequence } ?: emptyList() - override fun readFrom(sessionId: SessionId, fromSequence: Long): List { - return streams[sessionId] - ?.filter { it.sequence > fromSequence } + override fun readFrom(sessionId: SessionId, fromSequence: Long): List = + streams[sessionId] + ?.filter { it.sessionSequence > fromSequence } ?.toList() ?: emptyList() - } - override fun lastSequence(sessionId: SessionId): Long? { - return sequences[sessionId]?.get() - } + override fun lastSequence(sessionId: SessionId): Long? = sequences[sessionId]?.get() override fun subscribe(sessionId: SessionId): Flow = subscriptions.computeIfAbsent(sessionId) { MutableSharedFlow(replay = 0, extraBufferCapacity = 64) } + override fun subscribeAll(): Flow = globalFlow.asSharedFlow() + + override suspend fun lastGlobalSequence(): Long = globalSeq.get() + override fun allEvents(): Sequence = streams.values .flatMap { stream -> synchronized(stream) { stream.toList() } } @@ -66,17 +78,20 @@ class InMemoryEventStore : EventStore { override fun allSessionIds(): Set = sequences.keys private fun doAppend(event: NewEvent, stream: MutableList): StoredEvent { - val seq = sequences.computeIfAbsent(event.metadata.sessionId) { AtomicLong(0) } + val sessionSeq = sequences.computeIfAbsent(event.metadata.sessionId) { AtomicLong(0) } .incrementAndGet() + val globalSeqVal = globalSeq.incrementAndGet() - val stored = StoredEvent(metadata = event.metadata, sequence = seq, payload = event.payload) - // safety: enforce ordering invariant - if (stream.isNotEmpty() && stream.last().sequence >= stored.sequence) { + val stored = StoredEvent( + metadata = event.metadata, + sequence = globalSeqVal, + sessionSequence = sessionSeq, + payload = event.payload, + ) + if (stream.isNotEmpty() && stream.last().sessionSequence >= stored.sessionSequence) { error("sequence violation for session ${event.metadata.sessionId}") } stream.add(stored) - subscriptions[event.metadata.sessionId]?.tryEmit(stored) - return stored } } diff --git a/infrastructure/persistence/src/main/kotlin/com/correx/infrastructure/persistence/SqliteEventStore.kt b/infrastructure/persistence/src/main/kotlin/com/correx/infrastructure/persistence/SqliteEventStore.kt index 360cbdc1..616e1625 100644 --- a/infrastructure/persistence/src/main/kotlin/com/correx/infrastructure/persistence/SqliteEventStore.kt +++ b/infrastructure/persistence/src/main/kotlin/com/correx/infrastructure/persistence/SqliteEventStore.kt @@ -13,8 +13,10 @@ import com.correx.core.events.types.EventId import com.correx.core.events.types.SessionId import com.correx.infrastructure.persistence.util.JDBCHelper.transaction import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.channels.BufferOverflow import kotlinx.coroutines.flow.Flow import kotlinx.coroutines.flow.MutableSharedFlow +import kotlinx.coroutines.flow.asSharedFlow import kotlinx.coroutines.withContext import kotlinx.datetime.Instant import java.sql.Connection @@ -27,6 +29,11 @@ class SqliteEventStore( private val artifactStore: ArtifactStore, ) : EventStore { private val subscriptions = ConcurrentHashMap>() + private val globalFlow = MutableSharedFlow( + replay = 0, + extraBufferCapacity = 1024, + onBufferOverflow = BufferOverflow.SUSPEND, + ) init { connection.createStatement().use { stmt -> @@ -36,6 +43,7 @@ class SqliteEventStore( event_id TEXT PRIMARY KEY, session_id TEXT NOT NULL, sequence INTEGER NOT NULL, + session_sequence INTEGER NOT NULL, timestamp TEXT NOT NULL, schema_version INTEGER NOT NULL, causation_id TEXT, @@ -44,6 +52,12 @@ class SqliteEventStore( ); """.trimIndent(), ) + stmt.execute( + "CREATE UNIQUE INDEX IF NOT EXISTS events_sequence_uq ON events(sequence);" + ) + stmt.execute( + "CREATE UNIQUE INDEX IF NOT EXISTS events_session_seq_uq ON events(session_id, session_sequence);" + ) } } @@ -54,15 +68,14 @@ class SqliteEventStore( stored = connection.transaction { val existing = findByEventId(event.metadata.eventId) check(existing == null) { "duplicate event_id: ${event.metadata.eventId}" } - - val seq = nextSequence(event.metadata.sessionId) - + val globalSeqVal = nextGlobalSequence() + val sessionSeqVal = nextSessionSequence(event.metadata.sessionId) val s = StoredEvent( metadata = event.metadata, - sequence = seq, + sequence = globalSeqVal, + sessionSequence = sessionSeqVal, payload = event.payload, ) - insert(s) s } @@ -70,14 +83,13 @@ class SqliteEventStore( } val result = checkNotNull(stored) subscriptions[event.metadata.sessionId]?.tryEmit(result) + globalFlow.emit(result) return result } override suspend fun appendAll(events: List): List { if (events.isEmpty()) return emptyList() - val sessionId = events.first().metadata.sessionId - var stored: List = emptyList() artifactStore.flushBefore { withContext(Dispatchers.IO) { @@ -85,15 +97,14 @@ class SqliteEventStore( events.map { event -> val existing = findByEventId(event.metadata.eventId) check(existing == null) { "duplicate event_id: ${event.metadata.eventId}" } - - val seq = nextSequence(sessionId) - + val globalSeqVal = nextGlobalSequence() + val sessionSeqVal = nextSessionSequence(sessionId) val s = StoredEvent( metadata = event.metadata, - sequence = seq, + sequence = globalSeqVal, + sessionSequence = sessionSeqVal, payload = event.payload, ) - insert(s) s } @@ -102,6 +113,7 @@ class SqliteEventStore( } val flow = subscriptions[sessionId] if (flow != null) stored.forEach { flow.tryEmit(it) } + stored.forEach { globalFlow.emit(it) } return stored } @@ -110,11 +122,10 @@ class SqliteEventStore( """ SELECT * FROM events WHERE session_id = ? - ORDER BY sequence ASC + ORDER BY session_sequence ASC """.trimIndent(), ).use { ps -> ps.setString(1, sessionId.value) - ps.executeQuery().use { rs -> buildList { while (rs.next()) { @@ -129,13 +140,12 @@ class SqliteEventStore( """ SELECT * FROM events WHERE session_id = ? - AND sequence > ? - ORDER BY sequence ASC + AND session_sequence > ? + ORDER BY session_sequence ASC """.trimIndent(), ).use { ps -> ps.setString(1, sessionId.value) ps.setLong(2, fromSequence) - ps.executeQuery().use { rs -> buildList { while (rs.next()) { @@ -147,14 +157,9 @@ class SqliteEventStore( override fun lastSequence(sessionId: SessionId): Long? = connection.prepareStatement( - """ - SELECT MAX(sequence) - FROM events - WHERE session_id = ? - """.trimIndent(), + "SELECT MAX(session_sequence) FROM events WHERE session_id = ?" ).use { ps -> ps.setString(1, sessionId.value) - ps.executeQuery().use { rs -> if (rs.next()) rs.getLong(1).takeIf { !rs.wasNull() } else null @@ -164,6 +169,20 @@ class SqliteEventStore( override fun subscribe(sessionId: SessionId): Flow = subscriptions.computeIfAbsent(sessionId) { MutableSharedFlow(replay = 0, extraBufferCapacity = 64) } + override fun subscribeAll(): Flow = globalFlow.asSharedFlow() + + override suspend fun lastGlobalSequence(): Long = + withContext(Dispatchers.IO) { + connection.prepareStatement( + "SELECT COALESCE(MAX(sequence), 0) FROM events" + ).use { ps -> + ps.executeQuery().use { rs -> + rs.next() + rs.getLong(1) + } + } + } + override fun allEvents(): Sequence = connection.prepareStatement( """ @@ -191,13 +210,19 @@ class SqliteEventStore( // ---------- helpers ---------- - private fun nextSequence(sessionId: SessionId): Long = + private fun nextGlobalSequence(): Long = connection.prepareStatement( - """ - SELECT COALESCE(MAX(sequence), 0) - FROM events - WHERE session_id = ? - """.trimIndent(), + "SELECT COALESCE(MAX(sequence), 0) FROM events" + ).use { ps -> + ps.executeQuery().use { rs -> + rs.next() + rs.getLong(1) + 1 + } + } + + private fun nextSessionSequence(sessionId: SessionId): Long = + connection.prepareStatement( + "SELECT COALESCE(MAX(session_sequence), 0) FROM events WHERE session_id = ?" ).use { ps -> ps.setString(1, sessionId.value) ps.executeQuery().use { rs -> @@ -214,23 +239,24 @@ class SqliteEventStore( event_id, session_id, sequence, + session_sequence, timestamp, schema_version, causation_id, correlation_id, payload - ) VALUES (?, ?, ?, ?, ?, ?, ?, ?) + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) """.trimIndent(), ).use { ps -> ps.setString(1, event.metadata.eventId.value) ps.setString(2, event.metadata.sessionId.value) ps.setLong(3, event.sequence) - ps.setString(4, event.metadata.timestamp.toString()) - ps.setInt(5, event.metadata.schemaVersion) - ps.setString(6, event.metadata.causationId?.value) - ps.setString(7, event.metadata.correlationId?.value) - ps.setString(8, jsonSerializer.serialize(event.payload)) - + ps.setLong(4, event.sessionSequence) + ps.setString(5, event.metadata.timestamp.toString()) + ps.setInt(6, event.metadata.schemaVersion) + ps.setString(7, event.metadata.causationId?.value) + ps.setString(8, event.metadata.correlationId?.value) + ps.setString(9, jsonSerializer.serialize(event.payload)) ps.executeUpdate() } } @@ -243,12 +269,10 @@ class SqliteEventStore( """.trimIndent(), ).use { ps -> ps.setString(1, eventId.value) - ps.executeQuery().use { rs -> if (rs.next()) rs.toStoredEvent(jsonSerializer) else null } } - } private fun ResultSet.toStoredEvent(jsonSerializer: JsonEventSerializer): StoredEvent = @@ -262,5 +286,6 @@ private fun ResultSet.toStoredEvent(jsonSerializer: JsonEventSerializer): Stored correlationId = getString("correlation_id")?.let { CorrelationId(it) }, ), sequence = getLong("sequence"), + sessionSequence = getLong("session_sequence"), payload = jsonSerializer.deserialize(getString("payload")), ) diff --git a/testing/contracts/src/testFixtures/kotlin/com/correx/testing/contracts/fixtures/events/store/EventStoreContractTest.kt b/testing/contracts/src/testFixtures/kotlin/com/correx/testing/contracts/fixtures/events/store/EventStoreContractTest.kt index 49687588..7a679796 100644 --- a/testing/contracts/src/testFixtures/kotlin/com/correx/testing/contracts/fixtures/events/store/EventStoreContractTest.kt +++ b/testing/contracts/src/testFixtures/kotlin/com/correx/testing/contracts/fixtures/events/store/EventStoreContractTest.kt @@ -1,11 +1,16 @@ package com.correx.testing.contracts.fixtures.events.store +import com.correx.core.events.events.StoredEvent import com.correx.core.events.stores.EventStore import com.correx.core.events.types.EventId import com.correx.core.events.types.SessionId import com.correx.testing.contracts.utils.ConcurrencyRunner import com.correx.testing.fixtures.EventFixtures +import kotlinx.coroutines.delay +import kotlinx.coroutines.flow.collect +import kotlinx.coroutines.launch import kotlinx.coroutines.runBlocking +import kotlinx.coroutines.yield import org.junit.jupiter.api.Assertions import org.junit.jupiter.api.Test import org.junit.jupiter.api.assertThrows @@ -83,22 +88,22 @@ abstract class EventStoreContractTest { } @Test - fun `sequences are strictly increasing per session`(): Unit = runBlocking { + fun `sessionSequences are strictly increasing per session`(): Unit = runBlocking { val store = store() repeat(100) { store.append(EventFixtures.newEvent(EventId("e$it"), SessionId("s1"))) } - val seqs = store.read(SessionId("s1")).map { it.sequence } + val sessionSeqs = store.read(SessionId("s1")).map { it.sessionSequence } + sessionSeqs.zipWithNext().forEach { (a, b) -> Assertions.assertTrue(a < b) } - seqs.zipWithNext().forEach { (a, b) -> - Assertions.assertTrue(a < b) - } + val globalSeqs = store.read(SessionId("s1")).map { it.sequence } + globalSeqs.zipWithNext().forEach { (a, b) -> Assertions.assertTrue(a < b) } } @Test - fun `sequences are isolated per session`(): Unit = runBlocking { + fun `sessionSequences are isolated per session and global sequence spans all`(): Unit = runBlocking { val store = store() repeat(50) { @@ -106,11 +111,17 @@ abstract class EventStoreContractTest { store.append(EventFixtures.newEvent(EventId("b$it"), SessionId("B"))) } - val aSeq = store.read(SessionId("A")).map { it.sequence } - val bSeq = store.read(SessionId("B")).map { it.sequence } + val aSessionSeq = store.read(SessionId("A")).map { it.sessionSequence } + val bSessionSeq = store.read(SessionId("B")).map { it.sessionSequence } - Assertions.assertEquals(aSeq.sorted(), aSeq) - Assertions.assertEquals(bSeq.sorted(), bSeq) + Assertions.assertEquals(aSessionSeq.sorted(), aSessionSeq) + Assertions.assertEquals(bSessionSeq.sorted(), bSessionSeq) + Assertions.assertEquals(1L, aSessionSeq.first()) + Assertions.assertEquals(1L, bSessionSeq.first()) + + val allGlobal = (store.read(SessionId("A")) + store.read(SessionId("B"))) + .map { it.sequence }.sorted() + allGlobal.zipWithNext().forEach { (a, b) -> Assertions.assertTrue(a < b) } } @Test @@ -127,7 +138,7 @@ abstract class EventStoreContractTest { Assertions.assertEquals(50, result.size) - val seqs = result.map { it.sequence } + val seqs = result.map { it.sessionSequence } Assertions.assertEquals(seqs.sorted(), seqs) } @@ -149,7 +160,7 @@ abstract class EventStoreContractTest { val all = store.read(SessionId("s1")) - val seqs = all.map { it.sequence } + val seqs = all.map { it.sessionSequence } Assertions.assertEquals(seqs.sorted(), seqs) } @@ -179,7 +190,48 @@ abstract class EventStoreContractTest { val first = store.readFrom(SessionId("s1"), 3) val second = store.readFrom(SessionId("s1"), 7) - Assertions.assertTrue(first.all { it.sequence > 3 }) - Assertions.assertTrue(second.all { it.sequence > 7 }) + Assertions.assertTrue(first.all { it.sessionSequence > 3 }) + Assertions.assertTrue(second.all { it.sessionSequence > 7 }) + } + + @Test + fun `lastGlobalSequence returns 0 when store is empty`(): Unit = runBlocking { + val store = store() + Assertions.assertEquals(0L, store.lastGlobalSequence()) + } + + @Test + fun `lastGlobalSequence equals highest sequence after appends`(): Unit = runBlocking { + val store = store() + + store.append(EventFixtures.newEvent(EventId("x1"), SessionId("s1"))) + store.append(EventFixtures.newEvent(EventId("x2"), SessionId("s2"))) + store.append(EventFixtures.newEvent(EventId("x3"), SessionId("s1"))) + + val allSeqs = (store.read(SessionId("s1")) + store.read(SessionId("s2"))) + .map { it.sequence } + + Assertions.assertEquals(allSeqs.max(), store.lastGlobalSequence()) + } + + @Test + fun `subscribeAll emits events from all sessions`(): Unit = runBlocking { + val store = store() + val received = mutableListOf() + + val job = launch { + store.subscribeAll().collect { received.add(it) } + } + yield() + + store.append(EventFixtures.newEvent(EventId("p1"), SessionId("alpha"))) + store.append(EventFixtures.newEvent(EventId("q1"), SessionId("beta"))) + delay(100) + job.cancel() + + Assertions.assertEquals(2, received.size) + val sessionIds = received.map { it.metadata.sessionId.value }.toSet() + Assertions.assertTrue("alpha" in sessionIds) + Assertions.assertTrue("beta" in sessionIds) } } diff --git a/testing/deterministic/src/test/kotlin/RouterFacadeTest.kt b/testing/deterministic/src/test/kotlin/RouterFacadeTest.kt index d9625afe..e8b3cb16 100644 --- a/testing/deterministic/src/test/kotlin/RouterFacadeTest.kt +++ b/testing/deterministic/src/test/kotlin/RouterFacadeTest.kt @@ -606,7 +606,8 @@ class RouterFacadeTest { appendedEvents.add(event) val stored = StoredEvent( metadata = event.metadata, - sequence = nextSequence++, + sequence = nextSequence, + sessionSequence = nextSequence++, payload = event.payload, ) storedEvents[event.metadata.eventId] = stored @@ -635,5 +636,9 @@ class RouterFacadeTest { storedEvents.values.asSequence() override fun allSessionIds(): Set = storedEvents.values.map { it.metadata.sessionId }.toSet() + + override fun subscribeAll(): Flow = TODO("Not needed in this test context") + + override suspend fun lastGlobalSequence(): Long = TODO("Not needed in this test context") } } diff --git a/testing/fixtures/src/main/kotlin/com/correx/testing/fixtures/EventFixtures.kt b/testing/fixtures/src/main/kotlin/com/correx/testing/fixtures/EventFixtures.kt index 8502d888..748b0886 100644 --- a/testing/fixtures/src/main/kotlin/com/correx/testing/fixtures/EventFixtures.kt +++ b/testing/fixtures/src/main/kotlin/com/correx/testing/fixtures/EventFixtures.kt @@ -34,6 +34,7 @@ object EventFixtures { payload: EventPayload = ToolInvokedEvent("write"), sessionId: SessionId = SessionId("s1"), sequence: Long = 1L, + sessionSequence: Long = sequence, timestamp: Instant = Instant.parse("2026-01-01T00:00:00Z") ): StoredEvent { return StoredEvent( @@ -46,8 +47,8 @@ object EventFixtures { correlationId = null ), sequence = sequence, + sessionSequence = sessionSequence, payload = payload ) - } }