refactor: outbox становится live-only стримом, cursor вынесен в :cursor-api
Семантика outbox'а — теперь чистый live-канал:
* OutboxStore.conversationEvents(id) — больше не принимает after-cursor.
Catchup (replay) делает клиент: journal.listFlow(afterSeq) + подписка
на live. Outbox ответственен только за уведомления «что-то произошло».
* OutboxStore.oldestCursor()/currentCursor()/OutboxGapException — удалены.
Эпоха и offset живут ТОЛЬКО в CursorHolder/CursorStore, переживают
рестарт и инкремент для каждого commit.
* DurableEvent: commit принимает блок { cursor -> MessageRecord }
(cursor выдаёт CursorStore; клиент не вычисляет offset сам).
* MessageRecord больше не несёт cursor — это не его ответственность.
Новые модули:
* :cursor-api — Cursor(epoch, offset) + CursorHolder / MutableCursorHolder
* :cursor-ksqlite — KsqliteCursorHolder (персистентный)
* :cursor-inmemory — для тестов
* :client-sync — LocalSyncAgent (мини-агент поверх :client для десктопа)
Удалены:
* :sync-core — старая референсная реализация, заменена
cursor-разделением и :client.
* outbox-ksqlite — CursorStore/Schema уехали в :cursor-ksqlite.
* OffsetSequencer / PersistentOffsetSequencer / InMemoryOffsetSequencer.
standalone:
* ChatAgent/ConversationLoop/DurableLog/ToolDispatcher/ReflectionScheduler/
ConversationEvents — подписка через push-паттерн (collect событий).
* A2aBridge — currentCursor() и conversationEvents(after=) убраны.
* SqliteStores — cursor_offset удалён из schema v4; seedNextFromJournal
читает MAX(created_at).
* Main.kt — outboxSequencer → outboxCursorHolder; user→agent (:server)
transport удалён; debug-routes удалены; A2A остался.
* Тесты ChatAgentTest/PersistenceTest переписаны на push-паттерн
(subscribe-before-act, snapshot∪live = итоговое состояние). 25/25 + 19/19 ✅
server / client:
* Routes эпоху читают из CursorHolder; снимки несут Cursor? для catchup.
* AgentikAgent и HttpEventStore — те же подписки, без after-параметра.
* ReconnectingOutbox / ReconnectingOutboxTest — без изменений API.
* JournalStore API расширен count(after=Instant?) для unread-badge.
This commit is contained in:
@@ -7,7 +7,8 @@ plugins {
|
||||
// вызываются на каждом `append`, в одном проходе с amortized O(1) для стабильного
|
||||
// размера буфера.
|
||||
//
|
||||
// Зависимости: только `:outbox-api` (api → `:proto` транзитивно).
|
||||
// Зависимости: `:outbox-api` (api → `:proto`/`:cursor-api` транзитивно) и
|
||||
// `:cursor-inmemory` (default-holder для `InMemoryOutboxStore`).
|
||||
// Никакого I/O — pure in-memory.
|
||||
|
||||
kotlin {
|
||||
@@ -20,6 +21,7 @@ kotlin {
|
||||
sourceSets {
|
||||
commonMain.dependencies {
|
||||
api(project(":outbox-api"))
|
||||
implementation(project(":cursor-inmemory"))
|
||||
}
|
||||
commonTest.dependencies {
|
||||
implementation(kotlin("test"))
|
||||
|
||||
-34
@@ -1,34 +0,0 @@
|
||||
package pw.binom.agentik.outbox.inmemory
|
||||
|
||||
import kotlinx.coroutines.sync.Mutex
|
||||
import kotlinx.coroutines.sync.withLock
|
||||
import pw.binom.agentik.outbox.Cursor
|
||||
import pw.binom.agentik.outbox.OffsetSequencer
|
||||
|
||||
/**
|
||||
* In-memory [OffsetSequencer] — для тестов, dev-режима и ephemeral runtime.
|
||||
*
|
||||
* Эпоха генерируется случайно при создании и **не переживает** пересоздание
|
||||
* инстанса: новый store → новый epoch → клиент с прежним курсором получит
|
||||
* [pw.binom.agentik.outbox.OutboxGapException] и сделает resync. Для
|
||||
* production-агента нужен персистентный счётчик (см. `:journal-ksqlite`).
|
||||
*/
|
||||
class InMemoryOffsetSequencer(
|
||||
private val epochId: String = newEpoch(),
|
||||
initialOffset: Long = 0L,
|
||||
) : OffsetSequencer {
|
||||
|
||||
private val mutex = Mutex()
|
||||
private var counter: Long = initialOffset
|
||||
|
||||
override fun epoch(): String = epochId
|
||||
|
||||
override fun current(): Long = counter
|
||||
|
||||
override suspend fun reserve(): Long = mutex.withLock { counter++ }
|
||||
|
||||
companion object {
|
||||
/** Делегирует в [Cursor.newEpoch] — единый генератор epoch'а проекта. */
|
||||
fun newEpoch(): String = Cursor.newEpoch()
|
||||
}
|
||||
}
|
||||
-3
@@ -41,9 +41,6 @@ class InMemoryOnlineOutbox(
|
||||
|
||||
override fun onlineEvents(): Flow<OnlineEvent> = liveFlow
|
||||
|
||||
override fun onlineEvents(conversationId: String): Flow<OnlineEvent> =
|
||||
liveFlow.filter { it.conversationId == conversationId }
|
||||
|
||||
override suspend fun appendOnline(event: OnlineEvent) = liveFlow.emit(event)
|
||||
|
||||
override fun tryAppendOnline(event: OnlineEvent): Boolean = liveFlow.tryEmit(event)
|
||||
|
||||
+24
-131
@@ -1,161 +1,54 @@
|
||||
package pw.binom.agentik.outbox.inmemory
|
||||
|
||||
import kotlin.time.Clock
|
||||
import kotlin.time.Duration
|
||||
import kotlinx.coroutines.channels.BufferOverflow
|
||||
import kotlinx.coroutines.channels.Channel
|
||||
import kotlinx.coroutines.channels.SendChannel
|
||||
import kotlinx.coroutines.flow.Flow
|
||||
import kotlinx.coroutines.flow.channelFlow
|
||||
import kotlinx.coroutines.sync.Mutex
|
||||
import kotlinx.coroutines.sync.withLock
|
||||
import kotlinx.coroutines.flow.MutableSharedFlow
|
||||
import pw.binom.agentik.outbox.CommonEvent
|
||||
import pw.binom.agentik.outbox.Cursor
|
||||
import pw.binom.agentik.outbox.MutableOutboxStore
|
||||
import pw.binom.agentik.outbox.OffsetSequencer
|
||||
import pw.binom.agentik.outbox.OutboxGapException
|
||||
|
||||
/**
|
||||
* In-memory реализация [MutableOutboxStore] на `ArrayDeque` + [Mutex].
|
||||
* In-memory реализация [MutableOutboxStore] — live-шина durable-событий.
|
||||
*
|
||||
* ## Retention
|
||||
* - [maxMessages] `null` → неограниченно по количеству;
|
||||
* - [ttl] `null` → нет time-based eviction;
|
||||
* - **оба `null` → вечное хранилище в RAM**;
|
||||
* - любой non-null → граница применяется на каждом [append] (amortized O(1)).
|
||||
* Никакого хранения и никакого replay: [CommonEvent] публикуются подписчикам
|
||||
* «в моменте». Единственный буфер — внутренний у [MutableSharedFlow] для
|
||||
* развязки продюсера/подписчиков; при переполнении **старые дропаются**
|
||||
* ([BufferOverflow.DROP_OLDEST]), [append] не блокируется.
|
||||
*
|
||||
* ## Курсор и gap
|
||||
* Offset'ы берутся из [sequencer] ([reserveOffset] = [OffsetSequencer.reserve]).
|
||||
* [currentCursor] = offset последнего **записанного** события; [oldestCursor] =
|
||||
* `buffer.first().offset - 1` (или `lastOffset`, если буфер пуст). Подписка
|
||||
* `after` старше `oldestCursor` (или с чужой эпохой) бросает [OutboxGapException].
|
||||
* История событий для перезапроса живёт в
|
||||
* [pw.binom.agentik.journal.JournalStore]; курсор/resume для удалённого
|
||||
* доступа — в sync-слое. Здесь их нет.
|
||||
*
|
||||
* На старте `lastOffset = sequencer.current() - 1`: если счётчик персистентный
|
||||
* и равен N, то клиент с курсором N-1 (догнавший состояние до рестарта)
|
||||
* продолжает инкрементально, а клиент с курсором < N-1 получает gap и делает
|
||||
* resync. Так рестарт сервера не теряет события молча.
|
||||
*
|
||||
* ## Concurrency
|
||||
* Один [Mutex] защищает append/evict/подписки. Регистрация подписчика и снятие
|
||||
* snapshot'а идут **одним критическим участком** — это закрывает окно
|
||||
* «snapshot → live», в котором append мог потеряться: всё, что попадёт в буфер
|
||||
* после регистрации, доедет до подписчика через его [Channel]. Snapshot
|
||||
* итерируется и эмитится вне lock'а.
|
||||
*
|
||||
* ## Live-tail
|
||||
* Каждому подписчику — свой [Channel] с `DROP_OLDEST`: медленный подписчик
|
||||
* теряет только хвост live-потока и обязан сам сделать resync при обнаружении
|
||||
* gap'а по retention'у.
|
||||
* **Маршрутизация**: один общий [MutableSharedFlow]; [conversationEvents] /
|
||||
* [agentEvents] фильтруют его (см. [MutableOutboxStore]).
|
||||
*/
|
||||
class InMemoryOutboxStore(
|
||||
private val maxMessages: Int?,
|
||||
private val ttl: Duration?,
|
||||
private val clock: Clock = Clock.System,
|
||||
private val sequencer: OffsetSequencer = InMemoryOffsetSequencer(),
|
||||
liveBufferCapacity: Int = DEFAULT_LIVE_BUFFER_CAPACITY,
|
||||
) : MutableOutboxStore {
|
||||
|
||||
private val mutex = Mutex()
|
||||
private val buffer = ArrayDeque<CommonEvent>()
|
||||
private val subscribers = mutableSetOf<SendChannel<CommonEvent>>()
|
||||
|
||||
/**
|
||||
* Offset последнего **записанного** события. Инициализируется из счётчика:
|
||||
* `current() - 1` (для fresh-счётчика это `-1`).
|
||||
*/
|
||||
private var lastOffset: Long = sequencer.current() - 1
|
||||
private val liveFlow = MutableSharedFlow<CommonEvent>(
|
||||
replay = 0,
|
||||
extraBufferCapacity = liveBufferCapacity,
|
||||
onBufferOverflow = BufferOverflow.DROP_OLDEST,
|
||||
)
|
||||
|
||||
init {
|
||||
require(maxMessages == null || maxMessages > 0) {
|
||||
"maxMessages must be > 0 or null, got $maxMessages"
|
||||
require(liveBufferCapacity > 0) {
|
||||
"liveBufferCapacity must be > 0, got $liveBufferCapacity"
|
||||
}
|
||||
}
|
||||
|
||||
override suspend fun reserveOffset(): Long = sequencer.reserve()
|
||||
override fun events(): Flow<CommonEvent> = liveFlow
|
||||
|
||||
override suspend fun append(event: CommonEvent) {
|
||||
mutex.withLock {
|
||||
require(event.offset > lastOffset) {
|
||||
"Non-monotonic offset: got ${event.offset}, last=${lastOffset}"
|
||||
}
|
||||
buffer.addLast(event)
|
||||
lastOffset = event.offset
|
||||
subscribers.forEach { it.trySend(event) }
|
||||
evictLocked()
|
||||
}
|
||||
liveFlow.emit(event)
|
||||
}
|
||||
|
||||
private fun evictLocked() {
|
||||
val ttlValue = ttl
|
||||
if (ttlValue != null) {
|
||||
val cutoff = clock.now() - ttlValue
|
||||
while (true) {
|
||||
val head = buffer.firstOrNull() ?: break
|
||||
if (head.date >= cutoff) break
|
||||
buffer.removeFirst()
|
||||
}
|
||||
}
|
||||
val cap = maxMessages
|
||||
if (cap != null) {
|
||||
while (buffer.size > cap) {
|
||||
buffer.removeFirst()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
override fun events(after: Cursor?): Flow<CommonEvent> = channelFlow {
|
||||
val channel = Channel<CommonEvent>(
|
||||
capacity = LIVE_BUFFER_CAPACITY,
|
||||
onBufferOverflow = BufferOverflow.DROP_OLDEST,
|
||||
)
|
||||
val snapshot: List<CommonEvent> = mutex.withLock {
|
||||
val epoch = sequencer.epoch()
|
||||
if (after != null) {
|
||||
val floor = buffer.firstOrNull()?.let { it.offset - 1 } ?: lastOffset
|
||||
if (after.epoch != epoch || after.offset < floor || after.offset > lastOffset) {
|
||||
throw OutboxGapException(
|
||||
requested = after,
|
||||
current = Cursor(epoch, lastOffset),
|
||||
oldest = Cursor(epoch, floor),
|
||||
)
|
||||
}
|
||||
}
|
||||
subscribers += channel
|
||||
// `after == null` → live-only (без replay буфера). Чтобы получить
|
||||
// весь удержанный хвост, клиент передаёт `after = oldestCursor()`.
|
||||
if (after == null) emptyList() else buffer.filter { it.offset > after.offset }
|
||||
}
|
||||
try {
|
||||
snapshot.forEach { send(it) }
|
||||
for (event in channel) send(event)
|
||||
} finally {
|
||||
mutex.withLock { subscribers -= channel }
|
||||
channel.close()
|
||||
}
|
||||
}
|
||||
|
||||
override suspend fun currentCursor(): Cursor = mutex.withLock {
|
||||
Cursor(sequencer.epoch(), lastOffset)
|
||||
}
|
||||
|
||||
override suspend fun oldestCursor(): Cursor = mutex.withLock {
|
||||
val floor = buffer.firstOrNull()?.let { it.offset - 1 } ?: lastOffset
|
||||
Cursor(sequencer.epoch(), floor)
|
||||
}
|
||||
|
||||
/**
|
||||
* **Test-only helper** — снимок буфера в текущий момент.
|
||||
*
|
||||
* `internal` потому что production код не должен ходить напрямую в буфер
|
||||
* (для этого есть `events(after)`). Доступно только из `commonTest`.
|
||||
*/
|
||||
internal suspend fun snapshot(): List<CommonEvent> = mutex.withLock { buffer.toList() }
|
||||
|
||||
override fun close() {
|
||||
buffer.clear()
|
||||
subscribers.clear()
|
||||
// replay = 0 — чистить нечего; сам flow соберётся GC'ом при выходе ссылки.
|
||||
// Идемпотентно: повторный close() безопасен.
|
||||
}
|
||||
|
||||
private companion object {
|
||||
private const val LIVE_BUFFER_CAPACITY = 4096
|
||||
private const val DEFAULT_LIVE_BUFFER_CAPACITY = 64
|
||||
}
|
||||
}
|
||||
|
||||
+73
-248
@@ -1,292 +1,117 @@
|
||||
package pw.binom.agentik.outbox.inmemory
|
||||
|
||||
import kotlinx.coroutines.flow.toList
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlinx.coroutines.test.runCurrent
|
||||
import kotlinx.coroutines.test.runTest
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertFailsWith
|
||||
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
|
||||
// Импортируем напрямую из :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.Cursor
|
||||
import pw.binom.agentik.outbox.DurableEvent
|
||||
import pw.binom.agentik.outbox.OutboxGapException
|
||||
import kotlin.time.Instant
|
||||
|
||||
/**
|
||||
* Тесты [InMemoryOutboxStore] как live-шины durable-событий: доставка живым
|
||||
* подписчикам, отсутствие replay, фильтры по варианту/диалогу.
|
||||
*/
|
||||
class InMemoryOutboxStoreTest {
|
||||
|
||||
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 agentEvent(conversationId: String, at: Instant = Instant.fromEpochMilliseconds(1)) =
|
||||
CommonEvent.Agent(date = at, event = AgentEvent.Created(date = at, conversationId = conversationId))
|
||||
|
||||
private fun agentEvent(offset: Long, conversationId: String, at: Instant = Instant.fromEpochSeconds(offset)) =
|
||||
CommonEvent.Agent(
|
||||
date = at,
|
||||
offset = offset,
|
||||
event = AgentEvent.Created(date = at, conversationId = conversationId),
|
||||
)
|
||||
|
||||
private fun evtAt(clock: Clock, offset: Long, body: String): CommonEvent =
|
||||
CommonEvent.Agent(
|
||||
date = clock.now(),
|
||||
offset = offset,
|
||||
event = AgentEvent.Created(date = clock.now(), conversationId = body),
|
||||
)
|
||||
|
||||
private suspend fun InMemoryOutboxStore.ids() =
|
||||
snapshot().map { ((it as CommonEvent.Agent).event as AgentEvent.Created).conversationId }
|
||||
private fun conversationEvent(conversationId: String, at: Instant = Instant.fromEpochMilliseconds(1)) =
|
||||
CommonEvent.Conversation(date = at, conversationId = conversationId, event = DurableEvent.Interrupted(date = at))
|
||||
|
||||
@Test
|
||||
fun `append stores all events when both limits are null store-forever`() = runBlocking {
|
||||
val store = InMemoryOutboxStore(maxMessages = null, ttl = null)
|
||||
repeat(100) { i ->
|
||||
store.append(agentEvent(i.toLong(), "c-$i"))
|
||||
}
|
||||
assertEquals(100, store.snapshot().size)
|
||||
fun `events delivers appended events live`() = runTest {
|
||||
val store = InMemoryOutboxStore()
|
||||
val collected = mutableListOf<CommonEvent>()
|
||||
val job = backgroundScope.launch { store.events().collect { collected.add(it) } }
|
||||
runCurrent()
|
||||
|
||||
store.append(agentEvent("c1"))
|
||||
store.append(agentEvent("c2"))
|
||||
runCurrent()
|
||||
|
||||
assertEquals(listOf("c1", "c2"), collected.map { (it as CommonEvent.Agent).event.let { e -> (e as AgentEvent.Created).conversationId } })
|
||||
job.cancel()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `reserveOffset is monotonic`() = runBlocking {
|
||||
val store = InMemoryOutboxStore(maxMessages = null, ttl = null)
|
||||
assertEquals(0L, store.reserveOffset())
|
||||
assertEquals(1L, store.reserveOffset())
|
||||
assertEquals(2L, store.reserveOffset())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `maxMessages cap evicts oldest when exceeded`() = runBlocking {
|
||||
val store = InMemoryOutboxStore(maxMessages = 3, ttl = null)
|
||||
for (i in 1..5) store.append(agentEvent(i.toLong(), "c-$i"))
|
||||
assertEquals(listOf("c-3", "c-4", "c-5"), store.ids())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `ttl evicts events older than threshold`() = runBlocking {
|
||||
val clock = FixedClock()
|
||||
val store = InMemoryOutboxStore(maxMessages = null, ttl = 100.milliseconds, clock = clock)
|
||||
|
||||
store.append(evtAt(clock, 0, "old"))
|
||||
clock.advance(50.milliseconds)
|
||||
store.append(evtAt(clock, 1, "middle"))
|
||||
clock.advance(70.milliseconds)
|
||||
store.append(evtAt(clock, 2, "fresh"))
|
||||
|
||||
assertEquals(listOf("middle", "fresh"), store.ids())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `both maxMessages and ttl apply together`() = runBlocking {
|
||||
val clock = FixedClock()
|
||||
val store = InMemoryOutboxStore(maxMessages = 2, ttl = 100.milliseconds, clock = clock)
|
||||
|
||||
store.append(evtAt(clock, 0, "a"))
|
||||
clock.advance(20.milliseconds)
|
||||
store.append(evtAt(clock, 1, "b"))
|
||||
clock.advance(60.milliseconds)
|
||||
store.append(evtAt(clock, 2, "c"))
|
||||
|
||||
assertEquals(listOf("b", "c"), store.ids())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `currentCursor is last appended offset`() = runBlocking {
|
||||
val store = InMemoryOutboxStore(maxMessages = null, ttl = null)
|
||||
store.append(agentEvent(0, "e0"))
|
||||
store.append(agentEvent(1, "e1"))
|
||||
assertEquals(1L, store.currentCursor().offset)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `oldestCursor equals last offset when buffer is empty`() = runBlocking {
|
||||
val store = InMemoryOutboxStore(maxMessages = null, ttl = null)
|
||||
// Свежий счётчик: next offset = 0 → oldest = -1.
|
||||
assertEquals(-1L, store.oldestCursor().offset)
|
||||
assertEquals(store.currentCursor().epoch, store.oldestCursor().epoch)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `oldestCursor is first-minus-one when buffer is non-empty`() = runBlocking {
|
||||
val store = InMemoryOutboxStore(maxMessages = 2, ttl = null)
|
||||
store.append(agentEvent(0, "a"))
|
||||
store.append(agentEvent(1, "b"))
|
||||
store.append(agentEvent(2, "c"))
|
||||
// buffer = [1, 2]; oldest = 1 - 1 = 0.
|
||||
assertEquals(0L, store.oldestCursor().offset)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `events with null after is live-only (no replay)`() = runBlocking {
|
||||
val store = InMemoryOutboxStore(maxMessages = null, ttl = null)
|
||||
store.append(agentEvent(0, "e1"))
|
||||
store.append(agentEvent(1, "e2"))
|
||||
fun `events does not replay history`() = runTest {
|
||||
val store = InMemoryOutboxStore()
|
||||
store.append(agentEvent("old"))
|
||||
|
||||
val collected = mutableListOf<CommonEvent>()
|
||||
val done = CompletableDeferred<Unit>()
|
||||
val job = launch {
|
||||
store.events(after = null).collect { e ->
|
||||
collected.add(e)
|
||||
done.complete(Unit)
|
||||
}
|
||||
}
|
||||
delay(20)
|
||||
store.append(agentEvent(2, "e3"))
|
||||
done.await()
|
||||
val job = backgroundScope.launch { store.events().collect { collected.add(it) } }
|
||||
runCurrent()
|
||||
assertEquals(emptyList(), collected)
|
||||
|
||||
store.append(agentEvent("new"))
|
||||
runCurrent()
|
||||
assertEquals(1, collected.size)
|
||||
job.cancel()
|
||||
assertEquals(listOf("e3"), collected.map { ((it as CommonEvent.Agent).event as AgentEvent.Created).conversationId })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `events with cursor catches up then continues with live`() = runBlocking {
|
||||
val store = InMemoryOutboxStore(maxMessages = null, ttl = null)
|
||||
store.append(agentEvent(0, "e1"))
|
||||
store.append(agentEvent(1, "e2"))
|
||||
store.append(agentEvent(2, "e3"))
|
||||
fun `conversationEvents filters to conversation variant`() = runTest {
|
||||
val store = InMemoryOutboxStore()
|
||||
val collected = mutableListOf<CommonEvent.Conversation>()
|
||||
val job = backgroundScope.launch { store.conversationEvents().collect { collected.add(it) } }
|
||||
runCurrent()
|
||||
|
||||
val from = Cursor(store.currentCursor().epoch, 0L)
|
||||
val collected = mutableListOf<CommonEvent>()
|
||||
val done = CompletableDeferred<Unit>()
|
||||
val job = launch {
|
||||
store.events(after = from).collect { e ->
|
||||
collected.add(e)
|
||||
if (collected.size >= 3) done.complete(Unit)
|
||||
}
|
||||
}
|
||||
delay(20)
|
||||
store.append(agentEvent(3, "e4"))
|
||||
done.await()
|
||||
store.append(agentEvent("agent"))
|
||||
store.append(conversationEvent("c1"))
|
||||
runCurrent()
|
||||
|
||||
assertEquals(listOf("c1"), collected.map { it.conversationId })
|
||||
job.cancel()
|
||||
val ids = collected.map { ((it as CommonEvent.Agent).event as AgentEvent.Created).conversationId }
|
||||
assertEquals(listOf("e2", "e3", "e4"), ids)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `subscribe from currentCursor receives only newer events - no handoff loss`() = runBlocking {
|
||||
val store = InMemoryOutboxStore(maxMessages = null, ttl = null)
|
||||
store.append(agentEvent(0, "old"))
|
||||
fun `conversationEvents with id filters by conversation`() = runTest {
|
||||
val store = InMemoryOutboxStore()
|
||||
val collected = mutableListOf<CommonEvent.Conversation>()
|
||||
val job = backgroundScope.launch { store.conversationEvents("c1").collect { collected.add(it) } }
|
||||
runCurrent()
|
||||
|
||||
val collected = mutableListOf<CommonEvent>()
|
||||
val from = store.currentCursor()
|
||||
val job = launch { store.events(after = from).collect { collected.add(it) } }
|
||||
delay(50)
|
||||
store.append(agentEvent(1, "new"))
|
||||
delay(50)
|
||||
store.append(conversationEvent("c1"))
|
||||
store.append(conversationEvent("c2"))
|
||||
store.append(conversationEvent("c1"))
|
||||
runCurrent()
|
||||
|
||||
assertEquals(2, collected.size)
|
||||
assertEquals(listOf("c1", "c1"), collected.map { it.conversationId })
|
||||
job.cancel()
|
||||
|
||||
assertEquals(listOf(1L), collected.map { it.offset })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `gap exception when cursor older than oldest`() = runBlocking {
|
||||
val store = InMemoryOutboxStore(maxMessages = 2, ttl = null)
|
||||
store.append(agentEvent(0, "e0"))
|
||||
store.append(agentEvent(1, "e1"))
|
||||
store.append(agentEvent(2, "e2"))
|
||||
fun `agentEvents filters to agent variant`() = runTest {
|
||||
val store = InMemoryOutboxStore()
|
||||
val collected = mutableListOf<CommonEvent.Agent>()
|
||||
val job = backgroundScope.launch { store.agentEvents().collect { collected.add(it) } }
|
||||
runCurrent()
|
||||
|
||||
val tooOld = Cursor(store.currentCursor().epoch, -1L)
|
||||
assertFailsWith<OutboxGapException> {
|
||||
store.events(after = tooOld).collect { }
|
||||
}
|
||||
// Ровно на границе — ещё можно.
|
||||
val atFloor = Cursor(store.currentCursor().epoch, 0L)
|
||||
val got = mutableListOf<Long>()
|
||||
val job = launch { store.events(after = atFloor).collect { got.add(it.offset) } }
|
||||
delay(30)
|
||||
store.append(conversationEvent("c1"))
|
||||
store.append(agentEvent("agent"))
|
||||
runCurrent()
|
||||
|
||||
assertEquals(1, collected.size)
|
||||
assertEquals("agent", (collected[0].event as AgentEvent.Created).conversationId)
|
||||
job.cancel()
|
||||
assertEquals(listOf(1L, 2L), got)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `gap exception on epoch mismatch`(): Unit = runBlocking {
|
||||
val store = InMemoryOutboxStore(maxMessages = null, ttl = null)
|
||||
store.append(agentEvent(0, "e0"))
|
||||
val foreign = Cursor("some-other-epoch", 0L)
|
||||
assertFailsWith<OutboxGapException> {
|
||||
store.events(after = foreign).collect { }
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `gap exception exposes requested current and oldest`() = runBlocking {
|
||||
val store = InMemoryOutboxStore(maxMessages = 1, ttl = null)
|
||||
store.append(agentEvent(0, "e0"))
|
||||
store.append(agentEvent(1, "e1"))
|
||||
val tooOld = Cursor(store.currentCursor().epoch, -1L)
|
||||
val e = assertFailsWith<OutboxGapException> { store.events(tooOld).collect { } }
|
||||
assertEquals(tooOld, e.requested)
|
||||
assertEquals(1L, e.current.offset)
|
||||
assertEquals(0L, e.oldest.offset)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `conversationEvents default impl filters to conversation variant`() = runBlocking {
|
||||
val store = InMemoryOutboxStore(maxMessages = null, ttl = null)
|
||||
val now = Instant.fromEpochSeconds(0)
|
||||
store.append(CommonEvent.Agent(now, 0, AgentEvent.Created(now, "agent-event")))
|
||||
store.append(CommonEvent.Conversation(now, 1, "c-1", DurableEvent.Interrupted(now)))
|
||||
|
||||
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 = InMemoryOutboxStore(maxMessages = null, ttl = null)
|
||||
val now = Instant.fromEpochSeconds(0)
|
||||
store.append(CommonEvent.Conversation(now, 0, "c-1", DurableEvent.Interrupted(now)))
|
||||
store.append(CommonEvent.Conversation(now, 1, "c-2", DurableEvent.Interrupted(now)))
|
||||
store.append(CommonEvent.Conversation(now, 2, "c-1", DurableEvent.Interrupted(now)))
|
||||
|
||||
val c1 = store.snapshot()
|
||||
.filterIsInstance<CommonEvent.Conversation>()
|
||||
.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 = InMemoryOutboxStore(maxMessages = null, ttl = null)
|
||||
val now = Instant.fromEpochSeconds(0)
|
||||
store.append(CommonEvent.Agent(now, 0, AgentEvent.Created(now, "created")))
|
||||
store.append(CommonEvent.Conversation(now, 1, "c-1", DurableEvent.Interrupted(now)))
|
||||
|
||||
val agents = store.snapshot().filterIsInstance<CommonEvent.Agent>()
|
||||
assertEquals(1, agents.size)
|
||||
assertEquals("created", (agents[0].event as AgentEvent.Created).conversationId)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `append with non-monotonic offset throws`(): Unit = runBlocking {
|
||||
val store = InMemoryOutboxStore(maxMessages = null, ttl = null)
|
||||
store.append(agentEvent(5, "e5"))
|
||||
assertFailsWith<IllegalArgumentException> { store.append(agentEvent(5, "again")) }
|
||||
assertFailsWith<IllegalArgumentException> { store.append(agentEvent(4, "lower")) }
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `close clears buffer`() = runBlocking {
|
||||
val store = InMemoryOutboxStore(maxMessages = null, ttl = null)
|
||||
store.append(agentEvent(0, "e1"))
|
||||
fun `close is idempotent`() {
|
||||
val store = InMemoryOutboxStore()
|
||||
store.close()
|
||||
store.close()
|
||||
assertEquals(emptyList(), store.snapshot())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `negative maxMessages throws at construction`() {
|
||||
kotlin.runCatching { InMemoryOutboxStore(maxMessages = -1, ttl = null) }
|
||||
.onFailure { /* expected */ }
|
||||
.onSuccess { kotlin.test.fail("should have thrown") }
|
||||
fun `non-positive buffer capacity throws`() {
|
||||
assertFailsWith<IllegalArgumentException> { InMemoryOutboxStore(liveBufferCapacity = 0) }
|
||||
}
|
||||
}
|
||||
|
||||
private val Int.milliseconds: Duration get() = Duration.parse("${this}ms")
|
||||
|
||||
Reference in New Issue
Block a user