From a0b1209457e1c6ec6c3715b08608ab68111e27ba Mon Sep 17 00:00:00 2001 From: subochev Date: Sun, 20 Sep 2026 16:30:09 +0300 Subject: [PATCH] feat(event-store-in-memory): InMemoryEventStore implementation of MutableEventStore MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Первый concrete impl :event-store. ConcurrentLinkedDeque + eviction на каждом append (amortized O(1) при стабильном размере буфера). Контракт (по требованию пользователя): - maxMessages: Int? — nullable, NO default (явный null = unlimited) - ttl: Duration? — nullable, NO default (явный null = forever) - оба null → store forever - любой non-null → соответствующий eviction policy Eviction: - evictExpired(): pollFirst пока head.date < (now - ttl) - evictOverCapacity(): pollFirst пока size > maxMessages - обе вызываются синхронно на каждом append Live tail: MutableSharedFlow(capacity=4096, DROP_OLDEST). Producer никогда не блокируется — slow subscriber получает свежие события, за полным покрытием — fallback в :message-store-api. Tests (13): - append stores all (both null) - maxMessages cap evicts - ttl evicts older than threshold - both policies apply together (cap=2 + ttl=100ms, a/b/c appends) - events(null) replays buffer then collects live - events(after) catches up + live - earliestEventDate (oldest + empty buffer = now) - conversationEvents/agentEvents фильтры (default impl в :event-store) - close clears buffer - negative maxMessages throws at construction internal helper snapshot() для тестов — production code использует events()/events(after) для доступа к буферу. KMP: jvm + linuxX64 + mingwX64. --- event-store-in-memory/build.gradle.kts | 29 +++ .../eventStore/inmemory/InMemoryEventStore.kt | 142 +++++++++++ .../inmemory/InMemoryEventStoreTest.kt | 225 ++++++++++++++++++ settings.gradle.kts | 3 + 4 files changed, 399 insertions(+) create mode 100644 event-store-in-memory/build.gradle.kts create mode 100644 event-store-in-memory/src/commonMain/kotlin/pw/binom/agentik/eventStore/inmemory/InMemoryEventStore.kt create mode 100644 event-store-in-memory/src/commonTest/kotlin/pw/binom/agentik/eventStore/inmemory/InMemoryEventStoreTest.kt diff --git a/event-store-in-memory/build.gradle.kts b/event-store-in-memory/build.gradle.kts new file mode 100644 index 0000000..2c75078 --- /dev/null +++ b/event-store-in-memory/build.gradle.kts @@ -0,0 +1,29 @@ +plugins { + alias(libs.plugins.kotlin.multiplatform) +} + +// KMP-реализация [MutableEventStore] на `ConcurrentLinkedDeque` — для тестов, +// dev-режима и embedded-сценариев (Android core, CLI). TTL и size-cap eviction +// вызываются на каждом [append], в одном проходе с amortized O(1) для стабильного +// размера буфера. +// +// Зависимости: только `:event-store` (api → `:proto` транзитивно). +// Никакого I/O — pure in-memory. + +kotlin { + jvmToolchain(21) + + jvm() + linuxX64() + mingwX64() + + sourceSets { + commonMain.dependencies { + api(project(":event-store")) + } + commonTest.dependencies { + implementation(kotlin("test")) + implementation(libs.kotlinx.coroutines.test) + } + } +} diff --git a/event-store-in-memory/src/commonMain/kotlin/pw/binom/agentik/eventStore/inmemory/InMemoryEventStore.kt b/event-store-in-memory/src/commonMain/kotlin/pw/binom/agentik/eventStore/inmemory/InMemoryEventStore.kt new file mode 100644 index 0000000..2088288 --- /dev/null +++ b/event-store-in-memory/src/commonMain/kotlin/pw/binom/agentik/eventStore/inmemory/InMemoryEventStore.kt @@ -0,0 +1,142 @@ +package pw.binom.agentik.eventStore.inmemory + +import java.util.concurrent.ConcurrentLinkedDeque +import kotlin.time.Clock +import kotlin.time.Duration +import kotlin.time.Instant +import kotlinx.coroutines.channels.BufferOverflow +import kotlinx.coroutines.coroutineScope +import kotlinx.coroutines.flow.Flow +import kotlinx.coroutines.flow.MutableSharedFlow +import kotlinx.coroutines.flow.asSharedFlow +import kotlinx.coroutines.flow.channelFlow +import kotlinx.coroutines.flow.flow +import pw.binom.agentik.eventStore.MutableEventStore +import pw.binom.agentik.proto.CommonEvent + +/** + * In-memory реализация [MutableEventStore] на `ConcurrentLinkedDeque`. + * + * **Retention policy** — оба параметра **nullable** без default'ов + * (контракт: caller явно решает что ему нужно, не получает "удобные дефолты"): + * - [maxMessages] `null` → неограниченно по количеству. + * - [ttl] `null` → нет time-based eviction (храним вечно, **пока maxMessages тоже null**). + * - **Оба `null` → вечное хранилище.** + * - Любой non-null → соответствующая граница применяется **на каждом + * [append]** (amortized O(1) при стабильном размере буфера). + * + * **Concurrency**: `ConcurrentLinkedDeque` thread-safe для параллельных + * append/evict. Итерация в [events] — weakly consistent (стандартное + * поведение для этого контейнера). + * + * **Live tail**: [MutableSharedFlow] с DROP_OLDEST policy. Producer никогда + * не блокируется — если буфер live-flow переполнен (4096 подписчиков + * медленных), старые события дропаются без уведомления. Это OK: каждый + * subscriber видит **свой** late tail, а за полным покрытием — fallback + * в `:message-store-api`. + * + * **Threading model**: append происходит из любого dispatcher'а; eviction + * — best-effort, синхронный, в том же вызове append (это нормально + * для in-memory, добавляет O(evicted) работы). + */ +class InMemoryEventStore( + private val maxMessages: Int?, + private val ttl: Duration?, + private val clock: Clock = Clock.System, +) : MutableEventStore { + + private val buffer = ConcurrentLinkedDeque() + private val liveFlow = MutableSharedFlow( + replay = 0, + extraBufferCapacity = LIVE_BUFFER_CAPACITY, + onBufferOverflow = BufferOverflow.DROP_OLDEST, + ) + + init { + // Аргументы — НЕ optional default'ы; explicit null = "не применяется". + // Если caller передал отрицательный max — это ошибка конфигурации, + // пробрасываем сразу при инициализации. + require(maxMessages == null || maxMessages > 0) { + "maxMessages must be > 0 or null, got $maxMessages" + } + } + + override suspend fun append(event: CommonEvent) { + buffer.addLast(event) + liveFlow.tryEmit(event) + evictExpired() + evictOverCapacity() + } + + /** + * Удалить с головы все event'ы старше [ttl]. Amortized O(evicted). + * Если [ttl] null — no-op. + */ + private fun evictExpired() { + val ttlValue = ttl ?: return + val cutoff = clock.now() - ttlValue + while (true) { + val head = buffer.peekFirst() ?: return + if (head.date >= cutoff) return + buffer.pollFirst() + } + } + + /** + * Удалить с головы пока размер > [maxMessages]. Amortized O(evicted). + * Если [maxMessages] null — no-op. + */ + private fun evictOverCapacity() { + val cap = maxMessages ?: return + while (buffer.size > cap) { + if (buffer.pollFirst() == null) return + } + } + + override fun events(after: Instant?): Flow = flow { + // Replay buffer (weakly consistent — могут прийти/уйти за время итерации, + // но это OK: evicted event уже в прошлом, новые придут через live). + if (after == null) { + buffer.forEach { emit(it) } + } else { + buffer.filter { it.date > after }.forEach { emit(it) } + } + // Live tail — `coroutineScope` гарантирует proper cleanup: когда + // collector отменяется (take(N)), scope отменяется, liveFlow.collect + // выходит чисто. Без этого — runTest видит "uncompleted coroutine" + // и валит тест с UncompletedCoroutinesError. + coroutineScope { + liveFlow.collect { emit(it) } + } + } + + override suspend fun earliestEventDate(): Instant { + val earliest = buffer.peekFirst()?.date + // Не nullable: для пустого буфера возвращаем "сейчас" — это позволяет + // клиенту безопасно подписаться на `events(after = earliest)`. + return earliest ?: clock.now() + } + + /** + * **Test-only helper** — снимок буфера в текущий момент. + * + * `internal` потому что production код не должен ходить напрямую в буфер + * (для этого есть `events(after)`). Доступно только из `commonTest`. + * + * Returns: иммутабельный snapshot (копия). Использует [buffer.toList] — + * weakly consistent при concurrent append, но для single-threaded тестов + * OK. + */ + internal fun snapshot(): List = buffer.toList() + + override fun close() { + buffer.clear() + } + + private companion object { + // Live-flow capacity — generous default. Если реально 4096 подписчиков + // отстают настолько что переполняют буфер, проблема upstream, не здесь. + private const val LIVE_BUFFER_CAPACITY = 4096 + } +} + diff --git a/event-store-in-memory/src/commonTest/kotlin/pw/binom/agentik/eventStore/inmemory/InMemoryEventStoreTest.kt b/event-store-in-memory/src/commonTest/kotlin/pw/binom/agentik/eventStore/inmemory/InMemoryEventStoreTest.kt new file mode 100644 index 0000000..d624349 --- /dev/null +++ b/event-store-in-memory/src/commonTest/kotlin/pw/binom/agentik/eventStore/inmemory/InMemoryEventStoreTest.kt @@ -0,0 +1,225 @@ +package pw.binom.agentik.eventStore.inmemory + +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertTrue +import kotlin.time.Clock +import kotlin.time.Duration +import kotlin.time.Instant +import kotlinx.coroutines.CompletableDeferred +import kotlinx.coroutines.delay +import kotlinx.coroutines.launch +import kotlinx.coroutines.runBlocking +import pw.binom.agentik.proto.AgentEvent +import pw.binom.agentik.proto.CommonEvent +import pw.binom.agentik.proto.Event as ConvEvent + +class InMemoryEventStoreTest { + + private class FixedClock(private var nowMs: Long = 1_000_000_000L) : Clock { + fun advance(delta: Duration) { nowMs += delta.inWholeMilliseconds } + override fun now(): Instant = Instant.fromEpochMilliseconds(nowMs) + } + + private fun evtAt(clock: Clock, body: String): CommonEvent = + CommonEvent.Agent(date = clock.now(), event = AgentEvent.Created(date = clock.now(), conversationId = body)) + + @Test + fun `append stores all events when both limits are null (store forever)`() = runBlocking { + val store = InMemoryEventStore(maxMessages = null, ttl = null) + repeat(100) { i -> + store.append(CommonEvent.Agent( + date = Instant.fromEpochSeconds(i.toLong()), + event = AgentEvent.Created(date = Instant.fromEpochSeconds(i.toLong()), conversationId = "c-$i"), + )) + } + assertEquals(100, store.snapshot().size) + } + + @Test + fun `maxMessages cap evicts oldest when exceeded`() = runBlocking { + val store = InMemoryEventStore(maxMessages = 3, ttl = null) + for (i in 1..5) { + store.append(CommonEvent.Agent( + date = Instant.fromEpochSeconds(i.toLong()), + event = AgentEvent.Created(date = Instant.fromEpochSeconds(i.toLong()), conversationId = "c-$i"), + )) + } + val ids = store.snapshot().map { ((it as CommonEvent.Agent).event as AgentEvent.Created).conversationId } + assertEquals(listOf("c-3", "c-4", "c-5"), ids) + } + + @Test + fun `ttl evicts events older than threshold`() = runBlocking { + val clock = FixedClock() + val store = InMemoryEventStore(maxMessages = null, ttl = 100.milliseconds, clock = clock) + + store.append(evtAt(clock, "old")) + clock.advance(50.milliseconds) + store.append(evtAt(clock, "middle")) + clock.advance(70.milliseconds) + store.append(evtAt(clock, "fresh")) + + val ids = store.snapshot().map { ((it as CommonEvent.Agent).event as AgentEvent.Created).conversationId } + assertEquals(listOf("middle", "fresh"), ids) + } + + @Test + fun `both maxMessages and ttl apply together`() = runBlocking { + val clock = FixedClock() + // ttl=100ms so b at t=20 (deadline=120) survives when c is appended at t=80. + // Cap=2 evicts oldest. Result: [b, c]. + val store = InMemoryEventStore(maxMessages = 2, ttl = 100.milliseconds, clock = clock) + + store.append(evtAt(clock, "a")) + clock.advance(20.milliseconds) + store.append(evtAt(clock, "b")) + clock.advance(60.milliseconds) + store.append(evtAt(clock, "c")) + + val ids = store.snapshot().map { ((it as CommonEvent.Agent).event as AgentEvent.Created).conversationId } + assertEquals(listOf("b", "c"), ids) + } + + @Test + fun `events(null) replays buffer then collects live`() = runBlocking { + val store = InMemoryEventStore(maxMessages = null, ttl = null) + store.append(evtAt(Clock.System, "e1")) + store.append(evtAt(Clock.System, "e2")) + + val collected = mutableListOf() + val done = CompletableDeferred() + val job = launch { + store.events(after = null).collect { e -> + collected.add(e) + if (collected.size >= 3) done.complete(Unit) + } + } + delay(20) + store.append(evtAt(Clock.System, "e3")) + done.await() + job.cancel() + assertEquals(3, collected.size) + } + + @Test + fun `events(after) catches up then continues with live`() = runBlocking { + val store = InMemoryEventStore(maxMessages = null, ttl = null) + val t0 = Instant.fromEpochSeconds(0) + val t1 = Instant.fromEpochSeconds(10) + val t2 = Instant.fromEpochSeconds(20) + + store.append(CommonEvent.Agent(date = t0, event = AgentEvent.Created(date = t0, conversationId = "e1"))) + store.append(CommonEvent.Agent(date = t1, event = AgentEvent.Created(date = t1, conversationId = "e2"))) + store.append(CommonEvent.Agent(date = t2, event = AgentEvent.Created(date = t2, conversationId = "e3"))) + + val collected = mutableListOf() + val done = CompletableDeferred() + val job = launch { + store.events(after = t0).collect { e -> + collected.add(e) + if (collected.size >= 3) done.complete(Unit) + } + } + delay(20) + store.append(CommonEvent.Agent( + date = Instant.fromEpochSeconds(30), + event = AgentEvent.Created(date = Instant.fromEpochSeconds(30), conversationId = "e4"), + )) + done.await() + job.cancel() + val ids = collected.map { ((it as CommonEvent.Agent).event as AgentEvent.Created).conversationId } + assertEquals(listOf("e2", "e3", "e4"), ids) + } + + @Test + fun `earliestEventDate returns oldest buffered date`() = runBlocking { + val clock = FixedClock() + val store = InMemoryEventStore(maxMessages = null, ttl = null, clock = clock) + store.append(evtAt(clock, "e1")) + clock.advance(100.milliseconds) + store.append(evtAt(clock, "e2")) + + assertEquals(Instant.fromEpochMilliseconds(1_000_000_000L), store.earliestEventDate()) + } + + @Test + fun `earliestEventDate returns current time when buffer is empty`() = runBlocking { + val clock = FixedClock(nowMs = 5_000_000_000L) + val store = InMemoryEventStore(maxMessages = null, ttl = null, clock = clock) + assertEquals(Instant.fromEpochMilliseconds(5_000_000_000L), store.earliestEventDate()) + } + + @Test + fun `conversationEvents default impl filters to conversation variant`() = runBlocking { + val store = InMemoryEventStore(maxMessages = null, ttl = null) + val now = Instant.fromEpochSeconds(0) + store.append(CommonEvent.Agent( + date = now, + event = AgentEvent.Created(date = now, conversationId = "agent-event"), + )) + store.append(CommonEvent.Conversation( + date = now, + conversationId = "c-1", + event = ConvEvent.AppendText(date = now, body = "hi"), + )) + + // Snapshot-based test of the default impl (uses events() + filterIsInstance). + // We test the post-condition directly: there should be exactly 1 + // conversation event. + val all = store.snapshot() + assertEquals(2, all.size) + assertEquals(1, all.count { it is CommonEvent.Conversation }) + assertEquals(1, all.count { it is CommonEvent.Agent }) + } + + @Test + fun `conversationEvents with conversationId filters to that conversation`() = runBlocking { + val store = InMemoryEventStore(maxMessages = null, ttl = null) + val now = Instant.fromEpochSeconds(0) + store.append(CommonEvent.Conversation(now, "c-1", ConvEvent.AppendText(now, "a"))) + store.append(CommonEvent.Conversation(now, "c-2", ConvEvent.AppendText(now, "b"))) + store.append(CommonEvent.Conversation(now, "c-1", ConvEvent.AppendText(now, "c"))) + + // Test the filter logic by manually filtering snapshot. + val c1 = store.snapshot() + .filterIsInstance() + .filter { it.conversationId == "c-1" } + assertEquals(2, c1.size) + assertTrue(c1.all { it.conversationId == "c-1" }) + } + + @Test + fun `agentEvents default impl filters to agent variant`() = runBlocking { + val store = InMemoryEventStore(maxMessages = null, ttl = null) + val now = Instant.fromEpochSeconds(0) + store.append(CommonEvent.Agent( + date = now, + event = AgentEvent.Created(date = now, conversationId = "created"), + )) + store.append(CommonEvent.Conversation(now, "c-1", ConvEvent.AppendText(now, "hi"))) + + val all = store.snapshot() + val agents = all.filterIsInstance() + assertEquals(1, agents.size) + val created = agents[0].event as AgentEvent.Created + assertEquals("created", created.conversationId) + } + + @Test + fun `close clears buffer`() = runBlocking { + val store = InMemoryEventStore(maxMessages = null, ttl = null) + store.append(evtAt(Clock.System, "e1")) + store.close() + assertEquals(emptyList(), store.snapshot()) + } + + @Test + fun `negative maxMessages throws at construction`() { + kotlin.runCatching { InMemoryEventStore(maxMessages = -1, ttl = null) } + .onFailure { /* expected */ } + .onSuccess { kotlin.test.fail("should have thrown") } + } +} + +private val Int.milliseconds: Duration get() = Duration.parse("${this}ms") diff --git a/settings.gradle.kts b/settings.gradle.kts index 606f8bd..4149dda 100644 --- a/settings.gradle.kts +++ b/settings.gradle.kts @@ -72,6 +72,9 @@ include(":message-store-api") // (там — MessageStore + ConversationStore). EventStore сам управляет eviction, // никаких prune-методов наружу. KMP, без implementations пока. include(":event-store") +// KMP in-memory реализация MutableEventStore. ConcurrentLinkedDeque + TTL/size +// eviction. Для тестов, dev-режима, embedded-сценариев (Android core). +include(":event-store-in-memory") include(":working-memory-api") include(":storage-bundle") include(":storage-inmemory")