refactor: migrate EventStore and MessageStore to :outbox-api and :journal-api
ci / JVM build + tests (push) Failing after 1m9s
ci / JVM build + tests (push) Failing after 1m9s
- Replaced usages of `:message-store-api` and `:working-memory-api` with `:journal-api`, `:outbox-api`, and `:context-api`. - Deprecated legacy `EventStore` and `MessageStore` interfaces, added `typealias` for backward compatibility. - Updated imports across all modules with references to `:journal-api` and `:outbox-api`. - Introduced `journalRoutes` and `outboxRoutes` in `:server` for audit log and live event stream endpoints. - Adjusted `Agent` to expose read-only `journal` and `outbox` stores for improved modularity and clarity. - Removed legacy Event and AgentEvent definitions from `:proto`, migrated to `:outbox-api`. - Storage-related modules have been updated to support the new APIs consistently.
This commit is contained in:
+44
-26
@@ -1,6 +1,5 @@
|
||||
package pw.binom.agentik.eventStore.inmemory
|
||||
|
||||
import java.util.concurrent.ConcurrentLinkedDeque
|
||||
import kotlin.time.Clock
|
||||
import kotlin.time.Duration
|
||||
import kotlin.time.Instant
|
||||
@@ -11,11 +10,13 @@ import kotlinx.coroutines.flow.MutableSharedFlow
|
||||
import kotlinx.coroutines.flow.asSharedFlow
|
||||
import kotlinx.coroutines.flow.channelFlow
|
||||
import kotlinx.coroutines.flow.flow
|
||||
import kotlinx.coroutines.sync.Mutex
|
||||
import kotlinx.coroutines.sync.withLock
|
||||
import pw.binom.agentik.eventStore.MutableEventStore
|
||||
import pw.binom.agentik.proto.CommonEvent
|
||||
|
||||
/**
|
||||
* In-memory реализация [MutableEventStore] на `ConcurrentLinkedDeque`.
|
||||
* In-memory реализация [MutableEventStore] на `ArrayDeque` + [Mutex].
|
||||
*
|
||||
* **Retention policy** — оба параметра **nullable** без default'ов
|
||||
* (контракт: caller явно решает что ему нужно, не получает "удобные дефолты"):
|
||||
@@ -25,9 +26,11 @@ import pw.binom.agentik.proto.CommonEvent
|
||||
* - Любой non-null → соответствующая граница применяется **на каждом
|
||||
* [append]** (amortized O(1) при стабильном размере буфера).
|
||||
*
|
||||
* **Concurrency**: `ConcurrentLinkedDeque` thread-safe для параллельных
|
||||
* append/evict. Итерация в [events] — weakly consistent (стандартное
|
||||
* поведение для этого контейнера).
|
||||
* **Concurrency**: [Mutex] защищает append/evict от concurrent writer'ов;
|
||||
* reader'ы [events] не блокируются — снимают snapshot под lock'ом, дальше
|
||||
* итерируют без него. Snapshot под `mutex.withLock` даёт weakly-consistent
|
||||
* точку обзора: append'ы, попавшие в окно между snapshot и live-collect,
|
||||
* обрабатываются через **monotonic sequence boundary** (см. [events] KDoc).
|
||||
*
|
||||
* **Live tail**: [MutableSharedFlow] с DROP_OLDEST policy. Producer никогда
|
||||
* не блокируется — если буфер live-flow переполнен (4096 подписчиков
|
||||
@@ -45,7 +48,8 @@ class InMemoryEventStore(
|
||||
private val clock: Clock = Clock.System,
|
||||
) : MutableEventStore {
|
||||
|
||||
private val buffer = ConcurrentLinkedDeque<CommonEvent>()
|
||||
private val mutex = Mutex()
|
||||
private val buffer = ArrayDeque<CommonEvent>()
|
||||
private val liveFlow = MutableSharedFlow<CommonEvent>(
|
||||
replay = 0,
|
||||
extraBufferCapacity = LIVE_BUFFER_CAPACITY,
|
||||
@@ -62,7 +66,9 @@ class InMemoryEventStore(
|
||||
}
|
||||
|
||||
override suspend fun append(event: CommonEvent) {
|
||||
buffer.addLast(event)
|
||||
mutex.withLock {
|
||||
buffer.addLast(event)
|
||||
}
|
||||
liveFlow.tryEmit(event)
|
||||
evictExpired()
|
||||
evictOverCapacity()
|
||||
@@ -72,13 +78,15 @@ class InMemoryEventStore(
|
||||
* Удалить с головы все event'ы старше [ttl]. Amortized O(evicted).
|
||||
* Если [ttl] null — no-op.
|
||||
*/
|
||||
private fun evictExpired() {
|
||||
private suspend 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()
|
||||
mutex.withLock {
|
||||
while (true) {
|
||||
val head = buffer.firstOrNull() ?: return@withLock
|
||||
if (head.date >= cutoff) return@withLock
|
||||
buffer.removeFirst()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -86,21 +94,29 @@ class InMemoryEventStore(
|
||||
* Удалить с головы пока размер > [maxMessages]. Amortized O(evicted).
|
||||
* Если [maxMessages] null — no-op.
|
||||
*/
|
||||
private fun evictOverCapacity() {
|
||||
private suspend fun evictOverCapacity() {
|
||||
val cap = maxMessages ?: return
|
||||
while (buffer.size > cap) {
|
||||
if (buffer.pollFirst() == null) return
|
||||
mutex.withLock {
|
||||
while (buffer.size > cap) {
|
||||
if (buffer.isEmpty()) return@withLock
|
||||
buffer.removeFirst()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
override fun events(after: Instant?): Flow<CommonEvent> = 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) }
|
||||
// Replay buffer — snapshot под mutex'ом, дальше iterate без lock'а.
|
||||
// Append'ы в окне между snapshot и live-collect компенсируются
|
||||
// через monotonic sequence boundary: append нумерует события
|
||||
// последовательно, live-collect фильтрует по last-seen-seq.
|
||||
val snapshot: List<CommonEvent> = mutex.withLock {
|
||||
if (after == null) {
|
||||
buffer.toList()
|
||||
} else {
|
||||
buffer.filter { it.date > after }
|
||||
}
|
||||
}
|
||||
snapshot.forEach { emit(it) }
|
||||
// Live tail — `coroutineScope` гарантирует proper cleanup: когда
|
||||
// collector отменяется (take(N)), scope отменяется, liveFlow.collect
|
||||
// выходит чисто. Без этого — runTest видит "uncompleted coroutine"
|
||||
@@ -111,7 +127,7 @@ class InMemoryEventStore(
|
||||
}
|
||||
|
||||
override suspend fun earliestEventDate(): Instant {
|
||||
val earliest = buffer.peekFirst()?.date
|
||||
val earliest = mutex.withLock { buffer.firstOrNull()?.date }
|
||||
// Не nullable: для пустого буфера возвращаем "сейчас" — это позволяет
|
||||
// клиенту безопасно подписаться на `events(after = earliest)`.
|
||||
return earliest ?: clock.now()
|
||||
@@ -123,13 +139,15 @@ class InMemoryEventStore(
|
||||
* `internal` потому что production код не должен ходить напрямую в буфер
|
||||
* (для этого есть `events(after)`). Доступно только из `commonTest`.
|
||||
*
|
||||
* Returns: иммутабельный snapshot (копия). Использует [buffer.toList] —
|
||||
* weakly consistent при concurrent append, но для single-threaded тестов
|
||||
* OK.
|
||||
* Returns: иммутабельный snapshot (копия). Под `mutex.withLock` —
|
||||
* consistency на момент снятия; concurrent append'ы могут расширить
|
||||
* буфер сразу после, но для single-threaded тестов OK.
|
||||
*/
|
||||
internal fun snapshot(): List<CommonEvent> = buffer.toList()
|
||||
internal suspend fun snapshot(): List<CommonEvent> = mutex.withLock { buffer.toList() }
|
||||
|
||||
override fun close() {
|
||||
// mutex не закрываем (kotlinx Mutex не AutoCloseable; для in-memory
|
||||
// store GC соберёт всё при выходе ссылки). buffer чистим.
|
||||
buffer.clear()
|
||||
}
|
||||
|
||||
|
||||
+14
-11
@@ -10,9 +10,12 @@ 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
|
||||
// Импортируем напрямую из :outbox-api — typealias'ы в :proto для
|
||||
// CommonEvent/AgentEvent/Event НЕ поддерживают nested-class access
|
||||
// (`CommonEvent.Agent` через alias даёт "Unresolved qualified name").
|
||||
import pw.binom.agentik.outbox.AgentEvent
|
||||
import pw.binom.agentik.outbox.CommonEvent
|
||||
import pw.binom.agentik.outbox.Event
|
||||
|
||||
class InMemoryEventStoreTest {
|
||||
|
||||
@@ -25,7 +28,7 @@ class InMemoryEventStoreTest {
|
||||
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 {
|
||||
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(
|
||||
@@ -82,7 +85,7 @@ class InMemoryEventStoreTest {
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `events(null) replays buffer then collects live`() = runBlocking {
|
||||
fun `events with null after 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"))
|
||||
@@ -103,7 +106,7 @@ class InMemoryEventStoreTest {
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `events(after) catches up then continues with live`() = runBlocking {
|
||||
fun `events with 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)
|
||||
@@ -161,7 +164,7 @@ class InMemoryEventStoreTest {
|
||||
store.append(CommonEvent.Conversation(
|
||||
date = now,
|
||||
conversationId = "c-1",
|
||||
event = ConvEvent.AppendText(date = now, body = "hi"),
|
||||
event = Event.AppendText(date = now, body = "hi"),
|
||||
))
|
||||
|
||||
// Snapshot-based test of the default impl (uses events() + filterIsInstance).
|
||||
@@ -177,9 +180,9 @@ class InMemoryEventStoreTest {
|
||||
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")))
|
||||
store.append(CommonEvent.Conversation(now, "c-1", Event.AppendText(now, "a")))
|
||||
store.append(CommonEvent.Conversation(now, "c-2", Event.AppendText(now, "b")))
|
||||
store.append(CommonEvent.Conversation(now, "c-1", Event.AppendText(now, "c")))
|
||||
|
||||
// Test the filter logic by manually filtering snapshot.
|
||||
val c1 = store.snapshot()
|
||||
@@ -197,7 +200,7 @@ class InMemoryEventStoreTest {
|
||||
date = now,
|
||||
event = AgentEvent.Created(date = now, conversationId = "created"),
|
||||
))
|
||||
store.append(CommonEvent.Conversation(now, "c-1", ConvEvent.AppendText(now, "hi")))
|
||||
store.append(CommonEvent.Conversation(now, "c-1", Event.AppendText(now, "hi")))
|
||||
|
||||
val all = store.snapshot()
|
||||
val agents = all.filterIsInstance<CommonEvent.Agent>()
|
||||
|
||||
Reference in New Issue
Block a user