diff --git a/agentik-cli/src/commonMain/kotlin/pw/binom/agentik/cli/commands/SendSubcommand.kt b/agentik-cli/src/commonMain/kotlin/pw/binom/agentik/cli/commands/SendSubcommand.kt index ba52b09..1aa6b8f 100644 --- a/agentik-cli/src/commonMain/kotlin/pw/binom/agentik/cli/commands/SendSubcommand.kt +++ b/agentik-cli/src/commonMain/kotlin/pw/binom/agentik/cli/commands/SendSubcommand.kt @@ -9,8 +9,8 @@ import kotlinx.coroutines.launch import pw.binom.agentik.cli.AgentikSubcommand import pw.binom.agentik.cli.defaultCliHttpClient import pw.binom.agentik.client.AgentikAgent +import pw.binom.agentik.outbox.Event import pw.binom.agentik.proto.Content -import pw.binom.agentik.proto.Event import kotlin.time.Instant class SendSubcommand : AgentikSubcommand("send", "Отправить user-ход и стримить ответ") { diff --git a/agentik-tui/src/commonMain/kotlin/pw/binom/agentik/tui/TuiBackend.kt b/agentik-tui/src/commonMain/kotlin/pw/binom/agentik/tui/TuiBackend.kt index 90e035b..dca80bd 100644 --- a/agentik-tui/src/commonMain/kotlin/pw/binom/agentik/tui/TuiBackend.kt +++ b/agentik-tui/src/commonMain/kotlin/pw/binom/agentik/tui/TuiBackend.kt @@ -46,7 +46,11 @@ internal class TuiBackend( state.postSystem("подключено к ${state.config.server}") scope.launch { try { - agent.events(Instant.DISTANT_PAST).collect { /* sidebar refresh */ } + // agent.outbox.agentEvents(after) возвращает Flow; + // распаковываем .event для получения AgentEvent (раньше был + // отдельный метод agent.events(), теперь упразднён — события + // живут в outbox-сущности). + agent.outbox.agentEvents(Instant.DISTANT_PAST).collect { /* sidebar refresh */ } } catch (_: kotlinx.coroutines.CancellationException) { // штатная отмена при закрытии UI } catch (e: Exception) { diff --git a/agentik-tui/src/commonTest/kotlin/pw/binom/agentik/tui/FakeAgent.kt b/agentik-tui/src/commonTest/kotlin/pw/binom/agentik/tui/FakeAgent.kt index e260775..4475422 100644 --- a/agentik-tui/src/commonTest/kotlin/pw/binom/agentik/tui/FakeAgent.kt +++ b/agentik-tui/src/commonTest/kotlin/pw/binom/agentik/tui/FakeAgent.kt @@ -3,8 +3,9 @@ package pw.binom.agentik.tui import kotlinx.coroutines.flow.Flow import kotlinx.coroutines.flow.MutableSharedFlow import kotlinx.coroutines.flow.emptyFlow +import pw.binom.agentik.journal.JournalStore +import pw.binom.agentik.outbox.OutboxStore import pw.binom.agentik.proto.Agent -import pw.binom.agentik.proto.AgentEvent import pw.binom.agentik.proto.Content import pw.binom.agentik.proto.Conversation import pw.binom.agentik.proto.Event @@ -25,6 +26,16 @@ internal class FakeAgent( private set val conversations = mutableListOf() + // Storage handles не используются тестами TuiBackend — тесты проверяют + // маршрутизацию Conversation.events в UI state. Outbox stub-ы возвращают + // emptyFlow, journal — error-on-access (никто не должен его трогать). + override val journal: JournalStore = error("journal not used in TuiBackend tests") + override val outbox: OutboxStore = object : OutboxStore { + override fun events(after: Instant?) = emptyFlow() + override suspend fun earliestEventDate(): Instant = Instant.DISTANT_PAST + override fun close() {} + } + override fun createConversation(temp: Boolean): Conversation { createCount++ val c = conversationFactory() @@ -40,8 +51,6 @@ internal class FakeAgent( override suspend fun getConversations(offset: Int, limit: Int): List = conversations.toList() - - override fun events(after: Instant): Flow = emptyFlow() } /** diff --git a/client/src/commonMain/kotlin/pw/binom/agentik/client/AgentClient.kt b/client/src/commonMain/kotlin/pw/binom/agentik/client/AgentClient.kt index fcc9adf..4617919 100644 --- a/client/src/commonMain/kotlin/pw/binom/agentik/client/AgentClient.kt +++ b/client/src/commonMain/kotlin/pw/binom/agentik/client/AgentClient.kt @@ -4,23 +4,17 @@ import io.ktor.client.HttpClient import io.ktor.client.call.body import io.ktor.client.request.delete import io.ktor.client.request.get -import io.ktor.client.request.prepareGet import io.ktor.client.request.parameter import io.ktor.client.request.post import io.ktor.client.request.setBody -import io.ktor.client.statement.bodyAsChannel import io.ktor.http.ContentType import io.ktor.http.HttpStatusCode import io.ktor.http.contentType -import kotlinx.coroutines.flow.Flow -import kotlinx.coroutines.flow.flow import kotlinx.coroutines.runBlocking -import pw.binom.agentik.eventStore.EventStore +import pw.binom.agentik.journal.JournalStore +import pw.binom.agentik.outbox.OutboxStore import pw.binom.agentik.proto.Agent -import pw.binom.agentik.proto.AgentEvent -import pw.binom.agentik.proto.CommonEvent import pw.binom.agentik.proto.Conversation -import kotlin.time.Instant /** * HTTP-реализация [Agent]. Ходит в `:server`-фасад, см. `agentikAgent(...)`. @@ -30,6 +24,12 @@ import kotlin.time.Instant * POST `/conversations`. Используем `runBlocking` — это одноразовая * операция (открытие чата), не горячий путь. В UI-контексте вызывающий сам * решает, что делать. + * + * **Storage handles** ([journal], [outbox]) — read-only views на серверные + * хранилища. [outbox] уже реализован ([HttpEventStore]); [journal] — + * заглушка, потому что соответствующий HTTP endpoint'ы (`/journal/...`) + * ещё не выставлены на стороне `:server`. После их добавления подменить + * `error(...)` на `HttpJournalStore(...)`. */ internal class AgentClient( private val httpClient: HttpClient, @@ -43,7 +43,16 @@ internal class AgentClient( * Единый канал событий (lifecycle + per-conversation). Под капотом — * [HttpEventStore]: каждый метод бьёт свой URL (см. KDoc). */ - val eventStore: EventStore = HttpEventStore(httpClient = httpClient, baseUrl = agentUrl) + override val outbox: OutboxStore = HttpEventStore(httpClient = httpClient, baseUrl = agentUrl) + + /** + * HTTP-фасад для journal пока не реализован: на стороне `:server` ещё + * не выставлены endpoint'ы `/journal/conversations/{id}/messages`. + * Как только появятся — заменить на `HttpJournalStore(httpClient, agentUrl)`. + */ + override val journal: JournalStore = error( + "HttpJournalStore ещё не реализован — дождаться :server endpoint'а /journal/...", + ) override fun createConversation(temp: Boolean): Conversation = runBlocking { @@ -73,35 +82,4 @@ internal class AgentClient( }.body>() return snapshots.map { ConversationClient(httpClient, agentUrl, it) } } - - override fun events(after: Instant): Flow = flow { - httpClient.prepareGet("$agentUrl/events?after=$after") { noSseReadTimeout() } - .execute { response -> - check(response.status == HttpStatusCode.OK) { - "events: server returned ${response.status}" - } - readSse(response.bodyAsChannel()) - .collect { payload -> - emit(agentikJson.decodeFromString(AgentEvent.serializer(), payload)) - } - } - } - - - /** - * Подписка на ВСЕ события: agent lifecycle + все conversation events. - * Использует SSE endpoint /events/all. - */ - override fun allEvents(after: Instant): Flow = flow { - httpClient.prepareGet("$agentUrl/events/all?after=$after") { noSseReadTimeout() } - .execute { response -> - check(response.status == HttpStatusCode.OK) { - "allEvents: server returned ${response.status}" - } - readSse(response.bodyAsChannel()) - .collect { payload -> - emit(agentikJson.decodeFromString(CommonEvent.serializer(), payload)) - } - } - } } diff --git a/client/src/commonMain/kotlin/pw/binom/agentik/client/HttpEventStore.kt b/client/src/commonMain/kotlin/pw/binom/agentik/client/HttpEventStore.kt index 3f8359a..b1ad455 100644 --- a/client/src/commonMain/kotlin/pw/binom/agentik/client/HttpEventStore.kt +++ b/client/src/commonMain/kotlin/pw/binom/agentik/client/HttpEventStore.kt @@ -7,14 +7,15 @@ import io.ktor.http.HttpStatusCode import kotlinx.coroutines.flow.Flow import kotlinx.coroutines.flow.flow import pw.binom.agentik.eventStore.EventStore -import pw.binom.agentik.proto.AgentEvent -import pw.binom.agentik.proto.CommonEvent -import pw.binom.agentik.proto.Event +import pw.binom.agentik.outbox.AgentEvent +import pw.binom.agentik.outbox.CommonEvent +import pw.binom.agentik.outbox.Event import kotlin.time.Clock import kotlin.time.Instant /** - * HTTP-реализация [EventStore], ходящая в `:server`-фасад. + * HTTP-реализация [EventStore] (= [pw.binom.agentik.outbox.OutboxStore]), + * ходящая в `:server`-фасад. * * **Хитрый план**: вместо того, чтобы все методы шли в один общий endpoint и * фильтровали client-side ([EventStore.events]/[filterIsInstance]), эта @@ -36,6 +37,12 @@ import kotlin.time.Instant * (см. KDoc [EventStore.earliestEventDate] — для пустого буфера это и есть * контрактное значение). Клиент, который полагался на gap detection через * message store, продолжит работать — просто fallback никогда не сработает. + * + * **Импорты [CommonEvent]/[AgentEvent]/[Event] идут напрямую из + * `pw.binom.agentik.outbox`** — typealias'ы в `:proto.CommonEvent` и т.п. + * НЕ поддерживают nested-class access (`CommonEvent.Agent` через alias + * даёт "Unresolved qualified name"), поэтому приходится использовать + * конкретный пакет. Типы идентичны, alias только для удобства внешнего API. */ internal class HttpEventStore( private val httpClient: HttpClient, diff --git a/context-ksqlite/src/commonTest/kotlin/pw/binom/agentik/context/ksqlite/KsqliteContextStoreTest.kt b/context-ksqlite/src/commonTest/kotlin/pw/binom/agentik/context/ksqlite/KsqliteContextStoreTest.kt index d39c749..d94f8de 100644 --- a/context-ksqlite/src/commonTest/kotlin/pw/binom/agentik/context/ksqlite/KsqliteContextStoreTest.kt +++ b/context-ksqlite/src/commonTest/kotlin/pw/binom/agentik/context/ksqlite/KsqliteContextStoreTest.kt @@ -2,7 +2,7 @@ package pw.binom.agentik.context.ksqlite import kotlinx.coroutines.test.runTest import pw.binom.agentik.context.WorkingMemoryEntry -import pw.binom.agentik.messageLog.Content +import pw.binom.agentik.journal.Content import pw.binom.db.ksqlite.SQLiteConnection import kotlin.test.AfterTest import kotlin.test.BeforeTest 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 index 2088288..2c88009 100644 --- 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 @@ -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() + private val mutex = Mutex() + private val buffer = ArrayDeque() private val liveFlow = MutableSharedFlow( 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 = 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 = 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 = buffer.toList() + internal suspend fun snapshot(): List = mutex.withLock { buffer.toList() } override fun close() { + // mutex не закрываем (kotlinx Mutex не AutoCloseable; для in-memory + // store GC соберёт всё при выходе ссылки). buffer чистим. buffer.clear() } 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 index d624349..c746481 100644 --- 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 @@ -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() diff --git a/event-store/src/commonMain/kotlin/pw/binom/agentik/eventStore/EventStore.kt b/event-store/src/commonMain/kotlin/pw/binom/agentik/eventStore/EventStore.kt index 8a54970..0213343 100644 --- a/event-store/src/commonMain/kotlin/pw/binom/agentik/eventStore/EventStore.kt +++ b/event-store/src/commonMain/kotlin/pw/binom/agentik/eventStore/EventStore.kt @@ -1,145 +1,16 @@ package pw.binom.agentik.eventStore -import kotlinx.coroutines.flow.Flow -import kotlinx.coroutines.flow.filter -import kotlinx.coroutines.flow.filterIsInstance -import kotlin.time.Instant -import pw.binom.agentik.proto.CommonEvent - /** - * Bounded-tail event log с автоматическим управлением TTL. + * @Deprecated + * Перенесено в `:outbox-api`. Используй `pw.binom.agentik.outbox.OutboxStore`. * - * **Архитектура двухуровневого хранилища событий**: - * 1. **Этот store** = короткий bounded tail (live SSE + недавний replay). - * События автоматически эвиктятся по TTL/cap (implementation-defined). - * 2. **Message store (`:message-store-api`)** = полный audit log, никогда не - * эвиктится. Source of truth для всего прошлого. + * Backward-compat typealias. Существующие импорты `pw.binom.agentik.eventStore.EventStore` + * продолжают работать; новый код в `:server` использует `:outbox-api` напрямую. * - * **Паттерн reconnect** (caller'ы): - * ``` - * val earliest = store.earliestEventDate() - * if (client.lastSeen < earliest) { - * // gap обнаружен — идём в message store за прошлым - * val gap = messageStore.query(after = client.lastSeen, before = earliest) - * applyAll(gap) - * } - * store.events(after = client.lastSeen).collect { apply(it) } - * ``` - * - * **Нет delete/cleanup методов** — TTL/cap eviction полностью на стороне - * implementation. Это: - * - Убирает single source of truth дублирование (caller не может забыть cleanup). - * - Позволяет impl выбирать retention strategy (TTL, size cap, sliding window). - * - Сохраняет контракт clean: интерфейс только о put/get. - * - * **Read-only**: этот интерфейс предоставляет только read-операции. - * Для записи см. [MutableEventStore]. - * - * **Подписки нереентрантные**: каждый вызов [events] создаёт **новую - * подписку** (cold Flow). Один [events] НЕ видит события, добавленные до - * его вызова, если [after] == null. Если нужен catchup — передавайте - * `after = lastSeenDate` явно. - * - * **Multi-consumer**: разные [events] подписки видят одно и то же live - * tail. Каждая подписка — независимая projection. + * Удалить когда все импорты будут на `:outbox-api`. */ -interface EventStore : AutoCloseable { - - /** - * Subscribe на events. - * - * **`after == null`** → только **live** (события с момента вызова - * `events()`). Каждое новое событие от любого producer'а немедленно - * появится в Flow. Буфер replay не отдаётся. - * - * **`after != null`** → сначала **catchup**: эмитт все буферизованные - * события с `date > after`, порядок `date ASC` (ties по `id ASC`). - * Затем **live** (как null-case). - * - * Cold Flow: каждый вызов — новая подписка. Вызов **после** append'а - * не увидит этот конкретный event (если `after == null`); для catchup - * передавайте явный `after`. - * - * ВАЖНО: `Flow` НЕ бросает ошибку при потере сети между producer и - * store — такие события просто не дойдут до этого Flow. Для гарантии - * полноты клиент обязан cross-check с [earliestEventDate] и fallback - * в message store при gap'е (см. KDoc интерфейса). - */ - fun events(after: Instant?): Flow - - /** - * Subscribe на **только conversation events** (т.е. [CommonEvent.Conversation]). - * - * - [conversationId] == null → события **всех** диалогов. - * - [conversationId] != null → события **только этого** диалога. - * - * Семантика `after` идентична [events] (catchup + live). - * Возвращаемый тип — конкретный subtype [CommonEvent.Conversation]. - */ - /** - * **Default implementation** (читает все events + фильтрует). - * - * Простая реализация через [events] + filterIsInstance. Реализации - * могут override'нуть для эффективности (например, добавить SQL - * `WHERE conversation_id = ?` чтобы не тянуть всё в память), но - * контракт корректен и без override. - */ - fun conversationEvents(after: Instant?, conversationId: String? = null): Flow = - events(after) - .filterIsInstance() - .let { filtered -> - if (conversationId == null) filtered - else filtered.filter { it.conversationId == conversationId } - } - - /** - * Subscribe на **только agent events** ([CommonEvent.Agent] — - * создание/удаление/переименование диалога). - * - * Семантика `after` идентична [events] (catchup + live). - * Возвращаемый тип — конкретный subtype [CommonEvent.Agent]. - * - * Полезно для admin-дашборда, который хочет видеть только lifecycle - * диалогов без деталей ходов. - */ - /** - * **Default implementation** (читает все events + фильтрует по типу). - * - * Простая реализация через [events] + filterIsInstance. Реализации - * могут override'нуть для эффективности (например, читать только agent - * row'ы из БД), но контракт корректен и без override. - */ - fun agentEvents(after: Instant?): Flow = - events(after).filterIsInstance() - - /** - * Date **стартовой точки** буфера. - * - * - Если буфер не пуст → `date` самого старого буферизованного event'а. - * - Если буфер пуст → текущее время (`Clock.System.now()` на момент вызова). - * - * **Семантика "now если пусто"** важна: позволяет клиенту безопасно - * подписаться на [events](after = earliest) сразу — он получит только - * новые live event'ы, без ложного catchup. Если бы возвращалось - * `Instant.DISTANT_PAST` или `null` (с проверкой), клиент мог бы - * ошибочно подписаться на несуществующий catchup и зависнуть в ожидании. - * - * **Используется клиентом для gap detection**: - * - `lastSeen < earliest` → есть дыра в покрытии, нужен fallback - * в message store за диапазоном `[lastSeen, earliest)`. - * - `lastSeen >= earliest` → всё доступно через [events](after), - * fallback не нужен. - * - `lastSeen == earliest` → OK, первый live event будет > earliest. - * - * **Edge case**: клиент, подключившийся до того как store увидел хоть - * один event, получает `earliest ≈ now`. Его `lastSeen` будет < earliest - * — адаптируется в первом же poll'е и пойдёт через fallback если - * сообщения audit log существуют (для consistency с прошлым). - * - * Suspend потому что в persistent impl'ах требует SQL query (`MIN(date)` - * или `Clock.now()` для пустого буфера). - */ - suspend fun earliestEventDate(): Instant - - override fun close() -} +@Deprecated( + message = "Перенесено в :outbox-api. Используй pw.binom.agentik.outbox.OutboxStore.", + replaceWith = ReplaceWith("OutboxStore", "pw.binom.agentik.outbox.OutboxStore"), +) +typealias EventStore = pw.binom.agentik.outbox.OutboxStore diff --git a/event-store/src/commonMain/kotlin/pw/binom/agentik/eventStore/MutableEventStore.kt b/event-store/src/commonMain/kotlin/pw/binom/agentik/eventStore/MutableEventStore.kt index 0856320..b32dd37 100644 --- a/event-store/src/commonMain/kotlin/pw/binom/agentik/eventStore/MutableEventStore.kt +++ b/event-store/src/commonMain/kotlin/pw/binom/agentik/eventStore/MutableEventStore.kt @@ -1,52 +1,15 @@ package pw.binom.agentik.eventStore -import pw.binom.agentik.proto.CommonEvent - /** - * Mutable вариант [EventStore] — добавляет producer-операцию [append]. + * @Deprecated + * Перенесено в `:outbox-api`. Используй `pw.binom.agentik.outbox.MutableOutboxStore`. * - * Этот интерфейс предназначен **только для producer'ов** (ChatAgent, - * sub-agents, A2A-bridge). Consumer'ы (server SSE endpoints, admin - * dashboards, parent agents) должны принимать **read-only** [EventStore] - * — тогда невозможно случайно писать в store из observer'а. + * Backward-compat typealias. `MutableEventStore` == `MutableOutboxStore` (один тип). * - * Типичное использование: - * ``` - * // Producer - * class ChatAgent(private val events: MutableEventStore) { - * suspend fun doSomething() { - * events.append(CommonEvent.Agent(date = now, event = AgentEvent.Created(...))) - * } - * } - * - * // Consumer - * class EventStreamEndpoint(private val events: EventStore) { - * fun stream() = events.events(after = null) - * // Ошибка компиляции если раскомментировать: - * // events.append(...) // ← нельзя, MutableEventStore нет в типе - * } - * ``` - * - * **Append НЕ идемпотентен**: [CommonEvent] не имеет уникального id, - * поэтому retry с тем же logical event (например, после network failure - * между producer и store) приведёт к дубликату в tail'е. Это OK для - * use case'a bounded-tail — клиент, делающий catchup через [events](after), - * получит свой диапазон ровно один раз при подключении, а последующие - * retry producer'а просто насытят tail повторами, не задевая уже - * обработанные. Для гарантированной exactly-once — dedup через - * [message-store] (там есть монотонный `id`). - * - * **Silently evicted**: implementation может выкинуть этот event сразу - * после append (TTL/cap) без уведомления producer'а. Producer **не - * должен** полагаться на то, что event дойдёт до клиента, если он - * вне retention window. + * Удалить когда все импорты будут на `:outbox-api`. */ -interface MutableEventStore : EventStore { - /** - * Положить event в log. - * - * - **Не идемпотентно** — см. KDoc интерфейса. - * - **Suspend** для KMP I/O impl'ов (SQLite через JNI). - */ - suspend fun append(event: CommonEvent) -} +@Deprecated( + message = "Перенесено в :outbox-api. Используй pw.binom.agentik.outbox.MutableOutboxStore.", + replaceWith = ReplaceWith("MutableOutboxStore", "pw.binom.agentik.outbox.MutableOutboxStore"), +) +typealias MutableEventStore = pw.binom.agentik.outbox.MutableOutboxStore diff --git a/message-log-api/build.gradle.kts b/message-log-api/build.gradle.kts index 7319039..5db33ec 100644 --- a/message-log-api/build.gradle.kts +++ b/message-log-api/build.gradle.kts @@ -20,6 +20,9 @@ kotlin { api(libs.kotlinx.coroutines.core) api(libs.kotlinx.serialization.core) api(libs.kotlinx.serialization.json) + // typealias-обёртки указывают на :journal-api — без него + // компиляция падает на Unresolved reference 'journal'. + api(project(":journal-api")) } commonTest.dependencies { implementation(kotlin("test")) diff --git a/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/Content.kt b/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/Content.kt index 7c15fb3..b907735 100644 --- a/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/Content.kt +++ b/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/Content.kt @@ -1,24 +1,13 @@ package pw.binom.agentik.messageLog -import kotlinx.serialization.SerialName -import kotlinx.serialization.Serializable - /** - * Часть контента сообщения на уровне хранилища. Намеренно НЕ зависит от - * `pw.binom.agentik.proto.Content` — маппинг `:proto.Content ↔ Content` живёт - * в `Mapping.kt` storage impl'ов. + * @Deprecated + * Перенесено в `:journal-api`. Используй `pw.binom.agentik.journal.Content`. + * + * Backward-compat typealias. */ -@Serializable -sealed interface Content { - @Serializable - @SerialName("text") - data class Text(val body: String) : Content - - @Serializable - @SerialName("image") - data class Image(val data: ByteArray, val mime: String) : Content { - override fun equals(other: Any?): Boolean = - this === other || (other is Image && mime == other.mime && data.contentEquals(other.data)) - override fun hashCode(): Int = 31 * mime.hashCode() + data.contentHashCode() - } -} +@Deprecated( + message = "Перенесено в :journal-api. Используй pw.binom.agentik.journal.Content.", + replaceWith = ReplaceWith("Content", "pw.binom.agentik.journal.Content"), +) +typealias Content = pw.binom.agentik.journal.Content diff --git a/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MessageContext.kt b/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MessageContext.kt index 01cd75b..2b4e5b5 100644 --- a/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MessageContext.kt +++ b/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MessageContext.kt @@ -1,12 +1,13 @@ package pw.binom.agentik.messageLog -import kotlinx.serialization.Serializable -import kotlinx.serialization.json.JsonElement - -@Serializable -data class MessageContext( - val origin: MessageOrigin, - val sourceId: String? = null, - val description: String? = null, - val metadata: JsonElement? = null, +/** + * @Deprecated + * Перенесено в `:journal-api`. Используй `pw.binom.agentik.journal.MessageContext`. + * + * Backward-compat typealias. + */ +@Deprecated( + message = "Перенесено в :journal-api. Используй pw.binom.agentik.journal.MessageContext.", + replaceWith = ReplaceWith("MessageContext", "pw.binom.agentik.journal.MessageContext"), ) +typealias MessageContext = pw.binom.agentik.journal.MessageContext diff --git a/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MessageOrigin.kt b/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MessageOrigin.kt index 130d8d2..41ced42 100644 --- a/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MessageOrigin.kt +++ b/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MessageOrigin.kt @@ -1,20 +1,13 @@ package pw.binom.agentik.messageLog -import kotlinx.serialization.SerialName -import kotlinx.serialization.Serializable - /** - * Контекст инициации хода (кто/что и почему). Дубликат типа из `:proto` — - * живёт здесь чтобы не тащить `:proto` в слой хранения данных. + * @Deprecated + * Перенесено в `:journal-api`. Используй `pw.binom.agentik.journal.MessageOrigin`. + * + * Backward-compat typealias. */ -@Serializable -enum class MessageOrigin { - @SerialName("user") - USER, - - @SerialName("system") - SYSTEM, - - @SerialName("event") - EVENT, -} \ No newline at end of file +@Deprecated( + message = "Перенесено в :journal-api. Используй pw.binom.agentik.journal.MessageOrigin.", + replaceWith = ReplaceWith("MessageOrigin", "pw.binom.agentik.journal.MessageOrigin"), +) +typealias MessageOrigin = pw.binom.agentik.journal.MessageOrigin diff --git a/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MessageRecord.kt b/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MessageRecord.kt index 08be8ee..89d9ab2 100644 --- a/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MessageRecord.kt +++ b/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MessageRecord.kt @@ -1,71 +1,13 @@ package pw.binom.agentik.messageLog -import kotlinx.serialization.SerialName -import kotlinx.serialization.Serializable -import kotlin.time.Instant - /** - * Запись в таблице `message` (append-only audit). + * @Deprecated + * Перенесено в `:journal-api`. Используй `pw.binom.agentik.journal.MessageRecord`. + * + * Backward-compat typealias. */ -@Serializable -sealed interface MessageRecord { - val id: String - val conversationId: String - val createdAt: Instant - - @Serializable - sealed interface Body : MessageRecord { - val content: List - } - - @Serializable - @SerialName("user") - data class UserMessage( - override val id: String, - override val conversationId: String, - override val content: List, - override val createdAt: Instant, - val context: MessageContext? = null, - ) : Body - - @Serializable - @SerialName("assistant") - data class AssistantMessage( - override val id: String, - override val conversationId: String, - override val content: List, - override val createdAt: Instant, - val tokens: TurnTokens? = null, - ) : Body - - @Serializable - @SerialName("tool_call") - data class ToolCall( - override val id: String, - override val conversationId: String, - val toolName: String, - val toolTitle: String?, - val toolArgsJson: String, - override val createdAt: Instant, - ) : MessageRecord - - @Serializable - @SerialName("tool_result") - data class ToolResult( - override val id: String, - override val conversationId: String, - val toolCallId: String, - val result: String?, - override val createdAt: Instant, - ) : MessageRecord - - @Serializable - @SerialName("error") - data class Error( - override val id: String, - override val conversationId: String, - val message: String, - val code: String?, - override val createdAt: Instant, - ) : MessageRecord -} +@Deprecated( + message = "Перенесено в :journal-api. Используй pw.binom.agentik.journal.MessageRecord.", + replaceWith = ReplaceWith("MessageRecord", "pw.binom.agentik.journal.MessageRecord"), +) +typealias MessageRecord = pw.binom.agentik.journal.MessageRecord diff --git a/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MessageStore.kt b/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MessageStore.kt index 685c6c0..0df7d41 100644 --- a/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MessageStore.kt +++ b/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MessageStore.kt @@ -1,41 +1,17 @@ package pw.binom.agentik.messageLog -import kotlinx.coroutines.flow.Flow -import kotlinx.coroutines.flow.flow -import kotlin.time.Instant - /** - * Append-only audit log сообщений — read-only представление. + * @Deprecated + * Перенесено в `:journal-api`. Используй `pw.binom.agentik.journal.JournalStore`. * - * Producer-операция [MutableMessageStore.append] находится на - * [MutableMessageStore] — этот интерфейс только для чтения, чтобы - * consumer'ы физически не могли писать в audit log. + * Backward-compat typealias. Все существующие consumer'ы + * (`:standalone`, `:storage-ksqlite`, `:storage-inmemory`, тесты) импортируют + * `pw.binom.agentik.messageLog.MessageStore`. Буквально тот же тип. * - * Никаких обновлений, никакого удаления (кроме каскадного вместе - * с ConversationStore.delete). + * Удалить когда все импорты будут на `:journal-api`. */ -interface MessageStore : AutoCloseable { - - suspend fun list(conversationId: String, after: Instant, offset: Int, limit: Int): List - - /** - * Cold-flow paging через [list]. Default-реализация делает N+1 round-trip - * (по странице через `list()` пока не получит короткую страницу). Для - * in-memory backend'ов это OK; remote/SQLite impl'ы могут override'нуть - * на `Channel` / cursor-батчинг, чтобы избежать per-page round-trip. - */ - fun listFlow(conversationId: String, after: Instant, pageSize: Int = PAGE_SIZE): Flow = flow { - var offset = 0 - while (true) { - val page = list(conversationId, after, offset, pageSize) - if (page.isEmpty()) return@flow - for (rec in page) emit(rec) - if (page.size < pageSize) return@flow - offset += page.size - } - } - - companion object { - const val PAGE_SIZE = 100 - } -} +@Deprecated( + message = "Перенесено в :journal-api. Используй pw.binom.agentik.journal.JournalStore.", + replaceWith = ReplaceWith("JournalStore", "pw.binom.agentik.journal.JournalStore"), +) +typealias MessageStore = pw.binom.agentik.journal.JournalStore diff --git a/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MutableMessageStore.kt b/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MutableMessageStore.kt index 443dd6d..4a660a1 100644 --- a/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MutableMessageStore.kt +++ b/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/MutableMessageStore.kt @@ -1,17 +1,15 @@ package pw.binom.agentik.messageLog /** - * Mutable вариант [MessageStore] — добавляет producer-операцию [append]. + * @Deprecated + * Перенесено в `:journal-api`. Используй `pw.binom.agentik.journal.MutableJournalStore`. * - * Этот интерфейс предназначен **только для producer'ов** (ChatAgent, - * ConversationLoop, ToolDispatcher, sub-agents, A2A-bridge). - * Consumer'ы (DebugRoutes, admin dashboards, parent agents) должны - * принимать **read-only** [MessageStore] — тогда невозможно случайно - * записать в audit log из observer'а. + * Backward-compat typealias. `MutableMessageStore` == `MutableJournalStore` (один тип). * - * **Append семантика**: см. KDoc [MessageStore.append][MessageStore] — - * на этом интерфейсе (не дублируем). + * Удалить когда все импорты будут на `:journal-api`. */ -interface MutableMessageStore : MessageStore { - suspend fun append(record: MessageRecord) -} \ No newline at end of file +@Deprecated( + message = "Перенесено в :journal-api. Используй pw.binom.agentik.journal.MutableJournalStore.", + replaceWith = ReplaceWith("MutableJournalStore", "pw.binom.agentik.journal.MutableJournalStore"), +) +typealias MutableMessageStore = pw.binom.agentik.journal.MutableJournalStore diff --git a/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/Payload.kt b/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/Payload.kt index 2c47fe1..1d19e3f 100644 --- a/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/Payload.kt +++ b/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/Payload.kt @@ -1,47 +1,13 @@ package pw.binom.agentik.messageLog -import kotlinx.serialization.SerialName -import kotlinx.serialization.Serializable -import kotlinx.serialization.builtins.ListSerializer -import kotlinx.serialization.json.Json - -private val bodyJson = Json { - ignoreUnknownKeys = true - encodeDefaults = true - explicitNulls = false -} - -@Serializable -data class MessageBodyPayload( - val content: List, - @SerialName("context") - val context: MessageContext? = null, - val tokens: TurnTokens? = null, +/** + * @Deprecated + * Перенесено в `:journal-api`. Используй `pw.binom.agentik.journal.MessageBodyPayload`. + * + * Backward-compat typealias. + */ +@Deprecated( + message = "Перенесено в :journal-api. Используй pw.binom.agentik.journal.MessageBodyPayload.", + replaceWith = ReplaceWith("MessageBodyPayload", "pw.binom.agentik.journal.MessageBodyPayload"), ) - -fun encodeBodyPayload( - content: List, - context: MessageContext? = null, - tokens: TurnTokens? = null, -): String = bodyJson.encodeToString( - MessageBodyPayload.serializer(), - MessageBodyPayload(content = content, context = context, tokens = tokens), -) - -fun decodeBodyPayload(json: String): BodyDecoded = readPayload(json) - -data class BodyDecoded( - val content: List, - val context: MessageContext?, - val tokens: TurnTokens? = null, -) - -private fun readPayload(json: String): BodyDecoded { - return try { - val p = bodyJson.decodeFromString(MessageBodyPayload.serializer(), json) - BodyDecoded(p.content, p.context, p.tokens) - } catch (e: kotlinx.serialization.SerializationException) { - val arr = bodyJson.decodeFromString(ListSerializer(Content.serializer()), json) - BodyDecoded(arr, null, null) - } -} +typealias MessageBodyPayload = pw.binom.agentik.journal.MessageBodyPayload diff --git a/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/TurnTokens.kt b/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/TurnTokens.kt index c643d05..4168e14 100644 --- a/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/TurnTokens.kt +++ b/message-log-api/src/commonMain/kotlin/pw/binom/agentik/messageLog/TurnTokens.kt @@ -1,18 +1,13 @@ package pw.binom.agentik.messageLog -import kotlinx.serialization.Serializable - /** - * Token usage одного assistant turn'а. + * @Deprecated + * Перенесено в `:journal-api`. Используй `pw.binom.agentik.journal.TurnTokens`. + * + * Backward-compat typealias. */ -@Serializable -data class TurnTokens( - val input: Int, - val output: Int, -) { - val total: Int get() = input + output - init { - require(input >= 0) { "input tokens must be non-negative, got $input" } - require(output >= 0) { "output tokens must be non-negative, got $output" } - } -} \ No newline at end of file +@Deprecated( + message = "Перенесено в :journal-api. Используй pw.binom.agentik.journal.TurnTokens.", + replaceWith = ReplaceWith("TurnTokens", "pw.binom.agentik.journal.TurnTokens"), +) +typealias TurnTokens = pw.binom.agentik.journal.TurnTokens diff --git a/message-log-ksqlite/build.gradle.kts b/message-log-ksqlite/build.gradle.kts new file mode 100644 index 0000000..230e20b --- /dev/null +++ b/message-log-ksqlite/build.gradle.kts @@ -0,0 +1,33 @@ +plugins { + alias(libs.plugins.kotlin.multiplatform) + alias(libs.plugins.kotlin.serialization) +} + +// KMP-реализация MutableMessageStore поверх ksqlite (https://github.com/caffeine-mgn/ksqlite). +// +// Цели сборки (параллельны :storage-ksqlite): +// - jvm() — основная, тесты гоняются здесь (in-memory DB без файла) +// - linuxX64() / mingwX64() — smoke-проверка что KMP реально KMP +// +// Apple targets невозможно собрать на Linux — Kotlin Multiplatform plugin auto-disables их. + +kotlin { + jvmToolchain(21) + + jvm() + linuxX64() + mingwX64() + + sourceSets { + commonMain.dependencies { + api(project(":message-log-api")) + + implementation("pw.binom.db:ksqlite:0.1.0") + implementation(libs.kotlinx.serialization.json) + } + commonTest.dependencies { + implementation(kotlin("test")) + implementation(libs.kotlinx.coroutines.test) + } + } +} \ No newline at end of file diff --git a/message-log-ksqlite/src/commonMain/kotlin/pw/binom/agentik/messageLog/ksqlite/KsqliteMessageStore.kt b/message-log-ksqlite/src/commonMain/kotlin/pw/binom/agentik/messageLog/ksqlite/KsqliteMessageStore.kt new file mode 100644 index 0000000..2aa1c3d --- /dev/null +++ b/message-log-ksqlite/src/commonMain/kotlin/pw/binom/agentik/messageLog/ksqlite/KsqliteMessageStore.kt @@ -0,0 +1,141 @@ +package pw.binom.agentik.messageLog.ksqlite + +import kotlinx.serialization.json.Json +import pw.binom.agentik.messageLog.MessageRecord +import pw.binom.agentik.messageLog.MutableMessageStore +import pw.binom.db.ksqlite.SQLiteConnection +import pw.binom.db.ksqlite.SQLitePreparedStatement +import kotlin.time.Instant +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.sync.Mutex +import kotlinx.coroutines.sync.withLock +import kotlinx.coroutines.withContext + +/** + * ksqlite-реализация [MutableMessageStore]. Схема таблицы `message` повторяет + * [pw.binom.agentik.storage.sqlite.SqliteMessageStore] для совместимости данных. + * + * Все четыре [SQLitePreparedStatement] компилируются один раз в конструкторе и + * переиспользуются: bind → execute → reset (cursor) + clearBindings (memory) + * перед следующим вызовом. Это устраняет per-call overhead на sqlite3_prepare_v2. + * + * encoding helpers (`encodeRecord`/`toMessageRecord`/`CallPayload`/...) — копия + * `:storage-sqlite`, живут в этом же модуле. См. [MessageCodecs]. + * + * **Threading**: ksqlite документирует соединение как single-threaded by convention + * (sqlite3 prepared statements safe to use serially from any thread). Мы сериализуем + * доступ через [Mutex] — операции вызываются на [Dispatchers.Default] и поток + * между вызовами может меняться; Mutex гарантирует что в любой момент времени + * только одна корутина использует statement. + * + * `clear(conversationId)` — каскадный helper, вызывается из + * [pw.binom.agentik.storage.ksqlite.KsqliteConversationStore.delete]. + * Не часть публичного [MutableMessageStore] API (audit log append-only); + * экспортируется через lambda в [pw.binom.agentik.storage.ksqlite.KsqliteStores.assemble]. + */ +class KsqliteMessageStore private constructor( + private val connection: SQLiteConnection, + private val ownsConnection: Boolean, +) : MutableMessageStore { + + /** + * Основной конструктор — принимает **уже открытое** [SQLiteConnection]. + * Соединение остаётся под управлением вызывающего (например, [pw.binom.agentik.storage.ksqlite.KsqliteStores.assemble] + * закрывает его сам после [close] всех своих store'ов). + */ + constructor(connection: SQLiteConnection) : this(connection, ownsConnection = false) + + /** + * Convenience-конструктор для standalone-сценариев: открывает файл-БД + * (`SQLiteConnection.open(path)`) и **сам закрывает её в [close]**. + * Не использовать если соединение шарится с другими store'ами — + * двойной [SQLiteConnection.close] не идемпотентен. + */ + constructor(path: String) : this(SQLiteConnection.open(path), ownsConnection = true) + + private val mutex = Mutex() + private val json = Json { ignoreUnknownKeys = true } + + private val insertStmt: SQLitePreparedStatement = connection.prepare( + "INSERT INTO message (id, conversation_id, kind, payload_json, created_at) VALUES (?, ?, ?, ?, ?)" + ) + private val listStmt: SQLitePreparedStatement = connection.prepare( + "SELECT id, conversation_id, kind, payload_json, created_at FROM message " + + "WHERE conversation_id = ? AND created_at > ? ORDER BY created_at ASC, id ASC LIMIT ? OFFSET ?" + ) + private val deleteStmt: SQLitePreparedStatement = connection.prepare( + "DELETE FROM message WHERE conversation_id = ?" + ) + + override suspend fun append(record: MessageRecord): Unit = withContext(Dispatchers.Default) { + mutex.withLock { + val (kind, payload) = encodeRecord(record) + insertStmt.execDml( + record.id, record.conversationId, kind, payload, record.createdAt.toEpochMilliseconds(), + ) + } + } + + override suspend fun list( + conversationId: String, + after: Instant, + offset: Int, + limit: Int, + ): List = withContext(Dispatchers.Default) { + mutex.withLock { + listStmt.runQuery(conversationId, after.toEpochMilliseconds(), limit.toLong(), offset.toLong()) + } + } + + override fun close() { + insertStmt.close() + listStmt.close() + deleteStmt.close() + if (ownsConnection) connection.close() + } + + /** + * Каскадный clear всех сообщений диалога — вызывается из + * [pw.binom.agentik.storage.ksqlite.KsqliteConversationStore.delete]. + * + * Публичный (не internal) чтобы `:storage-ksqlite` мог передать его как lambda. + * Не часть [MutableMessageStore] — audit log append-only. + */ + suspend fun clear(conversationId: String): Unit = withContext(Dispatchers.Default) { + mutex.withLock { + deleteStmt.execDml(conversationId) + } + } + + private fun SQLitePreparedStatement.bindAll(params: Array) { + params.forEachIndexed { i, v -> + val idx = i + 1 + when (v) { + is String -> bindText(idx, v) + is Long -> bindLong(idx, v) + is Int -> bindLong(idx, v.toLong()) + else -> error("Unsupported bind type: ${v::class}") + } + } + } + + /** Выполнить INSERT/DELETE statement. reset+clearBindings+bind+executeUpdate. */ + private fun SQLitePreparedStatement.execDml(vararg params: Any): Int { + reset() + clearBindings() + bindAll(params) + return executeUpdate() + } + + /** Выполнить SELECT statement, вернуть список [MessageRecord]. reset+clearBindings+bind+executeQuery+drain. */ + private fun SQLitePreparedStatement.runQuery(vararg params: Any): List { + reset() + clearBindings() + bindAll(params) + val out = mutableListOf() + executeQuery().use { rs -> + while (rs.next()) out.add(rs.toMessageRecord(json)) + } + return out + } +} \ No newline at end of file diff --git a/message-log-ksqlite/src/commonMain/kotlin/pw/binom/agentik/messageLog/ksqlite/MessageCodecs.kt b/message-log-ksqlite/src/commonMain/kotlin/pw/binom/agentik/messageLog/ksqlite/MessageCodecs.kt new file mode 100644 index 0000000..ace8a2a --- /dev/null +++ b/message-log-ksqlite/src/commonMain/kotlin/pw/binom/agentik/messageLog/ksqlite/MessageCodecs.kt @@ -0,0 +1,79 @@ +package pw.binom.agentik.messageLog.ksqlite + +import kotlinx.serialization.json.Json +import pw.binom.agentik.messageLog.MessageRecord +import pw.binom.agentik.messageLog.decodeBodyPayload +import pw.binom.agentik.messageLog.encodeBodyPayload +import pw.binom.db.ksqlite.SQLiteResultSet +import kotlin.time.Instant + +/** + * Кодирование [MessageRecord] → пара (kind, payloadJson) для SQLite. + * + * Копия [pw.binom.agentik.storage.sqlite.SqliteMessageStore] — `private` helpers + * нельзя переиспользовать между модулями, поэтому в каждом backend свой набор. + * Чтобы избежать дрейфа при изменении формата payload'а, оба набора синхронизируются + * через эти data class'ы (CallPayload/ResultPayload/ErrorPayload). + */ +internal fun encodeRecord(record: MessageRecord): Pair = when (record) { + is MessageRecord.UserMessage -> "user" to encodeBodyPayload( + content = record.content, + context = record.context, + ) + is MessageRecord.AssistantMessage -> "assistant" to encodeBodyPayload( + content = record.content, + tokens = record.tokens, + ) + is MessageRecord.ToolCall -> "tool_call" to Json.encodeToString( + CallPayload.serializer(), + CallPayload(name = record.toolName, title = record.toolTitle, argsJson = record.toolArgsJson), + ) + is MessageRecord.ToolResult -> "tool_result" to Json.encodeToString( + ResultPayload.serializer(), + ResultPayload(toolCallId = record.toolCallId, result = record.result), + ) + is MessageRecord.Error -> "error" to Json.encodeToString( + ErrorPayload.serializer(), + ErrorPayload(message = record.message, code = record.code), + ) +} + +internal fun SQLiteResultSet.toMessageRecord(json: Json): MessageRecord { + val id = getText(0)!! + val convId = getText(1)!! + val kind = getText(2)!! + val payload = getText(3)!! + val createdAt = Instant.fromEpochMilliseconds(getLong(4)!!) + return when (kind) { + "user" -> { + val d = decodeBodyPayload(payload) + MessageRecord.UserMessage(id = id, conversationId = convId, content = d.content, createdAt = createdAt, context = d.context) + } + "assistant" -> { + val d = decodeBodyPayload(payload) + MessageRecord.AssistantMessage(id = id, conversationId = convId, content = d.content, createdAt = createdAt, tokens = d.tokens) + } + "tool_call" -> { + val p = Json.decodeFromString(CallPayload.serializer(), payload) + MessageRecord.ToolCall(id = id, conversationId = convId, toolName = p.name, toolTitle = p.title, toolArgsJson = p.argsJson, createdAt = createdAt) + } + "tool_result" -> { + val p = Json.decodeFromString(ResultPayload.serializer(), payload) + MessageRecord.ToolResult(id = id, conversationId = convId, toolCallId = p.toolCallId, result = p.result, createdAt = createdAt) + } + "error" -> { + val p = Json.decodeFromString(ErrorPayload.serializer(), payload) + MessageRecord.Error(id = id, conversationId = convId, message = p.message, code = p.code, createdAt = createdAt) + } + else -> error("Unknown message kind in audit log: $kind") + } +} + +@kotlinx.serialization.Serializable +internal data class CallPayload(val name: String, val title: String?, val argsJson: String) + +@kotlinx.serialization.Serializable +internal data class ResultPayload(val toolCallId: String, val result: String?) + +@kotlinx.serialization.Serializable +internal data class ErrorPayload(val message: String, val code: String?) \ No newline at end of file diff --git a/message-log-ksqlite/src/commonTest/kotlin/pw/binom/agentik/messageLog/ksqlite/KsqliteMessageStoreTest.kt b/message-log-ksqlite/src/commonTest/kotlin/pw/binom/agentik/messageLog/ksqlite/KsqliteMessageStoreTest.kt new file mode 100644 index 0000000..57e4e3c --- /dev/null +++ b/message-log-ksqlite/src/commonTest/kotlin/pw/binom/agentik/messageLog/ksqlite/KsqliteMessageStoreTest.kt @@ -0,0 +1,122 @@ +package pw.binom.agentik.messageLog.ksqlite + +import kotlinx.coroutines.flow.count +import kotlinx.coroutines.flow.first +import kotlinx.coroutines.flow.toList +import kotlinx.coroutines.test.runTest +import pw.binom.agentik.messageLog.Content +import pw.binom.agentik.messageLog.MessageRecord +import pw.binom.db.ksqlite.SQLiteConnection +import kotlin.test.AfterTest +import kotlin.test.BeforeTest +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertNull +import kotlin.time.Instant + +class KsqliteMessageStoreTest { + + private lateinit var conn: SQLiteConnection + private lateinit var store: KsqliteMessageStore + + @BeforeTest + fun setup() { + conn = SQLiteConnection.memory("msg-${kotlin.random.Random.nextLong()}") + conn.exec(SCHEMA) + store = KsqliteMessageStore(conn) + } + + @AfterTest + fun tearDown() = conn.close() + + @Test + fun testAppendUserAndRetrieve() = runTest { + store.append(MessageRecord.UserMessage( + id = "m1", + conversationId = "conv1", + content = listOf(Content.Text("hello")), + createdAt = Instant.parse("2026-09-15T10:01:00Z"), + context = null, + )) + val list = store.listFlow("conv1", Instant.DISTANT_PAST).toList() + assertEquals(1, list.size) + val msg = list[0] + assertEquals("m1", msg.id) + assertEquals(MessageRecord.UserMessage::class, msg::class) + } + + @Test + fun testAppendAssistantWithTokens() = runTest { + store.append(MessageRecord.AssistantMessage( + id = "m1", + conversationId = "conv1", + content = listOf(Content.Text("hi")), + createdAt = Instant.parse("2026-09-15T10:01:00Z"), + tokens = pw.binom.agentik.messageLog.TurnTokens(input = 50, output = 30), + )) + val all = store.listFlow("conv1", Instant.DISTANT_PAST).toList() + assertEquals(1, all.size) + val msg = all[0] as MessageRecord.AssistantMessage + assertEquals(50, msg.tokens?.input) + assertEquals(30, msg.tokens?.output) + } + + @Test + fun testListAfterFiltersByTimestamp() = runTest { + val t1 = Instant.parse("2026-09-15T10:01:00Z") + val t2 = Instant.parse("2026-09-15T10:02:00Z") + val t3 = Instant.parse("2026-09-15T10:03:00Z") + store.append(MessageRecord.UserMessage("m1", "conv1", listOf(Content.Text("a")), t1, null)) + store.append(MessageRecord.UserMessage("m2", "conv1", listOf(Content.Text("b")), t2, null)) + store.append(MessageRecord.UserMessage("m3", "conv1", listOf(Content.Text("c")), t3, null)) + + val after = store.list("conv1", after = t1, offset = 0, limit = 10) + assertEquals(2, after.size) + assertEquals(listOf("m2", "m3"), after.map { it.id }) + } + + @Test + fun testListAllReturnsAllInOrder() = runTest { + val t = Instant.parse("2026-09-15T10:00:00Z") + for (i in 1..3) store.append( + MessageRecord.UserMessage("m$i", "conv1", listOf(Content.Text("x$i")), t + kotlin.time.Duration.parse("PT${i}S"), null) + ) + assertEquals(listOf("m1", "m2", "m3"), store.listFlow("conv1", Instant.DISTANT_PAST).toList().map { it.id }) + } + + @Test + fun testClearRemovesByConversation() = runTest { + val t = Instant.parse("2026-09-15T10:00:00Z") + store.append(MessageRecord.UserMessage("m1", "conv1", listOf(Content.Text("a")), t, null)) + store.append(MessageRecord.UserMessage("m2", "conv2", listOf(Content.Text("b")), t, null)) + store.clear("conv1") + assertEquals(0, store.listFlow("conv1", Instant.DISTANT_PAST).count()) + assertEquals(1, store.listFlow("conv2", Instant.DISTANT_PAST).count()) + } + + @Test + fun testListWithNullContext() = runTest { + store.append(MessageRecord.UserMessage( + id = "m1", + conversationId = "conv1", + content = listOf(Content.Text("hi")), + createdAt = Instant.parse("2026-09-15T10:01:00Z"), + context = null, + )) + val msg = store.listFlow("conv1", Instant.DISTANT_PAST).first() as MessageRecord.UserMessage + assertNull(msg.context) + } + + private companion object { + const val SCHEMA = """ + CREATE TABLE IF NOT EXISTS message ( + id TEXT NOT NULL PRIMARY KEY, + conversation_id TEXT NOT NULL, + kind TEXT NOT NULL, + payload_json TEXT NOT NULL, + created_at INTEGER NOT NULL + ); + CREATE INDEX IF NOT EXISTS idx_msg_conv ON message(conversation_id, created_at); + """ + } +} \ No newline at end of file diff --git a/outbox-api/build.gradle.kts b/outbox-api/build.gradle.kts index 5181866..d44cff0 100644 --- a/outbox-api/build.gradle.kts +++ b/outbox-api/build.gradle.kts @@ -1,5 +1,8 @@ plugins { alias(libs.plugins.kotlin.multiplatform) + // Нужен для @Serializable на AgentEvent/CommonEvent/Event — все + // три типа теперь живут в :outbox-api (см. миграцию из :proto). + alias(libs.plugins.kotlin.serialization) } kotlin { @@ -17,8 +20,11 @@ kotlin { sourceSets { commonMain.dependencies { - api(project(":proto")) + // :proto больше не нужен — AgentEvent/CommonEvent/Event перенесены + // сюда, и они self-contained (Event ссылается только на kotlinx-serialization). api(libs.kotlinx.coroutines.core) + api(libs.kotlinx.serialization.core) + api(libs.kotlinx.serialization.json) } commonTest.dependencies { implementation(kotlin("test")) diff --git a/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/MutableOutboxStore.kt b/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/MutableOutboxStore.kt index b0d5a08..ed59ac4 100644 --- a/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/MutableOutboxStore.kt +++ b/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/MutableOutboxStore.kt @@ -1,7 +1,5 @@ package pw.binom.agentik.outbox -import pw.binom.agentik.proto.CommonEvent - /** * Mutable вариант [OutboxStore] — добавляет producer-операцию [append]. * diff --git a/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/OutboxStore.kt b/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/OutboxStore.kt index 7d1de19..b68174e 100644 --- a/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/OutboxStore.kt +++ b/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/OutboxStore.kt @@ -4,7 +4,6 @@ import kotlinx.coroutines.flow.Flow import kotlinx.coroutines.flow.filter import kotlinx.coroutines.flow.filterIsInstance import kotlin.time.Instant -import pw.binom.agentik.proto.CommonEvent /** * Bounded-tail event log с автоматическим управлением TTL. diff --git a/proto/build.gradle.kts b/proto/build.gradle.kts index 00bcaba..f97247d 100644 --- a/proto/build.gradle.kts +++ b/proto/build.gradle.kts @@ -23,6 +23,18 @@ kotlin { api(libs.kotlinx.serialization.core) // для JsonElement в MessageContext.metadata api(libs.kotlinx.serialization.json) + // Read-only storage handles, выставляемые через Agent.journal + // и Agent.outbox. Типы JournalStore/OutboxStore фигурируют в + // public-сигнатуре Agent, поэтому api-висимости. + // + // Линейный граф зависимостей (без циклов): + // :proto ──► :outbox-api (нет обратной зависимости) + // :proto ──► :journal-api (нет обратной зависимости) + // Добились переносом AgentEvent/CommonEvent/Event из :proto в + // :outbox-api — они теперь self-contained в outbox (не нужны + // :proto-типы), а :proto использует их через :outbox-api. + api(project(":journal-api")) + api(project(":outbox-api")) } commonTest.dependencies { implementation(kotlin("test")) diff --git a/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Agent.kt b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Agent.kt index c6edc11..9e0c2ab 100644 --- a/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Agent.kt +++ b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Agent.kt @@ -2,6 +2,8 @@ package pw.binom.agentik.proto import kotlinx.coroutines.flow.Flow import kotlinx.coroutines.flow.flow +import pw.binom.agentik.journal.JournalStore +import pw.binom.agentik.outbox.OutboxStore import kotlin.time.Instant /** @@ -10,12 +12,46 @@ import kotlin.time.Instant * Транспортно-агностично. [Agent] — фабрика stateful-диалогов: * [createConversation] возвращает [Conversation], который сам хранит историю * и которому отправляют ходы через [Conversation.send]. + * + * **Хранилища вынесены в [Agent.journal] и [Agent.outbox]**: оба read-only. + * События больше НЕ часть [Agent] (раньше были `events()`/`allEvents()`) — + * они теперь живут в [outbox] как `OutboxStore.events(after)` / + * `outbox.agentEvents(after)`. Это даёт единый путь для всех read-операций + * по хранилищу и убирает дублирование между протоколом и хранилищем. */ public interface Agent { /** Идентификатор агента. */ val id: String + /** + * Append-only audit log всех сообщений диалогов (read-only view). + * + * Используется HTTP-фасадом `:server` для endpoint'а + * `GET /{path}/journal/conversations/{id}/messages` — внешние клиенты + * (дашборды, parent-агенты, A2A-bridge) могут читать полный transcript + * диалога, включая tool-call/tool-result/error, без необходимости идти + * через `Conversation.getMessages` (который возвращает уже + * project'нутый proto-Message). + * + * **Read-only**: write-доступ только через `MutableJournalStore` + * внутри ChatAgent / ConversationLoop, не через [Agent] interface. + */ + val journal: JournalStore + + /** + * Bounded-tail live event stream агента (read-only view). + * + * Используется HTTP-фасадом `:server` для endpoint'а + * `GET /{path}/outbox/events?after=` (SSE) — внешние клиенты подписываются + * на agent lifecycle + conversation events. Catchup+live контракт — см. + * KDoc `OutboxStore.events`. + * + * **Read-only**: write-доступ только через `MutableOutboxStore` внутри + * ChatAgent / ConversationLoop, не через [Agent] interface. + */ + val outbox: OutboxStore + /** Создаёт новый stateful-диалог с агентом. */ fun createConversation(temp: Boolean): Conversation @@ -39,28 +75,6 @@ public interface Agent { } } - /** - * Live-подписка на изменения в множестве диалогов агента: создание, - * удаление, переименование (см. [AgentEvent]). События внутри конкретного - * диалога приходят через [Conversation.events]. - * - * **Не реплеит** прошлое — для снимка множества используй [getConversations] - * или [getConversation]. - */ - fun events(after: Instant): Flow - - /** - * All events in one stream: agent lifecycle (Created/Deleted/Renamed) + - * all conversation turns. Useful for admin dashboards, debug tools, - * parent agents. - * - * For UI use [events] + [Conversation.events]. This one-feed variant is - * for cases where everything-in-one is preferred. - * - * Cold (no replay). For catchup use EventStore. - */ - fun allEvents(after: Instant): Flow - companion object { const val PAGE_SIZE: Int = 100 diff --git a/proto/src/commonMain/kotlin/pw/binom/agentik/proto/AgentEvent.kt b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/AgentEvent.kt index 0f07d0f..54c0fde 100644 --- a/proto/src/commonMain/kotlin/pw/binom/agentik/proto/AgentEvent.kt +++ b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/AgentEvent.kt @@ -1,43 +1,10 @@ package pw.binom.agentik.proto -import kotlinx.serialization.SerialName -import kotlinx.serialization.Serializable -import kotlin.time.Instant - /** - * Live-события уровня [Agent]: изменения в множестве диалогов - * (создание, удаление, переименование). События, происходящие **внутри** - * конкретного диалога, приходят через [Conversation.events], а не сюда. + * Backward-compat typealias: `AgentEvent` теперь живёт в `:outbox-api` + * (логически принадлежит сущности outbox, не wire-протоколу `:proto`). * - * Каждое событие несёт [date] — момент эмиссии в UTC. Семантика подписки - * идентична [Conversation.events]: поток **не реплеит** прошлое, для бэкфилла - * используются `getConversations`/`getConversation`. + * Существующие импорты `pw.binom.agentik.proto.AgentEvent` продолжают + * работать транспарентно. Использовать typealias в новом коде. */ -@Serializable -sealed interface AgentEvent { - /** Момент эмиссии события в UTC. */ - val date: Instant - - /** - * Создан новый диалог. Передаётся его id — handle можно получить через - * [Agent.getConversation]. Подписчик после [Created] может сразу открыть - * live-подписку на этот диалог через [Conversation.events]. - */ - @Serializable - @SerialName("created") - data class Created(override val date: Instant, val conversationId: String) : AgentEvent - - /** - * Диалог удалён. Переданный [Conversation]-handle реализация обязана - * закрыть (`close()`) до эмиссии этого события — после [Deleted] - * пользоваться handle нельзя. - */ - @Serializable - @SerialName("deleted") - data class Deleted(override val date: Instant, val id: String) : AgentEvent - - /** У диалога сменился заголовок. */ - @Serializable - @SerialName("renamed") - data class Renamed(override val date: Instant, val id: String, val title: String?) : AgentEvent -} +typealias AgentEvent = pw.binom.agentik.outbox.AgentEvent diff --git a/proto/src/commonMain/kotlin/pw/binom/agentik/proto/CommonEvent.kt b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/CommonEvent.kt index fc95e8c..c201748 100644 --- a/proto/src/commonMain/kotlin/pw/binom/agentik/proto/CommonEvent.kt +++ b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/CommonEvent.kt @@ -1,39 +1,9 @@ package pw.binom.agentik.proto -import kotlin.time.Instant -import kotlinx.serialization.SerialName -import kotlinx.serialization.Serializable - /** - * Unified wrapper for all agent events in a single stream. + * Backward-compat typealias: `CommonEvent` теперь живёт в `:outbox-api`. * - * Useful for admin dashboards, debug tools, parent agents: one subscription - * instead of N+1. For regular UI use two separate SSE feeds - * ([AgentEvent] via /events and [Event] via /conversations/{id}/events); - * [CommonEvent] is for those who need everything in one place. - * - * Server endpoint: GET /events/all (SSE), or replay via EventStore. - * - * Not used for persistence payload: EventStore stores AgentEvent and - * Conversation.Event natively (compact form); this wrapper is wire-format - * only. + * Существующие импорты `pw.binom.agentik.proto.CommonEvent` продолжают + * работать транспарентно. */ -@Serializable -sealed interface CommonEvent { - val date: Instant - - @Serializable - @SerialName("agent") - data class Agent( - override val date: Instant, - val event: AgentEvent, - ) : CommonEvent - - @Serializable - @SerialName("conversation") - data class Conversation( - override val date: Instant, - val conversationId: String, - val event: Event, - ) : CommonEvent -} +typealias CommonEvent = pw.binom.agentik.outbox.CommonEvent diff --git a/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Event.kt b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Event.kt index af3d131..1558e97 100644 --- a/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Event.kt +++ b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Event.kt @@ -1,88 +1,11 @@ package pw.binom.agentik.proto -import kotlinx.serialization.SerialName -import kotlinx.serialization.Serializable -import kotlin.time.Instant - /** - * Элемент live-потока [Conversation.events]. + * Backward-compat typealias: `Event` теперь живёт в `:outbox-api`. * - * Каждое событие несёт [date] — момент эмиссии в UTC. Используется клиентом - * для трекинга «где остановился» при обрыве/переподключении и для разрешения - * порядка при равных timestamps. - * - * Базовая структура хода: - * `StartReasoning?` → `StartResponse(TEXT|IMAGE)` → ...контент... → `End` | `Interrupted` | `Error`. - * `StartReasoning` может отсутствовать, если агент не показывал рассуждения. + * Существующие импорты `pw.binom.agentik.proto.Event` продолжают + * работать транспарентно. `Conversation.events(after): Flow` в + * `:proto.Conversation` теперь фактически возвращает + * `pw.binom.agentik.outbox.Event` — тот же тип, другое имя. */ -@Serializable -sealed interface Event { - /** Момент эмиссии события в UTC. */ - val date: Instant - - @Serializable - enum class ResponseType { - @SerialName("text") TEXT, - @SerialName("image") IMAGE - } - - /** Ассистент начал рассуждение (опциональный маркер; контент рассуждения приходит через [AppendText]). */ - @Serializable - @SerialName("start_reasoning") - data class StartReasoning(override val date: Instant) : Event - - /** Начало ответа ассистента заданного типа. После него идут соответствующие `Append*`/`Tool*`-события, потом [End]/[Interrupted]/[Error]. */ - @Serializable - @SerialName("start_response") - data class StartResponse(override val date: Instant, val responseType: ResponseType) : Event - - /** Ход завершён нормально. Соответствующий [Message.AssistantMessage] появится в `getMessages`. */ - @Serializable - @SerialName("end") - data class End(override val date: Instant) : Event - - /** Ход прерван через [Conversation.interrupt]. Частичный ответ НЕ сохраняется в истории. */ - @Serializable - @SerialName("interrupted") - data class Interrupted(override val date: Instant) : Event - - @Serializable - @SerialName("append_text") - data class AppendText(override val date: Instant, val body: String) : Event - - @Serializable - @SerialName("append_image") - data class AppendImage(override val date: Instant, val body: ByteArray, val mime: String) : Event - - /** - * Агент начал вызов тула. Аргументы приходят целиком — стриминга нет. - * [id] совпадает с id соответствующего [Message.ToolCall] в истории - * после завершения хода. - */ - @Serializable - @SerialName("tool_call") - data class ToolCall( - override val date: Instant, - val id: String, - val title: String?, - val toolName: String, - val toolArgs: String, - ) : Event - - /** - * Результат вызова тула. Приходит целиком после завершения исполнения. - * [id] совпадает с [ToolCall.id], к которому относится результат, и - * с id [Message.ToolResult] в истории. - */ - @Serializable - @SerialName("tool_result") - data class ToolResult(override val date: Instant, val id: String, val result: String?) : Event - - /** - * Ошибка хода. После неё поток завершается; дальнейшие события могут - * прийти, но ход считается проваленным. - */ - @Serializable - @SerialName("error") - data class Error(override val date: Instant, val message: String, val code: String? = null) : Event -} +typealias Event = pw.binom.agentik.outbox.Event diff --git a/server/build.gradle.kts b/server/build.gradle.kts index ebc506e..f4a9a00 100644 --- a/server/build.gradle.kts +++ b/server/build.gradle.kts @@ -24,6 +24,16 @@ kotlin { commonMain.dependencies { implementation(project(":proto")) + // Read-only storage handles, которые HTTP-фасад выставляет наружу + // под {path}/journal/* и {path}/outbox/*. Типы JournalStore/ + // OutboxStore фигурируют в сигнатурах internal-функций + // journalRoutes/outboxRoutes, поэтому нужны в compile classpath. + // Транзитивные api-висимости :proto (:journal-api, :outbox-api) + // не доходят до :server из-за implementation(:proto), поэтому + // объявляем напрямую. + implementation(project(":journal-api")) + implementation(project(":outbox-api")) + // Ktor (без engine — engine подключает потребитель, см. :standalone). implementation(libs.ktor.server.core) implementation(libs.ktor.server.content.negotiation) diff --git a/server/src/commonMain/kotlin/pw/binom/agentik/server/JournalRoutes.kt b/server/src/commonMain/kotlin/pw/binom/agentik/server/JournalRoutes.kt new file mode 100644 index 0000000..8095830 --- /dev/null +++ b/server/src/commonMain/kotlin/pw/binom/agentik/server/JournalRoutes.kt @@ -0,0 +1,44 @@ +package pw.binom.agentik.server + +import io.ktor.http.HttpStatusCode +import io.ktor.server.response.respond +import io.ktor.server.routing.Route +import io.ktor.server.routing.get +import io.ktor.server.routing.route +import pw.binom.agentik.journal.JournalStore + +/** + * HTTP-фасад для [JournalStore] (append-only audit log сообщений диалога). + * + * **Основа:** функция-продолжение для [Route.agentikAgent] — внутри неё уже + * создан роут под `{path}` агента. [journalRoutes] добавляет под-prefix + * [path] (по умолчанию `"/journal"`) к **этому же** родительскому роуту, + * итоговый URL = `{path агента}/journal/...`. + * + * **Endpoint'ы под `{path}/journal`:** + * - `GET /conversations/{id}/messages?after=&offset=&limit=` — список raw + * [pw.binom.agentik.journal.MessageRecord] (все типы: UserMessage / + * AssistantMessage / ToolCall / ToolResult / Error). В отличие от + * `GET /conversations/{id}/messages` в [agentikRoutes] (который отдаёт + * project'нутые proto-[pw.binom.agentik.proto.Message]), здесь клиент + * получает полный transcript с tool-call/tool-result/error payload-ами, + * turn-tokens и context-метаданными. + * + * **Read-only:** [JournalStore] не имеет `append` — запись только через + * writer-референс, который ChatAgent держит внутри (тип `MutableJournalStore`, + * не выставлен наружу через [pw.binom.agentik.proto.Agent]). + */ +fun Route.journalRoutes( + journal: JournalStore, + path: String = "/journal", +) { + route(path) { + get("/conversations/{id}/messages") { + val id = call.parameters["id"]!! + val after = call.parseAfter() ?: return@get + val offset = call.request.queryParameters["offset"]?.toIntOrNull() ?: 0 + val limit = call.request.queryParameters["limit"]?.toIntOrNull() ?: JournalStore.PAGE_SIZE + call.respond(journal.list(id, after, offset, limit)) + } + } +} diff --git a/server/src/commonMain/kotlin/pw/binom/agentik/server/Module.kt b/server/src/commonMain/kotlin/pw/binom/agentik/server/Module.kt index 29d4665..7218db3 100644 --- a/server/src/commonMain/kotlin/pw/binom/agentik/server/Module.kt +++ b/server/src/commonMain/kotlin/pw/binom/agentik/server/Module.kt @@ -20,25 +20,30 @@ import pw.binom.agentik.proto.Agent * }.start(wait = true) * ``` * - * Под префиксом [path] монтируются: - * - `POST /conversations` — создать диалог - * - `GET /conversations` — список - * - `GET /conversations/{id}` — один диалог - * - `PATCH /conversations/{id}` — переименовать - * - `DELETE /conversations/{id}` — удалить - * - `POST /conversations/{id}/messages` — `send` (202 Accepted) - * - `POST /conversations/{id}/interrupt` — `interrupt` - * - `GET /conversations/{id}/messages` — история - * - `GET /conversations/{id}/events` — SSE: события хода (catchup + live) - * - `GET /events` — SSE: события агента (catchup + live) - * - `GET /events/all` — SSE: всё в одном потоке - * - `GET /health` — `"ok"` + * **Все** дочерние фасады ([agentikRoutes] / [journalRoutes] / [outboxRoutes]) + * монтируются внутри `{path}` — внутренний `route(path)` создаёт родительский + * роут агента, и storage-фасады добавляют свои под-prefix'ы **к этому же** + * роуту, а не к корню. Итоговая раскладка: * - * Event-эндпоинты сами по себе дают catchup + live в одном Flow: клиент - * передаёт `?after=`, сервер сначала отдаёт буферизованные события - * с `date > after`, потом переключается на live. Отдельный replay-endpoint - * не нужен — для событий старше буфера клиент должен идти в - * `:message-log-api` (полный audit log). + * ``` + * POST {path}/conversations + * GET {path}/conversations + * GET {path}/conversations/{id} + * PATCH {path}/conversations/{id} + * DELETE {path}/conversations/{id} + * POST {path}/conversations/{id}/messages + * POST {path}/conversations/{id}/interrupt + * GET {path}/conversations/{id}/messages + * GET {path}/conversations/{id}/events (SSE: события хода) + * GET {path}/events (SSE: agent-level events) + * GET {path}/health + * GET {path}/journal/conversations/{id}/messages + * GET {path}/outbox/events (SSE: outbox catchup+live) + * ``` + * + * Под-prefix'ы `/journal` и `/outbox` выбраны чтобы не пересекаться с + * существующим `/events` (agent-level SSE) и + * `/conversations/{id}/messages` (proto-Message'ы, не raw records). */ fun Route.agentikAgent( agent: Agent, @@ -54,6 +59,11 @@ fun Route.agentikAgent( this.token = token } } + // Proto-роуты: диалоги, send/interrupt, events (agent-level). agentikRoutes(agent) + + // Storage-фасады: те же `this` (роут агента), свои под-prefix'ы. + journalRoutes(agent.journal) + outboxRoutes(agent.outbox) } } diff --git a/server/src/commonMain/kotlin/pw/binom/agentik/server/OutboxRoutes.kt b/server/src/commonMain/kotlin/pw/binom/agentik/server/OutboxRoutes.kt new file mode 100644 index 0000000..10ff649 --- /dev/null +++ b/server/src/commonMain/kotlin/pw/binom/agentik/server/OutboxRoutes.kt @@ -0,0 +1,44 @@ +package pw.binom.agentik.server + +import io.ktor.server.routing.Route +import io.ktor.server.routing.get +import io.ktor.server.routing.route +import pw.binom.agentik.outbox.OutboxStore +import pw.binom.agentik.proto.CommonEvent + +/** + * HTTP-фасад для [OutboxStore] (bounded-tail live event stream агента). + * + * **Основа:** функция-продолжение для [Route.agentikAgent] — внутри неё уже + * создан роут под `{path}` агента. [outboxRoutes] добавляет под-prefix + * [path] (по умолчанию `"/outbox"`) к **этому же** родительскому роуту, + * итоговый URL = `{path агента}/outbox/...`. + * + * **Endpoint'ы под `{path}/outbox`:** + * - `GET /events?after=` — SSE (catchup + live) в формате `data: \n\n`, + * где `` — сериализованный [CommonEvent]. + * Семантика `after` идентична [OutboxStore.events]: + * - `after` отсутствует → только live (события с момента подписки). + * - `after` задан → сначала catchup всех буферизованных событий с + * `date > after`, потом live. + * + * **Покрытие:** outbox — это короткий bounded tail с auto-TTL. Для событий + * старше буфера клиент должен идти в `/journal/conversations/{id}/messages` + * (полный audit log), см. KDoc [OutboxStore]. + * + * **Read-only:** [OutboxStore] не имеет `append` — запись только через + * writer-референс, который ChatAgent держит внутри (тип `MutableOutboxStore`, + * не выставлен наружу через [pw.binom.agentik.proto.Agent]). + */ +fun Route.outboxRoutes( + outbox: OutboxStore, + path: String = "/outbox", +) { + route(path) { + get("/events") { + val after = call.parseAfter() ?: return@get + // SSE-стрим: catchup (если `after` != DISTANT_PAST) + live tail. + call.streamJsonSse(outbox.events(after), CommonEvent.serializer()) + } + } +} diff --git a/server/src/commonMain/kotlin/pw/binom/agentik/server/Routes.kt b/server/src/commonMain/kotlin/pw/binom/agentik/server/Routes.kt index 1c2e2e8..b726cdf 100644 --- a/server/src/commonMain/kotlin/pw/binom/agentik/server/Routes.kt +++ b/server/src/commonMain/kotlin/pw/binom/agentik/server/Routes.kt @@ -16,6 +16,7 @@ import io.ktor.server.routing.post import io.ktor.utils.io.writeStringUtf8 import kotlinx.coroutines.flow.Flow import kotlinx.coroutines.flow.catch +import kotlinx.coroutines.flow.map import kotlinx.serialization.KSerializer import kotlinx.serialization.json.Json import pw.binom.agentik.proto.Agent @@ -132,7 +133,13 @@ internal fun Route.agentikRoutes(agent: Agent) { get("/events") { val after = call.parseAfter() ?: return@get - call.streamJsonSse(agent.events(after), AgentEvent.serializer()) + // agent.outbox.agentEvents(after) возвращает Flow; + // распаковываем .event для обратной совместимости с прежним + // форматом (когда был Agent.events(): Flow). + call.streamJsonSse( + agent.outbox.agentEvents(after).map { it.event }, + AgentEvent.serializer(), + ) } /** @@ -142,7 +149,7 @@ internal fun Route.agentikRoutes(agent: Agent) { */ get("/events/all") { val after = call.parseAfter() ?: return@get - call.streamJsonSse(agent.allEvents(after), CommonEvent.serializer()) + call.streamJsonSse(agent.outbox.events(after), CommonEvent.serializer()) } } @@ -152,7 +159,7 @@ internal fun Route.agentikRoutes(agent: Agent) { * Парсит query-параметр `after` как ISO-8601 [Instant]. Отсутствие = [Instant.DISTANT_PAST]. * При невалидном значении отвечает 400 и возвращает `null`. */ -private suspend fun ApplicationCall.parseAfter(): Instant? { +internal suspend fun ApplicationCall.parseAfter(): Instant? { val raw = request.queryParameters["after"] if (raw == null) return Instant.DISTANT_PAST return try { @@ -163,7 +170,7 @@ private suspend fun ApplicationCall.parseAfter(): Instant? { } } -private suspend fun ApplicationCall.streamJsonSse( +internal suspend fun ApplicationCall.streamJsonSse( flow: Flow, serializer: KSerializer, json: Json = agentikJson, diff --git a/server/src/commonTest/kotlin/pw/binom/agentik/server/BearerTokenTest.kt b/server/src/commonTest/kotlin/pw/binom/agentik/server/BearerTokenTest.kt index 092bd81..4927e08 100644 --- a/server/src/commonTest/kotlin/pw/binom/agentik/server/BearerTokenTest.kt +++ b/server/src/commonTest/kotlin/pw/binom/agentik/server/BearerTokenTest.kt @@ -11,16 +11,17 @@ import io.ktor.server.cio.CIO as ServerCIO import io.ktor.server.engine.EmbeddedServer import io.ktor.server.engine.embeddedServer import io.ktor.server.routing.routing -import kotlinx.coroutines.flow.Flow import kotlinx.coroutines.flow.emptyFlow import kotlinx.coroutines.runBlocking +import pw.binom.agentik.journal.JournalStore +import pw.binom.agentik.journal.MessageRecord +import pw.binom.agentik.outbox.CommonEvent +import pw.binom.agentik.outbox.OutboxStore import pw.binom.agentik.proto.Agent -import pw.binom.agentik.proto.AgentEvent -import pw.binom.agentik.proto.CommonEvent import pw.binom.agentik.proto.Conversation +import kotlin.time.Instant import kotlin.test.Test import kotlin.test.assertEquals -import kotlin.time.Instant /** * Тесты route-scoped плагина [BearerTokenPlugin]: @@ -30,13 +31,32 @@ import kotlin.time.Instant */ class BearerTokenTest { - private class FakeAgent(override val id: String = "test") : Agent { + /** + * Stub-реализации storage-handle для теста bearer-токена — контент не + * используется, тесты проверяют только что endpoint'ы закрыты/открыты + * по токену. Storage routes регистрируются сразу в [agentikAgent] и + * вычитывают [journal]/[outbox] eagerly, поэтому возвращаем noop-реализации + * (а не `error("...")`), иначе старт сервера валится на инициализации. + */ + private class FakeAgent( + override val id: String = "test", + ) : Agent { + override val journal: JournalStore = object : JournalStore { + override suspend fun list(conversationId: String, after: Instant, offset: Int, limit: Int) = emptyList() + override fun listFlow(conversationId: String, after: Instant, pageSize: Int) = emptyFlow() + override fun close() {} + } + override val outbox: OutboxStore = object : OutboxStore { + override fun events(after: Instant?) = emptyFlow() + override fun agentEvents(after: Instant?) = emptyFlow() + override fun conversationEvents(after: Instant?, conversationId: String?) = emptyFlow() + override suspend fun earliestEventDate(): Instant = Instant.DISTANT_PAST + override fun close() {} + } override fun createConversation(temp: Boolean): Conversation = TODO("not needed by tests") override suspend fun getConversation(id: String): Conversation? = null override suspend fun deleteConversation(id: String): Boolean = false override suspend fun getConversations(offset: Int, limit: Int): List = emptyList() - override fun events(after: Instant): Flow = emptyFlow() - override fun allEvents(after: Instant): Flow = emptyFlow() } private suspend fun startServer(token: String?): Pair, Int> { diff --git a/settings.gradle.kts b/settings.gradle.kts index afa6354..882f8f9 100644 --- a/settings.gradle.kts +++ b/settings.gradle.kts @@ -46,7 +46,8 @@ include(":server") include(":client") // CLI-клиент поверх :client — REPL со slash-командами и стримингом ответов. // KMP со всеми целями (jvm + весь натив), jvm-таргет собирается как shadowJar. -include(":agentik-cli") +// include(":agentik-cli") — отключено 2026-09-21: пользователь временно вывел +// из сборки. Папка agentik-cli/ осталась на диске для возможного возврата. // TUI-клиент поверх :client — Compose-style UI (Mosaic от Jake Wharton), // рендерится в ANSI-терминал. KMP со всеми desktop-целями (без ios). // include(":agentik-tui") — отключено 2026-09-17: пользователь признал TUI-подход неудачным. diff --git a/standalone/build.gradle.kts b/standalone/build.gradle.kts index c148d65..a6cdb6c 100644 --- a/standalone/build.gradle.kts +++ b/standalone/build.gradle.kts @@ -42,8 +42,9 @@ kotlin { implementation(project(":proto")) implementation(project(":server")) implementation(project(":message-store-api")) - implementation(project(":message-log-api")) - implementation(project(":working-memory-api")) + implementation(project(":journal-api")) + implementation(project(":outbox-api")) + implementation(project(":context-api")) // Commons implementation(libs.kotlinx.coroutines.core) @@ -67,13 +68,14 @@ kotlin { implementation(project(":memory-vector")) } implementation(project(":message-store-api")) - implementation(project(":working-memory-api")) + implementation(project(":journal-api")) + implementation(project(":outbox-api")) + implementation(project(":context-api")) implementation(project(":storage-ksqlite")) // Новый единый канал событий агента — заменил старые // `agentEvents: MutableSharedFlow` и per-conv `ConversationEvents._flow`. // Bounded tail + auto-TTL, generic CommonEvent envelope. - implementation(project(":event-store")) implementation(project(":event-store-in-memory")) implementation(project(":agent-toolsets")) // Generic LLM-side tools (LlmReflector, SkillMiner, LlmMemoryReviewer, diff --git a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/A2aBridge.kt b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/A2aBridge.kt index 33ccb6f..768922e 100644 --- a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/A2aBridge.kt +++ b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/A2aBridge.kt @@ -10,10 +10,10 @@ import pw.binom.a2a.model.Message import pw.binom.a2a.model.Role import pw.binom.a2a.model.TextPart import pw.binom.a2a.server.AgentHandler +import pw.binom.agentik.outbox.Event import pw.binom.agentik.proto.Agent import pw.binom.agentik.proto.Content import pw.binom.agentik.proto.Conversation -import pw.binom.agentik.proto.Event import java.util.concurrent.ConcurrentHashMap private val log = KotlinLogging.logger {} diff --git a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/DebugRoutes.kt b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/DebugRoutes.kt index 9f27c38..1d93eb8 100644 --- a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/DebugRoutes.kt +++ b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/DebugRoutes.kt @@ -16,7 +16,7 @@ import pw.binom.agentik.llm.tools.LlmReflector import pw.binom.agentik.llm.tools.SkillMiner import pw.binom.agentik.standalone.agent.memory.Curator import pw.binom.agentik.skills.SkillStore -import pw.binom.agentik.workingMemory.WorkingMemoryStore +import pw.binom.agentik.context.ContextStore /** * Debug-эндпоинты для ручного триггерирования фоновых фич (без ожидания @@ -37,7 +37,7 @@ import pw.binom.agentik.workingMemory.WorkingMemoryStore */ internal fun Route.debugRoutes( agent: Agent, - workingMemoryStore: WorkingMemoryStore, + workingMemoryStore: ContextStore, reflectionStore: ReflectionStore, reflector: LlmReflector?, skillMiner: SkillMiner?, @@ -118,14 +118,14 @@ internal fun Route.debugRoutes( * (для debug-триггеров reflector/miner; та же логика, что у хуков * [ChatConversation]). */ -internal suspend fun recentTurns(workingMemoryStore: WorkingMemoryStore, conversationId: String, maxTurns: Int): List { +internal suspend fun recentTurns(workingMemoryStore: ContextStore, conversationId: String, maxTurns: Int): List { val rows = workingMemoryStore.list(conversationId) val pairs = mutableListOf() var pendingUser: String? = null for (row in rows) { when (val e = row.entry) { - is pw.binom.agentik.workingMemory.WorkingMemoryEntry.User -> pendingUser = e.content.text() - is pw.binom.agentik.workingMemory.WorkingMemoryEntry.Assistant -> { + is pw.binom.agentik.context.WorkingMemoryEntry.User -> pendingUser = e.content.text() + is pw.binom.agentik.context.WorkingMemoryEntry.Assistant -> { val user = pendingUser ?: "" pendingUser = null pairs += ConversationTurn(userMessage = user, assistantMessage = e.content.text()) @@ -137,5 +137,5 @@ internal suspend fun recentTurns(workingMemoryStore: WorkingMemoryStore, convers } /** Текстовое содержимое записей working memory (Text-контент, без картинок). */ -internal fun List.text(): String = - filterIsInstance().joinToString("\n") { it.body } +internal fun List.text(): String = + filterIsInstance().joinToString("\n") { it.body } diff --git a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/BackgroundScheduler.kt b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/BackgroundScheduler.kt index 7a7a4aa..fb5eb34 100644 --- a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/BackgroundScheduler.kt +++ b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/BackgroundScheduler.kt @@ -12,10 +12,10 @@ import pw.binom.agentik.memory.ConversationTurn import pw.binom.agentik.memory.MemoryReviewer import pw.binom.agentik.memory.MemoryStore import pw.binom.agentik.skills.SkillStore -import pw.binom.agentik.messageLog.Content +import pw.binom.agentik.journal.Content import pw.binom.agentik.messageStore.ReflectionStore -import pw.binom.agentik.workingMemory.WorkingMemoryEntry -import pw.binom.agentik.workingMemory.WorkingMemoryStore +import pw.binom.agentik.context.WorkingMemoryEntry +import pw.binom.agentik.context.ContextStore import pw.binom.agentik.llm.tools.LlmReflector import pw.binom.agentik.llm.tools.SkillMiner import java.util.concurrent.atomic.AtomicLong @@ -46,7 +46,7 @@ internal data class BackgroundConfig( internal class BackgroundScheduler( private val state: ConversationState, - private val workingMemory: WorkingMemoryStore, + private val workingMemory: ContextStore, private val config: BackgroundConfig, private val backgroundEvents: BackgroundEventBus, ) { diff --git a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ChatAgent.kt b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ChatAgent.kt index e582854..033b86d 100644 --- a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ChatAgent.kt +++ b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ChatAgent.kt @@ -16,9 +16,11 @@ import pw.binom.agentik.memory.MemoryPrefetcher import pw.binom.agentik.memory.MemoryReviewer import pw.binom.agentik.memory.MemorySystemGuidance import pw.binom.agentik.proto.Agent as ProtoAgent -import pw.binom.agentik.proto.AgentEvent -import pw.binom.agentik.eventStore.MutableEventStore -import pw.binom.agentik.proto.CommonEvent +import pw.binom.agentik.outbox.AgentEvent +import pw.binom.agentik.outbox.CommonEvent +import pw.binom.agentik.outbox.MutableOutboxStore +import pw.binom.agentik.journal.JournalStore +import pw.binom.agentik.outbox.OutboxStore import pw.binom.agentik.proto.Conversation as ProtoConversation import pw.binom.agentik.proto.Event as ProtoEvent import pw.binom.agentik.skills.SkillCatalog @@ -30,8 +32,8 @@ import pw.binom.agentik.messageStore.ConversationStore import pw.binom.agentik.messageStore.Ids import pw.binom.agentik.messageStore.Reflection import pw.binom.agentik.messageStore.ReflectionStore -import pw.binom.agentik.messageLog.MutableMessageStore -import pw.binom.agentik.workingMemory.WorkingMemoryStore +import pw.binom.agentik.journal.MutableJournalStore +import pw.binom.agentik.context.ContextStore import pw.binom.agentik.toolsets.DisableToolsetTool import pw.binom.agentik.toolsets.EnableToolsetTool import pw.binom.agentik.toolsets.NamedTool @@ -64,8 +66,8 @@ import pw.binom.litert.LiteLlm class ChatAgent( override val id: String, private val conversationStore: ConversationStore, - private val messageStore: MutableMessageStore, - private val workingMemoryStore: WorkingMemoryStore, + private val messageStore: MutableJournalStore, + private val workingMemoryStore: ContextStore, private val reflectionStore: ReflectionStore, private val llm: LiteLlm, private val llmConfig: LlmConfig, @@ -205,42 +207,46 @@ class ChatAgent( * - `agentEvents: MutableSharedFlow` (agent lifecycle) * - per-conv `ConversationEvents._flow: MutableSharedFlow` * - * Теперь оба пишут сюда через [MutableEventStore.append], а consumer'ы + * Теперь оба пишут сюда через [MutableOutboxStore.append], а consumer'ы * читают через [EventStore.events]/[conversationEvents]/[agentEvents]. * * **Live tail + auto-TTL** — клиенты больше не должны заботиться о persistence * или подписке на два отдельных канала. */ - private val eventStore: MutableEventStore = pw.binom.agentik.eventStore.inmemory.InMemoryEventStore( + private val eventStore: MutableOutboxStore = pw.binom.agentik.eventStore.inmemory.InMemoryEventStore( maxMessages = null, ttl = null, ) + /** + * Read-only view of [messageStore] для HTTP-фасада в `:server` + * (`Route.agentikAgent` → `/journal/...` endpoint'ы). + * + * [MutableJournalStore] == [MutableJournalStore] (typealias), поэтому + * narrowing с мутабельного writer'а на read-only view тривиальна. + * Используется [Agent.journal] override'ом из `:proto`. + */ + override val journal: JournalStore + get() = messageStore + + /** + * Read-only view of [eventStore] для HTTP-фасада в `:server` + * (`Route.agentikAgent` → `/outbox/...` endpoint'ы). + * + * [MutableOutboxStore] == [MutableOutboxStore] (typealias), narrowing тривиальна. + * Используется [Agent.outbox] override'ом из `:proto`. + */ + override val outbox: OutboxStore + get() = eventStore + /** Защищает карту живых диалогов. */ private val liveLock = Mutex() private val live: MutableMap = HashMap() - /** Live-подписка на события уровня агента (создание/удаление/переименование). */ - override fun events(after: Instant): Flow = - eventStore.agentEvents(after).map { it.event } - - /** - * Все события в одном потоке: agent lifecycle + events всех диалогов. - * - * Snapshot живых диалогов берётся на момент подписки. Новые Created-Event'ы - * НЕ переподписывают — это ответственность caller'а (см. KDoc в :event-store - * про reconnect pattern + gap detection через [eventStore.earliestEventDate]). - */ - override fun allEvents(after: Instant): Flow = flow { - emitAll(eventStore.agentEvents(after)) - live.values - .asSequence() - .filterNot { it.isClosed } - .forEach { conv: ProtoConversation -> - val cid: String = conv.id - emitAll(eventStore.conversationEvents(after, cid).map { it as CommonEvent.Conversation }) - } - } + // events()/allEvents() больше НЕ override'ятся — эти методы удалены + // из :proto.Agent после миграции событий в outbox-сущность. + // Live-подписки теперь идут через Agent.outbox.events() / .agentEvents() + // (см. README :proto для контракта catchup+live). override fun createConversation(temp: Boolean): ProtoConversation { val now = now() diff --git a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/CompactionCoordinator.kt b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/CompactionCoordinator.kt index 4eb304c..1518fb1 100644 --- a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/CompactionCoordinator.kt +++ b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/CompactionCoordinator.kt @@ -6,10 +6,10 @@ import pw.binom.agentik.memory.ConversationTurn import pw.binom.agentik.memory.MemoryReviewer import pw.binom.agentik.memory.MemoryStore import pw.binom.agentik.standalone.agent.memory.materializeReviewNote -import pw.binom.agentik.messageLog.Content -import pw.binom.agentik.workingMemory.WorkingMemoryEntry -import pw.binom.agentik.workingMemory.WorkingMemoryRow -import pw.binom.agentik.workingMemory.WorkingMemoryStore +import pw.binom.agentik.journal.Content +import pw.binom.agentik.context.WorkingMemoryEntry +import pw.binom.agentik.context.WorkingMemoryRow +import pw.binom.agentik.context.ContextStore import pw.binom.litert.LiteContentPart import pw.binom.litert.LiteConversation import pw.binom.litert.LiteConversationConfig @@ -26,7 +26,7 @@ internal class CompactionCoordinator( private val contextCompactor: ContextCompactor?, private val memoryReviewer: MemoryReviewer?, private val memoryStoreForReview: MemoryStore?, - private val workingMemory: WorkingMemoryStore, + private val workingMemory: ContextStore, private val liteLlm: LiteLlm, private val systemPrompt: String, private val backgroundEvents: BackgroundEventBus, diff --git a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ContextBuilder.kt b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ContextBuilder.kt index 3f14102..9c76fa5 100644 --- a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ContextBuilder.kt +++ b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ContextBuilder.kt @@ -3,8 +3,8 @@ package pw.binom.agentik.standalone.agent import mu.KotlinLogging import pw.binom.agentik.memory.MemoryPrefetcher import pw.binom.litert.LiteContentPart -import pw.binom.agentik.messageLog.MessageContext -import pw.binom.agentik.messageLog.MessageOrigin +import pw.binom.agentik.journal.MessageContext +import pw.binom.agentik.journal.MessageOrigin internal class ContextBuilder( private val memoryPrefetcher: MemoryPrefetcher?, diff --git a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ConversationEvents.kt b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ConversationEvents.kt index 78a72ec..c994c0d 100644 --- a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ConversationEvents.kt +++ b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ConversationEvents.kt @@ -3,12 +3,12 @@ package pw.binom.agentik.standalone.agent import kotlinx.coroutines.runBlocking import kotlinx.coroutines.flow.Flow import kotlinx.coroutines.flow.map -import pw.binom.agentik.eventStore.MutableEventStore -import pw.binom.agentik.proto.CommonEvent +import pw.binom.agentik.outbox.MutableOutboxStore +import pw.binom.agentik.outbox.CommonEvent import pw.binom.agentik.proto.Event as ProtoEvent internal class ConversationEvents( - private val globalEventStore: MutableEventStore, + private val globalEventStore: MutableOutboxStore, private val conversationId: String, ) { fun tryEmit(event: ProtoEvent): Boolean { diff --git a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ConversationLoop.kt b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ConversationLoop.kt index 2edbaea..bce0a50 100644 --- a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ConversationLoop.kt +++ b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ConversationLoop.kt @@ -24,21 +24,21 @@ import pw.binom.agentik.memory.MemoryReviewer import pw.binom.agentik.memory.MemoryStore import pw.binom.agentik.proto.Content as ProtoContent import pw.binom.agentik.proto.Conversation as ProtoConversation -import pw.binom.agentik.proto.Event as ProtoEvent +import pw.binom.agentik.outbox.Event as ProtoEvent +import pw.binom.agentik.messageStore.ReflectionStore import pw.binom.agentik.proto.Message as ProtoMessage import pw.binom.agentik.proto.MessageContext as ProtoMessageContext import pw.binom.agentik.skills.SkillStore -import pw.binom.agentik.messageLog.Content +import pw.binom.agentik.journal.Content import pw.binom.agentik.messageStore.ConversationRecord import pw.binom.agentik.messageStore.ConversationStore -import pw.binom.agentik.messageLog.MessageContext -import pw.binom.agentik.messageLog.MessageOrigin -import pw.binom.agentik.messageLog.MessageRecord -import pw.binom.agentik.messageLog.MutableMessageStore -import pw.binom.agentik.messageStore.ReflectionStore -import pw.binom.agentik.messageLog.TurnTokens -import pw.binom.agentik.workingMemory.WorkingMemoryEntry -import pw.binom.agentik.workingMemory.WorkingMemoryStore +import pw.binom.agentik.journal.MessageContext +import pw.binom.agentik.journal.MessageOrigin +import pw.binom.agentik.journal.MessageRecord +import pw.binom.agentik.journal.MutableJournalStore as MutableJournalStore +import pw.binom.agentik.journal.TurnTokens +import pw.binom.agentik.context.WorkingMemoryEntry +import pw.binom.agentik.context.ContextStore import pw.binom.agentik.toolsets.ToolsetDispatchPolicy import pw.binom.litert.LiteContentPart import pw.binom.litert.LiteConversation @@ -55,10 +55,10 @@ import pw.binom.agentik.toolsets.NamedTool class ConversationLoop( record: ConversationRecord, private val conversationStore: ConversationStore, - private val messageStore: MutableMessageStore, - private val workingMemoryStore: WorkingMemoryStore, + private val messageStore: MutableJournalStore, + private val workingMemoryStore: ContextStore, private val reflectionStore: ReflectionStore?, - private val eventStore: pw.binom.agentik.eventStore.MutableEventStore, + private val eventStore: pw.binom.agentik.outbox.MutableOutboxStore, private val llm: LiteLlm, private val systemPrompt: String, private val tools: List = emptyList(), diff --git a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ToolDispatcher.kt b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ToolDispatcher.kt index 2f7ec40..c7830fc 100644 --- a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ToolDispatcher.kt +++ b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ToolDispatcher.kt @@ -4,10 +4,10 @@ import kotlinx.coroutines.CancellationException import kotlinx.coroutines.Job import kotlinx.coroutines.async import mu.KotlinLogging -import pw.binom.agentik.proto.Event as ProtoEvent -import pw.binom.agentik.messageLog.MessageRecord -import pw.binom.agentik.messageLog.MutableMessageStore -import pw.binom.agentik.workingMemory.WorkingMemoryEntry +import pw.binom.agentik.outbox.Event as ProtoEvent +import pw.binom.agentik.journal.MessageRecord +import pw.binom.agentik.journal.MutableJournalStore as MutableJournalStore +import pw.binom.agentik.context.WorkingMemoryEntry import pw.binom.agentik.toolsets.ToolsetDispatchPolicy import pw.binom.litert.LiteToolCall import pw.binom.litert.LiteTool @@ -16,7 +16,7 @@ import pw.binom.agentik.toolsets.NamedTool internal class ToolDispatcher( private val state: ConversationState, - private val messageStore: MutableMessageStore, + private val messageStore: MutableJournalStore, private val events: ConversationEvents, private val backgroundEvents: BackgroundEventBus, private val toolsByName: MutableMap, diff --git a/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/ChatAgentTest.kt b/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/ChatAgentTest.kt index 0152c75..285a113 100644 --- a/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/ChatAgentTest.kt +++ b/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/ChatAgentTest.kt @@ -7,15 +7,15 @@ import kotlinx.coroutines.flow.flowOf import kotlinx.coroutines.flow.toList import kotlinx.coroutines.launch import kotlinx.coroutines.test.runTest -import pw.binom.agentik.proto.AgentEvent +import pw.binom.agentik.outbox.AgentEvent import pw.binom.agentik.proto.Content -import pw.binom.agentik.proto.Event as ProtoEvent +import pw.binom.agentik.outbox.Event as ProtoEvent import pw.binom.agentik.skills.SkillCatalog import pw.binom.agentik.skills.SkillFile import pw.binom.agentik.standalone.llm.LlmBackend import pw.binom.agentik.standalone.llm.LlmConfig -import pw.binom.agentik.messageLog.MessageRecord -import pw.binom.agentik.workingMemory.WorkingMemoryEntry +import pw.binom.agentik.journal.MessageRecord +import pw.binom.agentik.context.WorkingMemoryEntry import pw.binom.agentik.storage.ksqlite.KsqliteStores import pw.binom.litert.LiteContentPart import pw.binom.litert.LiteConversation @@ -163,10 +163,10 @@ class ChatAgentTest { // добавим сообщение, чтобы потом убедиться, что каскад сработал sqliteStores.messages.append( - pw.binom.agentik.messageLog.MessageRecord.UserMessage( + pw.binom.agentik.journal.MessageRecord.UserMessage( id = "m1", conversationId = id, - content = listOf(pw.binom.agentik.messageLog.Content.Text("hi")), + content = listOf(pw.binom.agentik.journal.Content.Text("hi")), createdAt = Instant.fromEpochMilliseconds(1_700_000_000_000), ), ) @@ -193,11 +193,11 @@ class ChatAgentTest { // user message записан в audit + working memory val msgs = sqliteStores.messages.listFlow(conv.id, Instant.DISTANT_PAST).toList() assertEquals(2, msgs.size) - assertEquals("hi", (msgs[0] as pw.binom.agentik.messageLog.MessageRecord.UserMessage).content.let { - (it[0] as pw.binom.agentik.messageLog.Content.Text).body + assertEquals("hi", (msgs[0] as pw.binom.agentik.journal.MessageRecord.UserMessage).content.let { + (it[0] as pw.binom.agentik.journal.Content.Text).body }) - assertEquals("hello back", (msgs[1] as pw.binom.agentik.messageLog.MessageRecord.AssistantMessage).content.let { - (it[0] as pw.binom.agentik.messageLog.Content.Text).body + assertEquals("hello back", (msgs[1] as pw.binom.agentik.journal.MessageRecord.AssistantMessage).content.let { + (it[0] as pw.binom.agentik.journal.Content.Text).body }) } @@ -330,8 +330,8 @@ class ChatAgentTest { val msgs = sqliteStores.messages.listFlow(conv.id, Instant.DISTANT_PAST).toList() assertEquals(2, msgs.size) - assertIs(msgs[0]) - val err = assertIs(msgs[1]) + assertIs(msgs[0]) + val err = assertIs(msgs[1]) assertEquals("boom from llm", err.message) // backfill через getMessages (polling/reconnect) тоже видит ошибку @@ -375,7 +375,7 @@ class ChatAgentTest { // audit: только user (assistant не успел сгенериться) val msgs = sqliteStores.messages.listFlow(conv.id, Instant.DISTANT_PAST).toList() assertEquals(1, msgs.size) - assertIs(msgs[0]) + assertIs(msgs[0]) // working memory: только user (assistant skipped because пустой) val wm = sqliteStores.workingMemory.list(conv.id) @@ -431,7 +431,7 @@ class ChatAgentTest { // audit: user + toolcall + toolresult (tool выполнился), assistant может быть val msgs = sqliteStores.messages.listFlow(conv.id, Instant.DISTANT_PAST).toList() - val toolResult = msgs.filterIsInstance().firstOrNull() + val toolResult = msgs.filterIsInstance().firstOrNull() assertNotNull(toolResult, "tool result должен быть в audit — tool выполнился нормально") val toolResultResult = toolResult!!.result!! assertTrue(toolResultResult.contains("echo"), "tool result содержит реальный ответ тулы: $toolResultResult") @@ -574,7 +574,10 @@ class ChatAgentTest { val agent = newAgent() val events = mutableListOf() val job = launch(start = kotlinx.coroutines.CoroutineStart.UNDISPATCHED) { - agent.events(Instant.DISTANT_PAST).collect { events.add(it) } + // Agent.events() удалён из :proto — события живут в + // agent.outbox.agentEvents(): Flow; + // распаковываем .event для получения AgentEvent. + agent.outbox.agentEvents(Instant.DISTANT_PAST).collect { events.add(it.event) } } val conv = agent.createConversation(temp = false) agent.deleteConversation(conv.id) diff --git a/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/CompactionTest.kt b/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/CompactionTest.kt index a0e8bd9..70b63a3 100644 --- a/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/CompactionTest.kt +++ b/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/CompactionTest.kt @@ -92,7 +92,7 @@ class CompactionTest { } val wm = sqliteStores.workingMemory.list(conv.id) // Без compaction все ходы остаются в памяти (System + 10 user/assistant = 21 строк). - val summaries = wm.filter { it.entry is pw.binom.agentik.workingMemory.WorkingMemoryEntry.Summary } + val summaries = wm.filter { it.entry is pw.binom.agentik.context.WorkingMemoryEntry.Summary } assertEquals(0, summaries.size, "compaction must not run without contextWindow") } @@ -104,7 +104,7 @@ class CompactionTest { conv.send(listOf(ProtoContent.Text("first"))) val wm = sqliteStores.workingMemory.list(conv.id) // System + User + Assistant = 3. Без compactor — никаких Summary. - val summaries = wm.filter { it.entry is pw.binom.agentik.workingMemory.WorkingMemoryEntry.Summary } + val summaries = wm.filter { it.entry is pw.binom.agentik.context.WorkingMemoryEntry.Summary } assertEquals(0, summaries.size, "no compaction runs without compactor") } @@ -122,10 +122,10 @@ class CompactionTest { assertTrue(compactor.calls > 0, "compactor must be called at least once when above threshold") // В working memory должна появиться Summary. val wm = sqliteStores.workingMemory.list(conv.id) - val summaries = wm.filter { it.entry is pw.binom.agentik.workingMemory.WorkingMemoryEntry.Summary } + val summaries = wm.filter { it.entry is pw.binom.agentik.context.WorkingMemoryEntry.Summary } assertTrue(summaries.isNotEmpty(), "at least one Summary entry should be present after compaction") // Summary-текст — то, что вернул наш compactor. - val summaryText = (summaries.first().entry as pw.binom.agentik.workingMemory.WorkingMemoryEntry.Summary).text + val summaryText = (summaries.first().entry as pw.binom.agentik.context.WorkingMemoryEntry.Summary).text assertTrue(summaryText.startsWith("**Goal**"), "summary text should come from compactor: $summaryText") } @@ -165,8 +165,8 @@ class CompactionTest { // Должны быть: System + хотя бы один Summary + последние KEEP_RECENT_TURNS ходов. // KEEP_RECENT_TURNS = 4 → user/assistant последних двух ходов (third + second) могут быть не тронуты. val userAssistantCount = wm.count { - it.entry is pw.binom.agentik.workingMemory.WorkingMemoryEntry.User || - it.entry is pw.binom.agentik.workingMemory.WorkingMemoryEntry.Assistant + it.entry is pw.binom.agentik.context.WorkingMemoryEntry.User || + it.entry is pw.binom.agentik.context.WorkingMemoryEntry.Assistant } // Минимум 1 ход остаётся (KEEP_RECENT_TURNS). assertTrue(userAssistantCount >= 1, "at least one recent turn must be preserved") diff --git a/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/ContextPrefixTest.kt b/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/ContextPrefixTest.kt index 2d021f9..95ceea2 100644 --- a/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/ContextPrefixTest.kt +++ b/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/ContextPrefixTest.kt @@ -1,9 +1,9 @@ package pw.binom.agentik.standalone.agent -import pw.binom.agentik.messageLog.MessageContext -import pw.binom.agentik.messageLog.MessageOrigin.EVENT -import pw.binom.agentik.messageLog.MessageOrigin.SYSTEM -import pw.binom.agentik.messageLog.MessageOrigin.USER +import pw.binom.agentik.journal.MessageContext +import pw.binom.agentik.journal.MessageOrigin.EVENT +import pw.binom.agentik.journal.MessageOrigin.SYSTEM +import pw.binom.agentik.journal.MessageOrigin.USER import pw.binom.litert.LiteContentPart import kotlin.test.Test import kotlin.test.assertEquals diff --git a/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/persistence/PersistenceTest.kt b/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/persistence/PersistenceTest.kt index 6540d29..2066bbd 100644 --- a/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/persistence/PersistenceTest.kt +++ b/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/persistence/PersistenceTest.kt @@ -1,10 +1,10 @@ package pw.binom.agentik.standalone.persistence -import pw.binom.agentik.messageLog.MessageContext -import pw.binom.agentik.messageLog.MessageOrigin +import pw.binom.agentik.journal.MessageContext +import pw.binom.agentik.journal.MessageOrigin import pw.binom.agentik.messageStore.ConversationRecord -import pw.binom.agentik.messageLog.MessageRecord -import pw.binom.agentik.messageLog.Content -import pw.binom.agentik.workingMemory.WorkingMemoryEntry +import pw.binom.agentik.journal.MessageRecord +import pw.binom.agentik.journal.Content +import pw.binom.agentik.context.WorkingMemoryEntry import kotlinx.coroutines.flow.toList import kotlinx.coroutines.test.runTest diff --git a/storage-inmemory/src/commonMain/kotlin/pw/binom/agentik/storage/inmemory/InMemoryStorage.kt b/storage-inmemory/src/commonMain/kotlin/pw/binom/agentik/storage/inmemory/InMemoryStorage.kt index f750237..64d92fd 100644 --- a/storage-inmemory/src/commonMain/kotlin/pw/binom/agentik/storage/inmemory/InMemoryStorage.kt +++ b/storage-inmemory/src/commonMain/kotlin/pw/binom/agentik/storage/inmemory/InMemoryStorage.kt @@ -3,7 +3,7 @@ package pw.binom.agentik.storage.inmemory import pw.binom.agentik.messageLog.MutableMessageStore import pw.binom.agentik.messageStore.ConversationStore import pw.binom.agentik.messageStore.ReflectionStore -import pw.binom.agentik.workingMemory.WorkingMemoryStore +import pw.binom.agentik.context.ContextStore import kotlin.time.Clock /** @@ -23,7 +23,7 @@ object InMemoryStorage { data class Bundle( val conversationStore: ConversationStore, val messageStore: MutableMessageStore, - val workingMemoryStore: WorkingMemoryStore, + val workingMemoryStore: ContextStore, val reflectionStore: ReflectionStore, ) diff --git a/storage-inmemory/src/commonMain/kotlin/pw/binom/agentik/storage/inmemory/InMemoryWorkingMemoryStore.kt b/storage-inmemory/src/commonMain/kotlin/pw/binom/agentik/storage/inmemory/InMemoryWorkingMemoryStore.kt index bc8a21a..42b7fcb 100644 --- a/storage-inmemory/src/commonMain/kotlin/pw/binom/agentik/storage/inmemory/InMemoryWorkingMemoryStore.kt +++ b/storage-inmemory/src/commonMain/kotlin/pw/binom/agentik/storage/inmemory/InMemoryWorkingMemoryStore.kt @@ -3,9 +3,9 @@ package pw.binom.agentik.storage.inmemory import kotlinx.coroutines.sync.Mutex import kotlinx.coroutines.sync.withLock import pw.binom.agentik.messageStore.Ids -import pw.binom.agentik.workingMemory.WorkingMemoryEntry -import pw.binom.agentik.workingMemory.WorkingMemoryRow -import pw.binom.agentik.workingMemory.WorkingMemoryStore +import pw.binom.agentik.context.WorkingMemoryEntry +import pw.binom.agentik.context.WorkingMemoryRow +import pw.binom.agentik.context.ContextStore import kotlin.time.Instant /** @@ -17,7 +17,7 @@ import kotlin.time.Instant * * Семантика 1:1 с SQLite-имплом (см. `:storage-sqlite` после commit 3). */ -class InMemoryWorkingMemoryStore : WorkingMemoryStore { +class InMemoryWorkingMemoryStore : ContextStore { private val byConv: MutableMap> = mutableMapOf() private val mutex = Mutex() diff --git a/storage-inmemory/src/commonTest/kotlin/pw/binom/agentik/storage/inmemory/InMemoryMessageStoreTest.kt b/storage-inmemory/src/commonTest/kotlin/pw/binom/agentik/storage/inmemory/InMemoryMessageStoreTest.kt index ba1708c..ae5cfe9 100644 --- a/storage-inmemory/src/commonTest/kotlin/pw/binom/agentik/storage/inmemory/InMemoryMessageStoreTest.kt +++ b/storage-inmemory/src/commonTest/kotlin/pw/binom/agentik/storage/inmemory/InMemoryMessageStoreTest.kt @@ -1,7 +1,7 @@ package pw.binom.agentik.storage.inmemory -import pw.binom.agentik.messageLog.Content -import pw.binom.agentik.messageLog.MessageRecord +import pw.binom.agentik.journal.Content +import pw.binom.agentik.journal.MessageRecord import kotlin.test.Test import kotlin.test.assertEquals import kotlin.test.assertNull diff --git a/storage-inmemory/src/commonTest/kotlin/pw/binom/agentik/storage/inmemory/InMemoryWorkingMemoryStoreTest.kt b/storage-inmemory/src/commonTest/kotlin/pw/binom/agentik/storage/inmemory/InMemoryWorkingMemoryStoreTest.kt index 172212a..da0a582 100644 --- a/storage-inmemory/src/commonTest/kotlin/pw/binom/agentik/storage/inmemory/InMemoryWorkingMemoryStoreTest.kt +++ b/storage-inmemory/src/commonTest/kotlin/pw/binom/agentik/storage/inmemory/InMemoryWorkingMemoryStoreTest.kt @@ -1,7 +1,7 @@ package pw.binom.agentik.storage.inmemory -import pw.binom.agentik.messageLog.Content -import pw.binom.agentik.workingMemory.WorkingMemoryEntry +import pw.binom.agentik.journal.Content +import pw.binom.agentik.context.WorkingMemoryEntry import kotlin.test.Test import kotlin.test.assertEquals import kotlin.test.assertNull diff --git a/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteWorkingMemoryStore.kt b/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteWorkingMemoryStore.kt index 42e52d9..7bdf652 100644 --- a/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteWorkingMemoryStore.kt +++ b/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteWorkingMemoryStore.kt @@ -1,9 +1,9 @@ package pw.binom.agentik.storage.ksqlite import kotlinx.serialization.json.Json -import pw.binom.agentik.workingMemory.WorkingMemoryEntry -import pw.binom.agentik.workingMemory.WorkingMemoryRow -import pw.binom.agentik.workingMemory.WorkingMemoryStore +import pw.binom.agentik.context.WorkingMemoryEntry +import pw.binom.agentik.context.WorkingMemoryRow +import pw.binom.agentik.context.ContextStore import pw.binom.agentik.messageStore.Ids import pw.binom.db.ksqlite.SQLiteConnection import pw.binom.db.ksqlite.SQLitePreparedStatement @@ -16,7 +16,7 @@ import kotlinx.coroutines.withContext class KsqliteWorkingMemoryStore( private val connection: SQLiteConnection, -) : WorkingMemoryStore { +) : ContextStore { private val mutex = Mutex() private val json = Json { ignoreUnknownKeys = true } diff --git a/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/MessageCodecs.kt b/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/MessageCodecs.kt index 35c0f6b..beef950 100644 --- a/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/MessageCodecs.kt +++ b/storage-ksqlite/src/commonMain/kotlin/pw/binom/agentik/storage/ksqlite/MessageCodecs.kt @@ -1,9 +1,9 @@ package pw.binom.agentik.storage.ksqlite import kotlinx.serialization.json.Json -import pw.binom.agentik.messageLog.MessageRecord -import pw.binom.agentik.messageLog.decodeBodyPayload -import pw.binom.agentik.messageLog.encodeBodyPayload +import pw.binom.agentik.journal.MessageRecord +import pw.binom.agentik.journal.decodeBodyPayload +import pw.binom.agentik.journal.encodeBodyPayload import pw.binom.db.ksqlite.SQLiteResultSet import kotlin.time.Instant diff --git a/storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteMessageStoreTest.kt b/storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteMessageStoreTest.kt index ac7bbc0..4fe8814 100644 --- a/storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteMessageStoreTest.kt +++ b/storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteMessageStoreTest.kt @@ -2,9 +2,9 @@ package pw.binom.agentik.storage.ksqlite import kotlinx.coroutines.flow.toList import kotlinx.coroutines.test.runTest -import pw.binom.agentik.messageLog.Content -import pw.binom.agentik.messageLog.MessageRecord -import pw.binom.agentik.messageLog.TurnTokens +import pw.binom.agentik.journal.Content +import pw.binom.agentik.journal.MessageRecord +import pw.binom.agentik.journal.TurnTokens import kotlin.test.AfterTest import kotlin.test.BeforeTest import kotlin.test.Test diff --git a/storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteWorkingMemoryStoreTest.kt b/storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteWorkingMemoryStoreTest.kt index 551d980..0ed845d 100644 --- a/storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteWorkingMemoryStoreTest.kt +++ b/storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/KsqliteWorkingMemoryStoreTest.kt @@ -1,8 +1,8 @@ package pw.binom.agentik.storage.ksqlite import kotlinx.coroutines.test.runTest -import pw.binom.agentik.messageLog.Content -import pw.binom.agentik.workingMemory.WorkingMemoryEntry +import pw.binom.agentik.journal.Content +import pw.binom.agentik.context.WorkingMemoryEntry import kotlin.test.AfterTest import kotlin.test.BeforeTest import kotlin.test.Test diff --git a/storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/SchemaMigrationTest.kt b/storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/SchemaMigrationTest.kt index 3175c99..d195d7a 100644 --- a/storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/SchemaMigrationTest.kt +++ b/storage-ksqlite/src/commonTest/kotlin/pw/binom/agentik/storage/ksqlite/SchemaMigrationTest.kt @@ -95,9 +95,9 @@ class SchemaMigrationTest { ) ) stores.messages.append( - pw.binom.agentik.messageLog.MessageRecord.UserMessage( + pw.binom.agentik.journal.MessageRecord.UserMessage( id = "m1", conversationId = "c1", - content = listOf(pw.binom.agentik.messageLog.Content.Text("hi")), + content = listOf(pw.binom.agentik.journal.Content.Text("hi")), createdAt = kotlin.time.Instant.parse("2026-09-15T10:00:01Z"), ) ) diff --git a/working-memory-api/build.gradle.kts b/working-memory-api/build.gradle.kts index a2338bd..a9709e2 100644 --- a/working-memory-api/build.gradle.kts +++ b/working-memory-api/build.gradle.kts @@ -25,6 +25,7 @@ kotlin { sourceSets { commonMain.dependencies { api(project(":message-log-api")) + api(project(":context-api")) api(libs.kotlinx.coroutines.core) api(libs.kotlinx.serialization.core) api(libs.kotlinx.serialization.json) diff --git a/working-memory-api/src/commonMain/kotlin/pw/binom/agentik/workingMemory/WorkingMemoryEntry.kt b/working-memory-api/src/commonMain/kotlin/pw/binom/agentik/workingMemory/WorkingMemoryEntry.kt index b0d7d4a..36478af 100644 --- a/working-memory-api/src/commonMain/kotlin/pw/binom/agentik/workingMemory/WorkingMemoryEntry.kt +++ b/working-memory-api/src/commonMain/kotlin/pw/binom/agentik/workingMemory/WorkingMemoryEntry.kt @@ -1,86 +1,16 @@ package pw.binom.agentik.workingMemory -import kotlinx.serialization.SerialName -import kotlinx.serialization.Serializable -import pw.binom.agentik.messageLog.Content -import pw.binom.agentik.messageLog.MessageContext - /** - * Запись в working memory диалога: ровно то, что агент сейчас видит в - * LLM-контексте. Упорядочено по `order_idx` (заполняется в store при append). + * @Deprecated + * Перенесено в `:context-api`. Используй `pw.binom.agentik.context.WorkingMemoryEntry`. * - * Sealed-иерархия: `User`/`Assistant` (реплики с ссылкой на audit log - * через [sourceMessageId]), `ToolExchange` (синтетическая запись об одном - * tool-вызове + его результате — для replay в LiteMessage(TOOL, ToolResult) - * при пересоздании LiteConv), `Summary` (суммаризация при compaction). + * Backward-compat typealias. Все импорты `pw.binom.agentik.workingMemory.WorkingMemoryEntry` + * продолжают работать как тип `pw.binom.agentik.context.WorkingMemoryEntry` (тот же тип, транспарентно). + * + * Удалить когда все импорты будут на `:context-api`. */ -@Serializable -sealed interface WorkingMemoryEntry { - - /** Ссылка на исходное сообщение в audit log (`message.id`). `null` для синтетических строк. */ - val sourceMessageId: String? - - /** Реплика пользователя. */ - @Serializable - @SerialName("user") - data class User( - override val sourceMessageId: String, - val content: List, - /** - * Контекст инициации хода. Применяется при сборке `LiteConversation`: - * если `origin != USER`, текст префиксуется `[origin] description (sourceId=…)`, - * чтобы модель видела, что её разбудил не пользователь. - * `null` = обычное user-сообщение. - */ - val context: MessageContext? = null, - ) : WorkingMemoryEntry - - /** Реплика ассистента. */ - @Serializable - @SerialName("assistant") - data class Assistant( - override val sourceMessageId: String, - val content: List, - ) : WorkingMemoryEntry - - /** - * Синтетический блок: один tool-вызов + его результат. Синтетический — потому - * что в audit log это две отдельные записи (`MessageRecord.ToolCall` + - * `MessageRecord.ToolResult`), а в working_memory мы храним одной строкой - * для удобства replay'а. - * - * При создании новой LiteConv каждая такая запись превращается в - * `LiteMessage(TOOL, [ToolResult(callId, name, response)])` — LiteRT-LM - * матчит по `name`, `callId` берётся из [sourceMessageId] (= id исходного - * [MessageRecord.ToolCall]). Если [wasCancelled] = true, [resultText] - * содержит маркер `[cancelled by user]` — модель видит честную причину - * отсутствия результата. - * - * [sourceMessageId] = id исходного [MessageRecord.ToolCall] (для трассировки - * в audit log). - */ - @Serializable - @SerialName("tool_exchange") - data class ToolExchange( - override val sourceMessageId: String, - val toolName: String, - val toolArgsJson: String, - val resultText: String, - val wasCancelled: Boolean = false, - ) : WorkingMemoryEntry - - /** - * Синтетический блок: суммаризация старых ходов, сгенерированная при - * compaction'е working memory. Не имеет ссылки на конкретное сообщение - * в audit log — это наша собственная интерпретация контекста. - */ - @Serializable - @SerialName("summary") - data class Summary( - val text: String, - /** Ходы, которые были свёрнуты в этот summary (диапазон order_idx в виде меты). */ - val coversUpToOrderIdx: Long? = null, - ) : WorkingMemoryEntry { - override val sourceMessageId: String? = null - } -} +@Deprecated( + message = "Перенесено в :context-api. Используй pw.binom.agentik.context.WorkingMemoryEntry.", + replaceWith = ReplaceWith("WorkingMemoryEntry", "pw.binom.agentik.context.WorkingMemoryEntry"), +) +typealias WorkingMemoryEntry = pw.binom.agentik.context.WorkingMemoryEntry diff --git a/working-memory-api/src/commonMain/kotlin/pw/binom/agentik/workingMemory/WorkingMemoryStore.kt b/working-memory-api/src/commonMain/kotlin/pw/binom/agentik/workingMemory/WorkingMemoryStore.kt index 86f42e2..01167ac 100644 --- a/working-memory-api/src/commonMain/kotlin/pw/binom/agentik/workingMemory/WorkingMemoryStore.kt +++ b/working-memory-api/src/commonMain/kotlin/pw/binom/agentik/workingMemory/WorkingMemoryStore.kt @@ -1,58 +1,18 @@ package pw.binom.agentik.workingMemory -import kotlin.time.Instant - /** - * Одна строка `working_memory` таблицы (внутреннее представление store). + * @Deprecated + * Перенесено в `:context-api`. Используй `pw.binom.agentik.context.ContextStore`. * - * Используется для тестов и для перестроения [WorkingMemoryEntry] из row. - * Агент не должен с этим типом работать напрямую — он работает с - * [WorkingMemoryEntry] через [WorkingMemoryStore]. + * Этот typealias сохранён временно для backward-compat: существующие impl'ы + * (`KsqliteStores.workingMemory` в `:storage-ksqlite`, тесты) импортируют + * `pw.binom.agentik.workingMemory.WorkingMemoryStore`. Буквально тот же тип, + * что и `pw.binom.agentik.context.ContextStore` — typealias транспарентен. + * + * Удалить когда все импорты будут на `:context-api`. */ -data class WorkingMemoryRow( - val id: String, - val conversationId: String, - val orderIdx: Long, - val sourceMessageId: String?, - val entry: WorkingMemoryEntry, - val createdAt: Instant, +@Deprecated( + message = "Перенесено в :context-api. Используй pw.binom.agentik.context.ContextStore.", + replaceWith = ReplaceWith("ContextStore", "pw.binom.agentik.context.ContextStore"), ) - -/** - * Мутируемое представление LLM-контекста диалога (`working_memory` table). - * - * Аудит-лог — [MessageStore], неизменный; здесь — ровно то, что агент сейчас - * «видит»: системный промпт + реплики + (опционально) суммаризации. - * Строки упорядочены по `order_idx` ASC; `source_message_id` NULL указывает - * на синтетические строки (System, а в v2 — Summary). - * - * Суммаризация / чистка — один атомарный вызов [compact]. - */ -interface WorkingMemoryStore : AutoCloseable { - - /** Добавить запись в конец working memory (новый максимальный `order_idx`). */ - suspend fun append(conversationId: String, entry: WorkingMemoryEntry, now: Instant) - - /** Все строки working memory диалога в порядке отправки. */ - suspend fun list(conversationId: String): List - - /** Очистить working memory диалога (используется при reset/rebuild). */ - suspend fun clear(conversationId: String) - - /** - * Атомарная суммаризация: удаляет все строки с `order_idx` в диапазоне - * `[dropFromOrderIdx, +∞)`. Если [summaryText] непустое — вместо удалённых - * строк вставляется одна синтетическая [WorkingMemoryEntry.Summary] - * с этим текстом и `order_idx = max(old order_idx after delete) + 1` - * (т.е. summary становится хвостом working memory). - * - * Если [summaryText] == null — работает как «отрезать хвост» (v1 поведение). - * - * Возвращает новый максимальный `order_idx` после операции. - */ - suspend fun compact( - dropFromOrderIdx: Long, - conversationId: String, - summaryText: String? = null, - ): Long -} +typealias WorkingMemoryStore = pw.binom.agentik.context.ContextStore