diff --git a/agentik-cli/src/commonMain/kotlin/pw/binom/agentik/cli/commands/MsgsSubcommand.kt b/agentik-cli/src/commonMain/kotlin/pw/binom/agentik/cli/commands/MsgsSubcommand.kt index d0120c0..a2ecacf 100644 --- a/agentik-cli/src/commonMain/kotlin/pw/binom/agentik/cli/commands/MsgsSubcommand.kt +++ b/agentik-cli/src/commonMain/kotlin/pw/binom/agentik/cli/commands/MsgsSubcommand.kt @@ -5,7 +5,7 @@ import kotlinx.cli.default import pw.binom.agentik.cli.AgentikSubcommand import pw.binom.agentik.cli.defaultCliHttpClient import pw.binom.agentik.client.AgentikAgent -import pw.binom.agentik.proto.Content +import pw.binom.agentik.content.Content import pw.binom.agentik.proto.Message import kotlin.time.Instant 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 909f094..31c3938 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 @@ -12,7 +12,7 @@ import pw.binom.agentik.cli.defaultCliHttpClient import pw.binom.agentik.client.AgentikAgent import pw.binom.agentik.outbox.Event import pw.binom.agentik.outbox.OnlineEvent -import pw.binom.agentik.proto.Content +import pw.binom.agentik.content.Content 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 317c163..53693da 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 @@ -5,7 +5,7 @@ import kotlinx.coroutines.Job import kotlinx.coroutines.flow.collect import kotlinx.coroutines.launch import pw.binom.agentik.proto.Agent -import pw.binom.agentik.proto.Content +import pw.binom.agentik.content.Content import pw.binom.agentik.proto.Conversation import pw.binom.agentik.outbox.Event import pw.binom.agentik.outbox.OnlineEvent 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 a1e1c2d..9671a3c 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 @@ -8,11 +8,11 @@ import pw.binom.agentik.journal.JournalStore import pw.binom.agentik.outbox.OutboxStore import pw.binom.agentik.proto.Agent import pw.binom.agentik.proto.AgentInfo -import pw.binom.agentik.proto.Content +import pw.binom.agentik.content.Content import pw.binom.agentik.proto.Conversation import pw.binom.agentik.outbox.Event import pw.binom.agentik.proto.Message -import pw.binom.agentik.proto.MessageContext +import pw.binom.agentik.content.MessageContext import kotlin.time.Instant /** diff --git a/agentik-tui/src/commonTest/kotlin/pw/binom/agentik/tui/TuiBackendTest.kt b/agentik-tui/src/commonTest/kotlin/pw/binom/agentik/tui/TuiBackendTest.kt index 30f3d00..44487b7 100644 --- a/agentik-tui/src/commonTest/kotlin/pw/binom/agentik/tui/TuiBackendTest.kt +++ b/agentik-tui/src/commonTest/kotlin/pw/binom/agentik/tui/TuiBackendTest.kt @@ -3,7 +3,7 @@ package pw.binom.agentik.tui import kotlinx.coroutines.ExperimentalCoroutinesApi import kotlinx.coroutines.test.runCurrent import kotlinx.coroutines.test.runTest -import pw.binom.agentik.proto.Content +import pw.binom.agentik.content.Content import pw.binom.agentik.outbox.Event import kotlin.test.Test import kotlin.test.assertEquals diff --git a/build.gradle.kts b/build.gradle.kts index 3e9969a..764a98e 100644 --- a/build.gradle.kts +++ b/build.gradle.kts @@ -40,6 +40,8 @@ val binomRepoPassword = (findProperty("binom.repo.password") as String? ?: "").t // beforeEvaluate — поздно). val moduleDescriptions: Map = mapOf( "proto" to "agentik :proto — stateful KMP protocol (Agent/Conversation/Message/Event) replacing AG-UI; типы и контракт без сетевой логики.", + "content-api" to "agentik :content-api — общие типы содержимого сообщения (Content/MessageContext/MessageOrigin/TurnTokens) для :proto, :journal-api, :outbox-api.", + "outbox-api" to "agentik :outbox-api — durable (Event) и live-only (OnlineEvent) потоки событий диалога + OutboxStore/OnlineOutbox.", "skills" to "agentik :skills — парсер opencode-style SKILL.md / *.yaml (YAML-frontmatter + markdown body); загружается в system prompt.", "server" to "agentik :server — Ktor-фасад, экспонирующий Agent по HTTP+JSON+SSE под путём /agentik.", "client" to "agentik :client — Ktor-клиент (HTTP+JSON+SSE), превращающий /agentik в Agent/Conversation из :proto.", diff --git a/client/README.md b/client/README.md index 57801aa..7b78f56 100644 --- a/client/README.md +++ b/client/README.md @@ -72,9 +72,11 @@ dependencies { ```kotlin import pw.binom.agentik.client.AgentikAgent -import pw.binom.agentik.proto.Content -import pw.binom.agentik.proto.Event +import pw.binom.agentik.content.Content +import pw.binom.agentik.outbox.Event +import pw.binom.agentik.outbox.OnlineEvent import io.ktor.client.engine.cio.CIO +import kotlinx.coroutines.launch import kotlinx.coroutines.runBlocking import kotlin.time.Clock @@ -87,21 +89,37 @@ fun main() = runBlocking { token = "s3cret", // или null, если не нужен ) - // 2. Открыть диалог, отправить сообщение. val conv = agent.createConversation(temp = false) - conv.send(listOf(Content.Text("Привет"))) - // 3. Собирать streaming-ответ. - conv.events(after = Clock.System.now()).collect { ev -> - when (ev) { - is Event.StartResponse -> println("[start]") - is Event.AppendText -> print(ev.body) - is Event.End -> println("[end]") - is Event.Error -> println("[error: ${ev.message}]") - else -> Unit + // 2. Два независимых потока событий диалога: + // durable (outbox) — целые события, с курсором после переподключения; + // online (OnlineOutbox) — стриминг ответа, только live (без курсора). + launch { + agent.outbox.conversationEvents(after = Clock.System.now(), conversationId = conv.id) + .collect { ce -> + when (val ev = ce.event) { + is Event.AssistantMessage -> println("[answer ready: ${ev.content}]") + is Event.Interrupted -> println("[interrupted]") + is Event.Error -> println("[error: ${ev.message}]") + else -> Unit + } + } + } + launch { + agent.onlineOutbox.onlineEvents(conv.id).collect { ev -> + when (ev) { + is OnlineEvent.Working -> println("[working]") + is OnlineEvent.AppendText -> print(ev.body) + is OnlineEvent.StartResponse -> println("[start]") + is OnlineEvent.End -> println("\n[end]") + else -> Unit + } } } + // 3. Отправить ход (fire-and-forget — ответ придёт по подпискам выше). + conv.send(listOf(Content.Text("Привет"))) + // 4. Чистый shutdown. conv.close() agent.close() @@ -109,7 +127,14 @@ fun main() = runBlocking { ``` **Это весь клиент.** `:server` сам хранит историю, контекст, события. -Ты только получаешь типизированный `Flow` и рендеришь как хочешь. +Ты только получаешь два типизированных `Flow` и рендеришь как хочешь. + +> **Durable vs online.** `Event` (в `agent.outbox`) — «целые» события, их +> можно перезапросить по курсору `after`. `OnlineEvent` (в +> `agent.onlineOutbox`) — поток стриминга (`Working`/`End`/`AppendText`/ +> `AppendImage`), **никогда не сохраняется** и не реплеится: потерянный при +> обрыве фрагмент невосстановим, но целый ответ всегда придёт durable- +> `Event.AssistantMessage` и/или ляжет в journal. `HttpClient`, `applyAgentikDefaults`, выбор engine'а — всё скрыто внутри `AgentikAgent`. Один вызов — один готовый `Agent`. @@ -172,9 +197,11 @@ history.forEach { rec -> ```kotlin import pw.binom.agentik.client.AgentikAgent -import pw.binom.agentik.proto.Content -import pw.binom.agentik.proto.Event +import pw.binom.agentik.content.Content +import pw.binom.agentik.outbox.Event +import pw.binom.agentik.outbox.OnlineEvent import io.ktor.client.engine.cio.CIO +import kotlinx.coroutines.launch val agent = AgentikAgent( id = "agentik", @@ -183,16 +210,30 @@ val agent = AgentikAgent( ) val conv = agent.createConversation(temp = false) -conv.send(listOf(Content.Text("Привет, расскажи про себя"))) -conv.events(after = kotlin.time.Clock.System.now()).collect { ev -> - when (ev) { - is Event.AppendText -> print(ev.body) // streaming чанки - is Event.End -> println("\n--- end ---") - is Event.Error -> error("agent error: ${ev.message}") - else -> Unit +// durable-поток (с курсором): terminal-события хода. +launch { + agent.outbox.conversationEvents(after = kotlin.time.Clock.System.now(), conversationId = conv.id) + .collect { ce -> + when (ce.event) { + is Event.AssistantMessage -> println("\n--- answer ready ---") + is Event.Error -> error("agent error: ${(ce.event as Event.Error).message}") + else -> Unit + } + } +} +// online-поток (live-only): стриминг ответа. +launch { + agent.onlineOutbox.onlineEvents(conv.id).collect { ev -> + when (ev) { + is OnlineEvent.AppendText -> print(ev.body) // streaming чанки + is OnlineEvent.End -> println("\n--- end ---") + else -> Unit + } } } + +conv.send(listOf(Content.Text("Привет, расскажи про себя"))) ``` ## История с локальным кэшем @@ -207,8 +248,8 @@ conv.events(after = kotlin.time.Clock.System.now()).collect { ev -> ```kotlin import pw.binom.agentik.journal.inmemory.InMemoryJournalStore -import pw.binom.agentik.proto.Content -import pw.binom.agentik.proto.Event +import pw.binom.agentik.content.Content +import pw.binom.agentik.outbox.Event import kotlin.time.Instant class ChatSession( @@ -235,10 +276,11 @@ class ChatSession( after = Instant.DISTANT_PAST, ).collect { cache.append(it) } } - // 2. Live: на каждом `End` хода просим у сервера новые записи. + // 2. Live: на каждом завершённом ходе (durable AssistantMessage) + // просим у сервера новые записи. scope.launch { - agent.getConversation(conversationId)!!.events(Instant.DISTANT_PAST).collect { ev -> - if (ev is Event.End) { + agent.outbox.conversationEvents(Instant.DISTANT_PAST, conversationId).collect { ce -> + if (ce.event is Event.AssistantMessage) { val newest = cache.let { // last-seen курсор — последний createdAt в кэше it.list(conversationId, Instant.DISTANT_PAST, 0, 1).lastOrNull()?.createdAt @@ -342,22 +384,24 @@ escape-hatch. Сам `InMemoryMutableConversationStore` тоже доступе появится в кэше через refresh-блок выше. ```kotlin -import pw.binom.agentik.proto.Event +import pw.binom.agentik.outbox.OnlineEvent -agent.getConversation(convId)!!.events(Instant.DISTANT_PAST).collect { ev -> +agent.onlineOutbox.onlineEvents(convId).collect { ev -> when (ev) { - is Event.StartResponse -> println("[start]") - is Event.AppendText -> print(ev.body) - is Event.AppendImage -> showImage(ev.body) - is Event.ToolCall -> println("[tool: ${ev.toolName}]") - is Event.ToolResult -> println("[result]") - is Event.End -> println("[end]") - is Event.Error -> println("[error: ${ev.message}]") - else -> Unit + is OnlineEvent.Working -> println("[working]") + is OnlineEvent.StartResponse -> println("[start]") + is OnlineEvent.AppendText -> print(ev.body) + is OnlineEvent.AppendImage -> showImage(ev.body) + is OnlineEvent.End -> println("[end]") + else -> Unit } } ``` +Инструментальные вызовы и целый ответ — durable-поток +(`agent.outbox.conversationEvents`): `Event.ToolCall`/`Event.ToolResult` и +`Event.AssistantMessage`/`Event.Interrupted`/`Event.Error`. + ## Прерывание хода ```kotlin diff --git a/client/src/commonMain/kotlin/pw/binom/agentik/client/ConversationClient.kt b/client/src/commonMain/kotlin/pw/binom/agentik/client/ConversationClient.kt index 339c092..1f05098 100644 --- a/client/src/commonMain/kotlin/pw/binom/agentik/client/ConversationClient.kt +++ b/client/src/commonMain/kotlin/pw/binom/agentik/client/ConversationClient.kt @@ -10,10 +10,10 @@ import io.ktor.client.request.setBody import io.ktor.http.ContentType import io.ktor.http.contentType import kotlinx.serialization.Serializable -import pw.binom.agentik.proto.Content +import pw.binom.agentik.content.Content import pw.binom.agentik.proto.Conversation import pw.binom.agentik.proto.Message -import pw.binom.agentik.proto.MessageContext +import pw.binom.agentik.content.MessageContext import kotlin.time.Instant /** diff --git a/client/src/commonMain/kotlin/pw/binom/agentik/client/HttpOnlineOutbox.kt b/client/src/commonMain/kotlin/pw/binom/agentik/client/HttpOnlineOutbox.kt index e260f25..2741cbd 100644 --- a/client/src/commonMain/kotlin/pw/binom/agentik/client/HttpOnlineOutbox.kt +++ b/client/src/commonMain/kotlin/pw/binom/agentik/client/HttpOnlineOutbox.kt @@ -13,11 +13,16 @@ import pw.binom.agentik.outbox.OnlineOutbox * HTTP-реализация [OnlineOutbox] (= [pw.binom.agentik.outbox.OnlineOutbox]), * ходящая в `:server`-фасад. * - * **Endpoint**: [onlineEvents] → `GET {baseUrl}/conversations/{id}/online` - * (live-only SSE, без `after`). Сервер не реплеит — подписка получает - * только то, что эмитится после подключения. Reconnect-логика здесь не - * нужна: при обрыве поток просто закрывается, а потерянные дельты - * восстанавливаются из durable-истории ([HttpEventStore] + журнал). + * **Endpoints**: + * - [onlineEvents] (без аргумента) → `GET {baseUrl}/online` (live-only SSE, + * все диалоги агента); + * - [onlineEvents] с `conversationId` → `GET {baseUrl}/conversations/{id}/online` + * (live-only SSE, один диалог). + * + * Сервер не реплеит — подписка получает только то, что эмитится после + * подключения. Reconnect-логика здесь не нужна: при обрыве поток просто + * закрывается, а потерянные дельты восстанавливаются из durable-истории + * ([HttpEventStore] + журнал). */ internal class HttpOnlineOutbox( private val httpClient: HttpClient, @@ -26,8 +31,12 @@ internal class HttpOnlineOutbox( private val agentUrl: String = baseUrl.trimEnd('/') - override fun onlineEvents(conversationId: String): Flow = flow { - val url = "$agentUrl/conversations/$conversationId/online" + override fun onlineEvents(): Flow = stream("$agentUrl/online") + + override fun onlineEvents(conversationId: String): Flow = + stream("$agentUrl/conversations/$conversationId/online") + + private fun stream(url: String): Flow = flow { httpClient.prepareGet(url) { noReadTimeout() } .execute { response -> check(response.status == HttpStatusCode.OK) { diff --git a/client/src/commonTest/kotlin/pw/binom/agentik/client/ReconnectingOutboxTest.kt b/client/src/commonTest/kotlin/pw/binom/agentik/client/ReconnectingOutboxTest.kt index b0ed8f9..5980f5d 100644 --- a/client/src/commonTest/kotlin/pw/binom/agentik/client/ReconnectingOutboxTest.kt +++ b/client/src/commonTest/kotlin/pw/binom/agentik/client/ReconnectingOutboxTest.kt @@ -66,7 +66,7 @@ private fun testEvent(dateMs: Long): CommonEvent = CommonEvent.Conversation( date = Instant.fromEpochMilliseconds(dateMs), conversationId = "test", - event = Event.End(date = Instant.fromEpochMilliseconds(dateMs)), + event = Event.Interrupted(date = Instant.fromEpochMilliseconds(dateMs)), ) @OptIn(ExperimentalCoroutinesApi::class) diff --git a/content-api/README.md b/content-api/README.md new file mode 100644 index 0000000..b46707e --- /dev/null +++ b/content-api/README.md @@ -0,0 +1,33 @@ +# `:content-api` — общие типы содержимого сообщения + +Низкоуровневый KMP-модуль с типами, которые используются во всех слоях +agentik и раньше дублировались: + +- `Content` — часть содержимого сообщения: `Content.Text(body)`, + `Content.Image(data, mime)`. +- `MessageContext` — контекст инициации хода (`origin`, `description`, + `sourceId`, `metadata`). +- `MessageOrigin` — `USER` / `SYSTEM` / `EVENT`. +- `TurnTokens` — token usage одного assistant turn'а (`input`, `output`). + +## Зачем отдельный модуль + +`:proto` (wire-контракт), `:journal-api` (слой хранения) и `:outbox-api` +(события) должны ссылаться на **один и тот же** `Content`/`MessageContext`, +а не держать по собственной копии. Общий модуль убирает дубли и циклы: + +``` +:content-api ◄── :proto + ◄── :journal-api + ◄── :outbox-api +``` + +`:proto`/`:journal-api`/`:outbox-api` объявляют `api(project(":content-api"))`, +поэтому потребители (`:client`, `:server`, `:standalone`, ...) видят типы +транзитивно, но должны импортировать их напрямую из `pw.binom.agentik.content`. + +## Публикация + +Каталог `gradle/libs.versions.toml` → `agentik-content-api`. +`./gradlew :content-api:publish -Pversion=...` публикует все KMP-таргеты +(jvm + натив). diff --git a/content-api/build.gradle.kts b/content-api/build.gradle.kts new file mode 100644 index 0000000..82c52a1 --- /dev/null +++ b/content-api/build.gradle.kts @@ -0,0 +1,30 @@ +plugins { + alias(libs.plugins.kotlin.multiplatform) + alias(libs.plugins.kotlin.serialization) +} + +kotlin { + jvmToolchain(21) + jvm() + macosX64() + macosArm64() + iosX64() + iosArm64() + iosSimulatorArm64() + linuxX64() + linuxArm64() + mingwX64() + + sourceSets { + commonMain.dependencies { + api(libs.kotlinx.coroutines.core) + api(libs.kotlinx.serialization.core) + // для JsonElement в MessageContext.metadata + api(libs.kotlinx.serialization.json) + } + commonTest.dependencies { + implementation(kotlin("test")) + implementation(libs.kotlinx.coroutines.test) + } + } +} diff --git a/content-api/src/commonMain/kotlin/pw/binom/agentik/content/Content.kt b/content-api/src/commonMain/kotlin/pw/binom/agentik/content/Content.kt new file mode 100644 index 0000000..053c512 --- /dev/null +++ b/content-api/src/commonMain/kotlin/pw/binom/agentik/content/Content.kt @@ -0,0 +1,31 @@ +package pw.binom.agentik.content + +import kotlinx.serialization.SerialName +import kotlinx.serialization.Serializable + +/** + * Часть содержимого сообщения (пользовательского или агентского). + * + * Единый тип для всего проекта: используется и в wire-контракте ([pw.binom.agentik.proto]), + * и в слое хранения ([pw.binom.agentik.journal]), и в durable-событиях + * ([pw.binom.agentik.outbox.Event]). Вынесен в отдельный модуль `:content-api`, + * чтобы не дублировать его в каждом слое и не заводить циклов в графе. + */ +@Serializable +sealed interface Content { + @Serializable + @SerialName("text") + data class Text(val body: String) : Content + + /** + * Картинка. [mime] — MIME-тип, например `"image/png"`, `"image/jpeg"`. + */ + @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() + } +} diff --git a/content-api/src/commonMain/kotlin/pw/binom/agentik/content/MessageContext.kt b/content-api/src/commonMain/kotlin/pw/binom/agentik/content/MessageContext.kt new file mode 100644 index 0000000..653f127 --- /dev/null +++ b/content-api/src/commonMain/kotlin/pw/binom/agentik/content/MessageContext.kt @@ -0,0 +1,52 @@ +package pw.binom.agentik.content + +import kotlinx.serialization.Serializable +import kotlinx.serialization.json.JsonElement + +/** + * Контекст инициации сообщения: кто/что и почему вызвало этот ход. + * + * Примеры: + * ``` + * // cron-задача утренней сводки + * MessageContext( + * origin = MessageOrigin.EVENT, + * description = "scheduled cron 'morning-briefing'", + * sourceId = "cron-42", + * metadata = buildJsonObject { put("scheduledAt", "2026-09-14T08:00:00Z") }, + * ) + * + * // обычное сообщение из IRC + * MessageContext( + * origin = MessageOrigin.USER, + * sourceId = "irc-channel:agentik", + * description = "PRIVMSG from nick", + * ) + * ``` + * + * Семантический контракт: + * - origin != USER ⇒ [description] обязателен и должен быть человекочитаемым. + * - origin == USER ⇒ context может быть `null` (дефолт). + * + * Снапшот-стабильность wire-формата: поля сериализуются по именам, snake_case + * на enum'е [MessageOrigin] даёт `user`/`system`/`event`. Новые поля — + * non-breaking для старых клиентов. + */ +@Serializable +data class MessageContext( + val origin: MessageOrigin, + /** + * Короткая человекочитаемая фраза для LLM: попадает в working memory + * как префикс `[origin] description (sourceId=…)` к user-сообщению. + */ + val description: String? = null, + /** + * Идентификатор источника: id cron-job'а, webhook endpoint'а, имя канала IRC. + */ + val sourceId: String? = null, + /** + * Произвольный структурированный payload о событии. Никогда не попадает + * в LLM-нагрузку как сырой JSON — только логирование и пост-аналитика. + */ + val metadata: JsonElement? = null, +) diff --git a/content-api/src/commonMain/kotlin/pw/binom/agentik/content/MessageOrigin.kt b/content-api/src/commonMain/kotlin/pw/binom/agentik/content/MessageOrigin.kt new file mode 100644 index 0000000..3fae897 --- /dev/null +++ b/content-api/src/commonMain/kotlin/pw/binom/agentik/content/MessageOrigin.kt @@ -0,0 +1,19 @@ +package pw.binom.agentik.content + +import kotlinx.serialization.SerialName +import kotlinx.serialization.Serializable + +/** + * Кто/что инициировал ход (кто/что и почему). + */ +@Serializable +enum class MessageOrigin { + @SerialName("user") + USER, + + @SerialName("system") + SYSTEM, + + @SerialName("event") + EVENT, +} diff --git a/journal-api/src/commonMain/kotlin/pw/binom/agentik/journal/TurnTokens.kt b/content-api/src/commonMain/kotlin/pw/binom/agentik/content/TurnTokens.kt similarity index 92% rename from journal-api/src/commonMain/kotlin/pw/binom/agentik/journal/TurnTokens.kt rename to content-api/src/commonMain/kotlin/pw/binom/agentik/content/TurnTokens.kt index 8c98496..4179e6a 100644 --- a/journal-api/src/commonMain/kotlin/pw/binom/agentik/journal/TurnTokens.kt +++ b/content-api/src/commonMain/kotlin/pw/binom/agentik/content/TurnTokens.kt @@ -1,4 +1,4 @@ -package pw.binom.agentik.journal +package pw.binom.agentik.content import kotlinx.serialization.Serializable @@ -11,6 +11,7 @@ data class TurnTokens( 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" } diff --git a/proto/src/commonTest/kotlin/pw/binom/agentik/proto/MessageContextTest.kt b/content-api/src/commonTest/kotlin/pw/binom/agentik/content/MessageContextTest.kt similarity index 99% rename from proto/src/commonTest/kotlin/pw/binom/agentik/proto/MessageContextTest.kt rename to content-api/src/commonTest/kotlin/pw/binom/agentik/content/MessageContextTest.kt index 5aec7db..fd35705 100644 --- a/proto/src/commonTest/kotlin/pw/binom/agentik/proto/MessageContextTest.kt +++ b/content-api/src/commonTest/kotlin/pw/binom/agentik/content/MessageContextTest.kt @@ -1,4 +1,4 @@ -package pw.binom.agentik.proto +package pw.binom.agentik.content import kotlinx.serialization.json.Json import kotlinx.serialization.json.JsonPrimitive diff --git a/context-api/src/commonMain/kotlin/pw/binom/agentik/context/WorkingMemoryEntry.kt b/context-api/src/commonMain/kotlin/pw/binom/agentik/context/WorkingMemoryEntry.kt index 4b66b5c..20d84ef 100644 --- a/context-api/src/commonMain/kotlin/pw/binom/agentik/context/WorkingMemoryEntry.kt +++ b/context-api/src/commonMain/kotlin/pw/binom/agentik/context/WorkingMemoryEntry.kt @@ -2,8 +2,8 @@ package pw.binom.agentik.context import kotlinx.serialization.SerialName import kotlinx.serialization.Serializable -import pw.binom.agentik.journal.Content -import pw.binom.agentik.journal.MessageContext +import pw.binom.agentik.content.Content +import pw.binom.agentik.content.MessageContext /** * Запись в working memory диалога: ровно то, что агент сейчас видит в 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 460d1e5..e6b6bbb 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.journal.Content +import pw.binom.agentik.content.Content import pw.binom.db.ksqlite.SQLiteConnection import kotlin.test.AfterTest import kotlin.test.BeforeTest diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index ed33485..5b9cb2f 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -47,28 +47,35 @@ agentik ## 3. `:proto` — интерфейсы `Agent` (см. `proto/src/commonMain/.../Agent.kt`): -- `id: String` +- `id: String`, `info: AgentInfo` +- read-only сторы: `journal: JournalStore`, `outbox: OutboxStore`, + `onlineOutbox: OnlineOutbox`, `conversationStore: ConversationStore` - `createConversation(temp: Boolean): Conversation` - `suspend getConversation(id): Conversation?` -- `suspend getConversations(offset, limit): List` -- `getConversations(offset = 0): Flow` — cold-flow paging через suspend-версию, `PAGE_SIZE = 100`. -- `events(after: Instant): Flow` — replay-free, бэкфилл через snapshot. -- `deleteConversation(id): Boolean` +- `suspend deleteConversation(id): Boolean` +- `suspend renameConversation(id, title): Instant?` `Conversation`: -- `isSupportImageInput / Output / isTemporal: Boolean` +- `isSupportImageInput / Output / isTemporal: Boolean`, `title: String?` - `updatedAt: Instant` -- `send(content: List)` — write-only, ничего не возвращает. +- `send(content: List, context: MessageContext? = null)` — write-only, ничего не возвращает. - `interrupt()` — отмена активного хода. -- `events(after): Flow` — live, replay-free. -- `getMessages(after, offset, limit)` + `getMessages(after): Flow` — paging. +- `getMessages(after, offset, limit): List` — paging. - `rename(title)` — мутация, бампит `updatedAt`. - `AutoCloseable` — `close()` идемпотентен. -`Content = Text(body) | Image(data, mime)`. +`Content = Text(body) | Image(data, mime)` — из `:content-api` (там же `MessageContext`/`MessageOrigin`/`TurnTokens`). `Message = UserMessage | AssistantMessage | ToolCall | ToolResult | Error`. -`Event = StartReasoning | StartResponse | End | AppendText | AppendImage | ToolCall | ToolResult | Error`. -`AgentEvent = Created(conversationId) | Deleted(id) | Renamed(id, title)`. + +События разделены на два потока (оба в `:outbox-api`): +- **durable** `Event` (`outbox.conversationEvents(after, id)`): `UserMessage | + AssistantMessage | ToolCall | ToolResult | ToolFailed | Interrupted | Error | + ConversationClosing | CompactionTriggered` — перезапрашиваются по курсору; +- **online** `OnlineEvent` (`onlineOutbox.onlineEvents(id)`): `Working | End | + StartReasoning | StartResponse | AppendText | AppendImage` — live-only, + никогда не сохраняются. + +`AgentEvent = Created(conversationId) | Deleted(id) | Renamed(id, title) | Touched(...)`. Принцип: **агент — источник истины** для транскрипта и сессий. Клиент лишь рендерит Event-stream и кэширует историю. diff --git a/docs/STANDALONE.md b/docs/STANDALONE.md index 87bcdc3..a237ba4 100644 --- a/docs/STANDALONE.md +++ b/docs/STANDALONE.md @@ -63,9 +63,9 @@ │ 3. ensureLiteConversation: │ │ first turn → create from WM; │ │ next turns → reuse (KV-cache) │ - │ 4. sendStreamContents → emit │ - │ StartResponse / AppendText / │ - │ End │ + │ 4. sendStreamContents → online: │ + │ StartResponse/AppendText/End; │ + │ durable: AssistantMessage │ │ 5. audit + WM: append AssistantMessage│ └──────────────────────────────────────┘ │ │ @@ -151,7 +151,8 @@ fun main() { | `updatedAt: Instant` | последний `send`/`rename` | | `send(content: List)` | **fire-and-forget**: добавить user-сообщение, запустить ход, выйти | | `interrupt()` | остановить текущий ход (best-effort) | -| `events(after): Flow` | live-события хода (StartReasoning, StartResponse, AppendText, End, Interrupted, Error) | +| `outbox.conversationEvents(after, id): Flow` | durable-события (UserMessage, AssistantMessage, ToolCall/Result, Interrupted, Error) — скурсором | +| `onlineOutbox.onlineEvents(id): Flow` | live-only стриминг (Working, End, StartReasoning, StartResponse, AppendText, AppendImage) | | `getMessages(after, offset, limit)` | страница истории | | `rename(title)` | переименовать | | `close()` | освободить ресурсы | @@ -173,7 +174,7 @@ fun main() { Подробный контракт — в комментариях к `MessageRecord.kt` и `WorkingMemoryEntry.kt`. -**Ошибки хода персистятся.** Если ход провалился (LLM/движок недоступны — например, HTTP 400 от endpoint'а), `ChatConversation.failTurn` пишет терминальную запись `MessageRecord.Error` в audit и эмитит `Event.Error` + `Event.End`. Благодаря audit-записи ошибка видна не только подписчику live-SSE, но и клиенту, который делает backfill через `getMessages` (polling/переподключение): в истории будет `Message.Error(id, message, code?)`, а для этого user-сообщения не будет `AssistantMessage`. При ошибке стрима живой `LiteConversation` сбрасывается — следующий `send` пересоберёт его из `working_memory`. В working_memory `Error` не пишется (модель не должна видеть ошибки прошлых ходов). +**Ошибки хода персистятся.** Если ход провалился (LLM/движок недоступны — например, HTTP 400 от endpoint'а), `ChatConversation.failTurn` пишет терминальную запись `MessageRecord.Error` в audit и эмитит durable `Event.Error` + онлайн `OnlineEvent.End`. Благодаря audit-записи ошибка видна не только подписчику live-SSE, но и клиенту, который делает backfill через `getMessages` (polling/переподключение): в истории будет `Message.Error(id, message, code?)`, а для этого user-сообщения не будет `AssistantMessage`. При ошибке стрима живой `LiteConversation` сбрасывается — следующий `send` пересоберёт его из `working_memory`. В working_memory `Error` не пишется (модель не должна видеть ошибки прошлых ходов). ### `JournalStore` diff --git a/journal-api/build.gradle.kts b/journal-api/build.gradle.kts index 7319039..3ac2371 100644 --- a/journal-api/build.gradle.kts +++ b/journal-api/build.gradle.kts @@ -17,6 +17,7 @@ kotlin { sourceSets { commonMain.dependencies { + api(project(":content-api")) api(libs.kotlinx.coroutines.core) api(libs.kotlinx.serialization.core) api(libs.kotlinx.serialization.json) diff --git a/journal-api/src/commonMain/kotlin/pw/binom/agentik/journal/Content.kt b/journal-api/src/commonMain/kotlin/pw/binom/agentik/journal/Content.kt deleted file mode 100644 index 5771606..0000000 --- a/journal-api/src/commonMain/kotlin/pw/binom/agentik/journal/Content.kt +++ /dev/null @@ -1,24 +0,0 @@ -package pw.binom.agentik.journal - -import kotlinx.serialization.SerialName -import kotlinx.serialization.Serializable - -/** - * Часть контента сообщения на уровне хранилища. Намеренно НЕ зависит от - * `pw.binom.agentik.proto.Content` — маппинг `:proto.Content ↔ Content` живёт - * в `Mapping.kt` storage impl'ов. - */ -@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() - } -} diff --git a/journal-api/src/commonMain/kotlin/pw/binom/agentik/journal/MessageContext.kt b/journal-api/src/commonMain/kotlin/pw/binom/agentik/journal/MessageContext.kt deleted file mode 100644 index b8e14bc..0000000 --- a/journal-api/src/commonMain/kotlin/pw/binom/agentik/journal/MessageContext.kt +++ /dev/null @@ -1,12 +0,0 @@ -package pw.binom.agentik.journal - -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, -) diff --git a/journal-api/src/commonMain/kotlin/pw/binom/agentik/journal/MessageOrigin.kt b/journal-api/src/commonMain/kotlin/pw/binom/agentik/journal/MessageOrigin.kt deleted file mode 100644 index 5f16d9f..0000000 --- a/journal-api/src/commonMain/kotlin/pw/binom/agentik/journal/MessageOrigin.kt +++ /dev/null @@ -1,20 +0,0 @@ -package pw.binom.agentik.journal - -import kotlinx.serialization.SerialName -import kotlinx.serialization.Serializable - -/** - * Контекст инициации хода (кто/что и почему). Дубликат типа из `:proto` — - * живёт здесь чтобы не тащить `:proto` в слой хранения данных. - */ -@Serializable -enum class MessageOrigin { - @SerialName("user") - USER, - - @SerialName("system") - SYSTEM, - - @SerialName("event") - EVENT, -} diff --git a/journal-api/src/commonMain/kotlin/pw/binom/agentik/journal/MessageRecord.kt b/journal-api/src/commonMain/kotlin/pw/binom/agentik/journal/MessageRecord.kt index 0828f8e..ea184a3 100644 --- a/journal-api/src/commonMain/kotlin/pw/binom/agentik/journal/MessageRecord.kt +++ b/journal-api/src/commonMain/kotlin/pw/binom/agentik/journal/MessageRecord.kt @@ -2,6 +2,9 @@ package pw.binom.agentik.journal import kotlinx.serialization.SerialName import kotlinx.serialization.Serializable +import pw.binom.agentik.content.Content +import pw.binom.agentik.content.MessageContext +import pw.binom.agentik.content.TurnTokens import kotlin.time.Instant /** @@ -36,6 +39,12 @@ sealed interface MessageRecord { override val content: List, override val createdAt: Instant, val tokens: TurnTokens? = null, + /** + * Текст размышлений модели (chain-of-thought / reasoning), если провайдер + * его отдаёт. Опционально: null, если размышлений не было или провайдер + * их не раскрывает. + */ + val reasoning: String? = null, ) : Body @Serializable diff --git a/journal-api/src/commonMain/kotlin/pw/binom/agentik/journal/Payload.kt b/journal-api/src/commonMain/kotlin/pw/binom/agentik/journal/Payload.kt index 8d14b17..52937b6 100644 --- a/journal-api/src/commonMain/kotlin/pw/binom/agentik/journal/Payload.kt +++ b/journal-api/src/commonMain/kotlin/pw/binom/agentik/journal/Payload.kt @@ -4,6 +4,9 @@ import kotlinx.serialization.SerialName import kotlinx.serialization.Serializable import kotlinx.serialization.builtins.ListSerializer import kotlinx.serialization.json.Json +import pw.binom.agentik.content.Content +import pw.binom.agentik.content.MessageContext +import pw.binom.agentik.content.TurnTokens private val bodyJson = Json { ignoreUnknownKeys = true @@ -17,15 +20,17 @@ data class MessageBodyPayload( @SerialName("context") val context: MessageContext? = null, val tokens: TurnTokens? = null, + val reasoning: String? = null, ) fun encodeBodyPayload( content: List, context: MessageContext? = null, tokens: TurnTokens? = null, + reasoning: String? = null, ): String = bodyJson.encodeToString( MessageBodyPayload.serializer(), - MessageBodyPayload(content = content, context = context, tokens = tokens), + MessageBodyPayload(content = content, context = context, tokens = tokens, reasoning = reasoning), ) fun decodeBodyPayload(json: String): BodyDecoded = readPayload(json) @@ -34,14 +39,15 @@ data class BodyDecoded( val content: List, val context: MessageContext?, val tokens: TurnTokens? = null, + val reasoning: String? = null, ) private fun readPayload(json: String): BodyDecoded { return try { val p = bodyJson.decodeFromString(MessageBodyPayload.serializer(), json) - BodyDecoded(p.content, p.context, p.tokens) + BodyDecoded(p.content, p.context, p.tokens, p.reasoning) } catch (e: kotlinx.serialization.SerializationException) { val arr = bodyJson.decodeFromString(ListSerializer(Content.serializer()), json) - BodyDecoded(arr, null, null) + BodyDecoded(arr, null, null, null) } } diff --git a/journal-inmemory/src/commonTest/kotlin/pw/binom/agentik/journal/inmemory/InMemoryJournalStoreTest.kt b/journal-inmemory/src/commonTest/kotlin/pw/binom/agentik/journal/inmemory/InMemoryJournalStoreTest.kt index b30e7aa..1b1d607 100644 --- a/journal-inmemory/src/commonTest/kotlin/pw/binom/agentik/journal/inmemory/InMemoryJournalStoreTest.kt +++ b/journal-inmemory/src/commonTest/kotlin/pw/binom/agentik/journal/inmemory/InMemoryJournalStoreTest.kt @@ -14,7 +14,7 @@ class InMemoryJournalStoreTest { MessageRecord.UserMessage( id = id, conversationId = convId, - content = listOf(pw.binom.agentik.journal.Content.Text(text)), + content = listOf(pw.binom.agentik.content.Content.Text(text)), createdAt = at, ) diff --git a/journal-ksqlite/src/commonMain/kotlin/pw/binom/agentik/journal/ksqlite/MessageCodecs.kt b/journal-ksqlite/src/commonMain/kotlin/pw/binom/agentik/journal/ksqlite/MessageCodecs.kt index ad0e5f5..096d3c1 100644 --- a/journal-ksqlite/src/commonMain/kotlin/pw/binom/agentik/journal/ksqlite/MessageCodecs.kt +++ b/journal-ksqlite/src/commonMain/kotlin/pw/binom/agentik/journal/ksqlite/MessageCodecs.kt @@ -21,6 +21,7 @@ internal fun encodeRecord(record: MessageRecord): Pair = when (r is MessageRecord.AssistantMessage -> "assistant" to encodeBodyPayload( content = record.content, tokens = record.tokens, + reasoning = record.reasoning, ) is MessageRecord.ToolCall -> "tool_call" to Json.encodeToString( CallPayload.serializer(), @@ -49,7 +50,7 @@ internal fun SQLiteResultSet.toMessageRecord(json: Json): MessageRecord { } "assistant" -> { val d = decodeBodyPayload(payload) - MessageRecord.AssistantMessage(id = id, conversationId = convId, content = d.content, createdAt = createdAt, tokens = d.tokens) + MessageRecord.AssistantMessage(id = id, conversationId = convId, content = d.content, createdAt = createdAt, tokens = d.tokens, reasoning = d.reasoning) } "tool_call" -> { val p = Json.decodeFromString(CallPayload.serializer(), payload) diff --git a/journal-ksqlite/src/commonTest/kotlin/pw/binom/agentik/journal/ksqlite/KsqliteJournalStoreTest.kt b/journal-ksqlite/src/commonTest/kotlin/pw/binom/agentik/journal/ksqlite/KsqliteJournalStoreTest.kt index 9681926..e7f7d69 100644 --- a/journal-ksqlite/src/commonTest/kotlin/pw/binom/agentik/journal/ksqlite/KsqliteJournalStoreTest.kt +++ b/journal-ksqlite/src/commonTest/kotlin/pw/binom/agentik/journal/ksqlite/KsqliteJournalStoreTest.kt @@ -2,9 +2,9 @@ package pw.binom.agentik.journal.ksqlite import kotlinx.coroutines.flow.toList import kotlinx.coroutines.test.runTest -import pw.binom.agentik.journal.Content +import pw.binom.agentik.content.Content import pw.binom.agentik.journal.MessageRecord -import pw.binom.agentik.journal.TurnTokens +import pw.binom.agentik.content.TurnTokens import pw.binom.db.ksqlite.SQLiteConnection import kotlin.test.AfterTest import kotlin.test.BeforeTest diff --git a/journal-ksqlite/src/commonTest/kotlin/pw/binom/agentik/journal/ksqlite/SchemaMigrationTest.kt b/journal-ksqlite/src/commonTest/kotlin/pw/binom/agentik/journal/ksqlite/SchemaMigrationTest.kt index 96648eb..55c6615 100644 --- a/journal-ksqlite/src/commonTest/kotlin/pw/binom/agentik/journal/ksqlite/SchemaMigrationTest.kt +++ b/journal-ksqlite/src/commonTest/kotlin/pw/binom/agentik/journal/ksqlite/SchemaMigrationTest.kt @@ -1,7 +1,7 @@ package pw.binom.agentik.journal.ksqlite import kotlinx.coroutines.test.runTest -import pw.binom.agentik.journal.Content +import pw.binom.agentik.content.Content import pw.binom.agentik.journal.ConversationRecord import pw.binom.agentik.journal.MessageRecord import pw.binom.db.ksqlite.SQLiteConnection diff --git a/outbox-api/build.gradle.kts b/outbox-api/build.gradle.kts index d44cff0..89d3690 100644 --- a/outbox-api/build.gradle.kts +++ b/outbox-api/build.gradle.kts @@ -20,8 +20,10 @@ kotlin { sourceSets { commonMain.dependencies { - // :proto больше не нужен — AgentEvent/CommonEvent/Event перенесены - // сюда, и они self-contained (Event ссылается только на kotlinx-serialization). + // :proto не нужен (Event/OnlineEvent/AgentEvent/CommonEvent живут + // здесь). Message-события несут общие типы содержимого из + // низкоуровневого :content-api (Content/MessageContext/TurnTokens). + api(project(":content-api")) api(libs.kotlinx.coroutines.core) api(libs.kotlinx.serialization.core) api(libs.kotlinx.serialization.json) diff --git a/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/Event.kt b/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/Event.kt index e265260..4bd4b8d 100644 --- a/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/Event.kt +++ b/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/Event.kt @@ -2,29 +2,25 @@ package pw.binom.agentik.outbox import kotlinx.serialization.SerialName import kotlinx.serialization.Serializable +import pw.binom.agentik.content.Content +import pw.binom.agentik.content.MessageContext +import pw.binom.agentik.content.TurnTokens import kotlin.time.Instant /** - * Элемент live-потока `Conversation.events(after)`. + * Элемент **durable**-потока диалога. * * Каждое событие несёт [date] — момент эмиссии в UTC. Используется клиентом - * для трекинга «где остановился» при обрыве/переподключении и для разрешения + * как курсор («где остановился») при обрыве/переподключении и для разрешения * порядка при равных timestamps. * - * Базовая структура хода: - * `Working` → `End` | `Interrupted` | `Error`, между ними — целые события - * ([ToolCall]/[ToolResult]/[ToolFailed]). + * Это «целые», сохраняемые события: сообщения ([UserMessage]/[AssistantMessage]), + * вызовы тулов ([ToolCall]/[ToolResult]/[ToolFailed]), терминаторы хода + * ([Interrupted]/[Error]) и lifecycle ([ConversationClosing]/[CompactionTriggered]). + * Их можно перезапросить по курсору (`after`). * - * **Стриминг ответа (дельты текста/картинок и маркеры фаз) — НЕ здесь.** - * Дельты токенов живут в [OnlineEvent] (live-only, не сохраняются). - * [Event] — только «целые» (durable) события, пригодные к перезапросу - * по курсору. - * - * `Working` — маркер «агент принял запрос и пошёл обрабатывать», эмитится - * синхронно в `Conversation.send()` ДО старта LLM-цикла (и возможной долгой - * очереди turnLock'а). Парный терминатор не нужен: `End`/`Interrupted`/`Error` - * уже закрывают ход. UI использует `Working` чтобы показать спиннер ещё до - * первого токена ответа. + * **Стриминг ответа и маркеры фаз хода — НЕ здесь.** Дельты текста/картинок + * и маркеры `Working`/`End` живут в [OnlineEvent] (live-only, не сохраняются). * * **История**: до 2026-09-21 жил в `:proto` как `pw.binom.agentik.proto.Event`; * при миграции в `:outbox-api` был оставлен typealias в `:proto` для @@ -39,24 +35,35 @@ sealed interface Event { val date: Instant /** - * Маркер «агент принял запрос и пошёл обрабатывать». Эмитится **до** - * [End]/[Interrupted]/[Error], синхронно из `Conversation.send()`, - * чтобы UI мог показать спиннер ещё до первого токена ответа. - * Терминатор хода ([End]/[Interrupted]/[Error]) — парный. + * Целое пользовательское сообщение хода. Эмитится в тот же момент, когда + * запись попадает в журнал (для персистентных диалогов). [id] совпадает + * с id соответствующего `MessageRecord.UserMessage`/`Message.UserMessage`. */ @Serializable - @SerialName("working") - data class Working(override val date: Instant) : Event + @SerialName("user_message") + data class UserMessage( + override val date: Instant, + val id: String, + val content: List, + val context: MessageContext? = null, + ) : Event - /** Ход завершён нормально. Соответствующий `Message.AssistantMessage` появится в `getMessages`. */ + /** + * Целое сообщение ассистента — итог хода. Эмитится при завершении хода, + * после того как текст ответа полностью собран. + * + * [reasoning] — текст размышлений модели (chain-of-thought), если провайдер + * его отдаёт; иначе `null`. [tokens] — расход токенов за ход. + */ @Serializable - @SerialName("end") - data class End(override val date: Instant) : Event - - /** Ход прерван через `Conversation.interrupt`. Частичный ответ НЕ сохраняется в истории. */ - @Serializable - @SerialName("interrupted") - data class Interrupted(override val date: Instant) : Event + @SerialName("assistant_message") + data class AssistantMessage( + override val date: Instant, + val id: String, + val content: List, + val reasoning: String? = null, + val tokens: TurnTokens? = null, + ) : Event /** * Агент начал вызов тула. Аргументы приходят целиком — стриминга нет. @@ -96,6 +103,14 @@ sealed interface Event { val result: String?, ) : Event + /** + * Ход прерван через `Conversation.interrupt`. Частичный ответ НЕ сохраняется + * в истории. + */ + @Serializable + @SerialName("interrupted") + data class Interrupted(override val date: Instant) : Event + /** * Ошибка хода. После неё поток завершается; дальнейшие события могут * прийти, но ход считается проваленным. diff --git a/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/OnlineEvent.kt b/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/OnlineEvent.kt index 0ef7f8f..8358fec 100644 --- a/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/OnlineEvent.kt +++ b/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/OnlineEvent.kt @@ -5,8 +5,8 @@ import kotlinx.serialization.Serializable import kotlin.time.Instant /** - * **Онлайн-события** диалога: стриминг ответа агента «в моменте» — - * дельты текста/картинок и маркеры фаз хода. + * **Онлайн-события** диалога: live-поток «в моменте» — маркеры фаз хода и + * стриминг ответа агента (дельты текста/картинок). * * Принципиальное отличие от [Event] (durable): * - **Никогда и нигде не сохраняются** — ни в буфер [OnlineOutbox], @@ -14,13 +14,13 @@ import kotlin.time.Instant * - **Только онлайн-подписка**: события, эмитнутые до подписки * (или в момент обрыва соединения), не реплеятся и не восстанавливаются. * Потерянный фрагмент не страшен — целый результат хода приходит - * durable-событием ([Event.End]) и/или лежит в журнале. + * durable-событием ([Event.AssistantMessage]) и/или лежит в журнале. * - **Нет курсора**: у потока нет `after`/`lastSeen` — курсор там, где * есть что реплеить. * - * Зачем разделять: дельты токенов — высокочастотный мусор, который, - * попав в durable store, копится в RAM (standalone-outbox растёт unbounded) - * и засоряет историю. В [Event] остаются только «целые» события, + * Зачем разделять: маркеры фаз и дельты токенов — высокочастотный мусор, + * который, попав в durable store, копится в RAM (standalone-outbox растёт + * unbounded) и засоряет историю. В [Event] остаются только «целые» события, * пригодные к перезапросу по курсору. */ @Serializable @@ -34,6 +34,20 @@ sealed interface OnlineEvent { @SerialName("image") IMAGE } + /** + * Маркер «агент принял запрос и пошёл обрабатывать». Эмитится **до** + * [End]/[Event.Interrupted]/[Event.Error], синхронно из `Conversation.send()`, + * чтобы UI мог показать спиннер ещё до первого токена ответа. + */ + @Serializable + @SerialName("working") + data class Working(override val date: Instant) : OnlineEvent + + /** Ход завершён (нормально либо оборван). Зеркало терминатора — см. [Event]. */ + @Serializable + @SerialName("end") + data class End(override val date: Instant) : OnlineEvent + /** Ассистент начал рассуждение (опциональный маркер; контент идёт через [AppendText]). */ @Serializable @SerialName("start_reasoning") diff --git a/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/OnlineOutbox.kt b/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/OnlineOutbox.kt index 3b813b3..127227d 100644 --- a/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/OnlineOutbox.kt +++ b/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/OnlineOutbox.kt @@ -14,7 +14,7 @@ import kotlinx.coroutines.flow.Flow * * Это осознанный компромисс: дельты токенов — высокочастотный мусор, * который в durable-сторе копился бы в RAM и засорял историю. Потеря - * фрагмента при обрыве не критична — целый ответ приходит [Event.End] + * фрагмента при обрыве не критична — целый ответ приходит [Event.AssistantMessage] * и/или лежит в [pw.binom.agentik.journal.JournalStore]. * * Read-only view: запись — через [MutableOnlineOutbox]. diff --git a/outbox-inmemory/src/commonMain/kotlin/pw/binom/agentik/outbox/inmemory/InMemoryOnlineOutbox.kt b/outbox-inmemory/src/commonMain/kotlin/pw/binom/agentik/outbox/inmemory/InMemoryOnlineOutbox.kt index 0d5953d..50d5ecb 100644 --- a/outbox-inmemory/src/commonMain/kotlin/pw/binom/agentik/outbox/inmemory/InMemoryOnlineOutbox.kt +++ b/outbox-inmemory/src/commonMain/kotlin/pw/binom/agentik/outbox/inmemory/InMemoryOnlineOutbox.kt @@ -60,6 +60,10 @@ class InMemoryOnlineOutbox( } private companion object { - private const val DEFAULT_LIVE_BUFFER_CAPACITY = 1 + // Достаточно, чтобы не терять подряд идущие маркеры фаз + // (Working/StartReasoning/StartResponse) и первые дельты, пока + // подписчик не возобновил сбор. Дельты токенов при переполнении + // всё равно дропаются (DROP_OLDEST) — потеря фрагмента допустима. + private const val DEFAULT_LIVE_BUFFER_CAPACITY = 64 } } diff --git a/outbox-inmemory/src/commonTest/kotlin/pw/binom/agentik/outbox/inmemory/InMemoryOutboxStoreTest.kt b/outbox-inmemory/src/commonTest/kotlin/pw/binom/agentik/outbox/inmemory/InMemoryOutboxStoreTest.kt index 9eb5f34..9cb2f0d 100644 --- a/outbox-inmemory/src/commonTest/kotlin/pw/binom/agentik/outbox/inmemory/InMemoryOutboxStoreTest.kt +++ b/outbox-inmemory/src/commonTest/kotlin/pw/binom/agentik/outbox/inmemory/InMemoryOutboxStoreTest.kt @@ -164,7 +164,7 @@ class InMemoryOutboxStoreTest { store.append(CommonEvent.Conversation( date = now, conversationId = "c-1", - event = Event.Working(date = now), + event = Event.Interrupted(date = now), )) // Snapshot-based test of the default impl (uses events() + filterIsInstance). @@ -180,9 +180,9 @@ class InMemoryOutboxStoreTest { fun `conversationEvents with conversationId filters to that conversation`() = runBlocking { val store = InMemoryOutboxStore(maxMessages = null, ttl = null) val now = Instant.fromEpochSeconds(0) - store.append(CommonEvent.Conversation(now, "c-1", Event.Working(now))) - store.append(CommonEvent.Conversation(now, "c-2", Event.Working(now))) - store.append(CommonEvent.Conversation(now, "c-1", Event.Working(now))) + store.append(CommonEvent.Conversation(now, "c-1", Event.Interrupted(now))) + store.append(CommonEvent.Conversation(now, "c-2", Event.Interrupted(now))) + store.append(CommonEvent.Conversation(now, "c-1", Event.Interrupted(now))) // Test the filter logic by manually filtering snapshot. val c1 = store.snapshot() @@ -200,7 +200,7 @@ class InMemoryOutboxStoreTest { date = now, event = AgentEvent.Created(date = now, conversationId = "created"), )) - store.append(CommonEvent.Conversation(now, "c-1", Event.Working(now))) + store.append(CommonEvent.Conversation(now, "c-1", Event.Interrupted(now))) val all = store.snapshot() val agents = all.filterIsInstance() diff --git a/proto/README.md b/proto/README.md index 909ed26..58611a1 100644 --- a/proto/README.md +++ b/proto/README.md @@ -59,57 +59,87 @@ target-specific артефакты + общий `kotlinMultiplatform`. ## Основные типы +`:proto` — **тонкий** контракт: сами интерфейсы `Agent`/`Conversation`/`Message`, +а общие типы содержимого и события живут в нижележащих модулях: +`Content`/`MessageContext`/`MessageOrigin`/`TurnTokens` — в `:content-api`, +`Event`/`OnlineEvent`/`OutboxStore`/`OnlineOutbox` — в `:outbox-api`. + ```kotlin -interface Agent { - fun id: String - suspend fun createConversation(title: String? = null): Conversation +interface Agent : AutoCloseable { + val id: String + val info: AgentInfo + val journal: JournalStore // append-only audit (read-only) + val outbox: OutboxStore // durable-события: catchup+live по курсору + val onlineOutbox: OnlineOutbox // live-only: стриминг, без курсора + val conversationStore: ConversationStore + fun createConversation(temp: Boolean): Conversation suspend fun getConversation(id: String): Conversation? - suspend fun getConversations(offset: Int = 0): Flow - suspend fun events(after: Instant): Flow // created/deleted/renamed + suspend fun deleteConversation(id: String): Boolean + suspend fun renameConversation(id: String, title: String?): Instant? } interface Conversation : AutoCloseable { val id: String val updatedAt: Instant + val isTemporal: Boolean + val title: String? val isSupportImageInput: Boolean val isSupportImageOutput: Boolean - suspend fun send(content: List): Flow // write+read вместе, как раньше - suspend fun events(after: Instant): Flow // отдельная live-подписка - suspend fun getMessages(offset: Int = 0): Flow - suspend fun rename(title: String): Boolean - fun interrupt() + suspend fun send(content: List, context: MessageContext? = null) // fire-and-forget + suspend fun getMessages(after: Instant, offset: Int, limit: Int): List + suspend fun rename(title: String) + suspend fun interrupt() } +// :content-api sealed interface Content { - class Text(val body: String) : Content - class Image(val data: ByteArray, val mime: String) : Content + data class Text(val body: String) : Content + data class Image(val data: ByteArray, val mime: String) : Content } +enum class MessageOrigin { USER, SYSTEM, EVENT } +data class MessageContext(origin: MessageOrigin, description: String?, sourceId: String?, metadata: JsonElement?) +// :proto sealed interface Message { val id: String val date: Instant - interface Body : Message { val content: List } - interface System : Message - class UserMessage(...) : Body - class AssistantMessage(...) : Body - class ToolCall(...) : System - class ToolResult(...) : System + class UserMessage(id, content: List, date, context: MessageContext?) : Message + class AssistantMessage(id, content: List, date, tokens: TurnTokens?, reasoning: String?) : Message + class ToolCall(...) : Message + class ToolResult(...) : Message + class Error(...) : Message } +// :outbox-api — durable (перезапрашиваются по курсору `after`) sealed interface Event { - enum ResponseType { TEXT, IMAGE } - class StartReasoning(...) : Event - class StartResponse(val type: ResponseType) : Event - class AppendText(val body: String) : Event - class AppendImage(val body: ByteArray, val mime: String) : Event - class End(...) : Event - class Interrupted(...) : Event - class Error(val message: String, val code: Int? = null) : Event + val date: Instant + class UserMessage(date, id, content: List, context: MessageContext?) : Event + class AssistantMessage(date, id, content: List, reasoning: String?, tokens: TurnTokens?) : Event class ToolCall(...) : Event class ToolResult(...) : Event + class ToolFailed(...) : Event + class Interrupted(date) : Event + class Error(date, message, code) : Event + class ConversationClosing(...) : Event + class CompactionTriggered(...) : Event +} + +// :outbox-api — online (live-only, НИКОГДА не сохраняются) +sealed interface OnlineEvent { + val date: Instant + enum ResponseType { TEXT, IMAGE } + class Working(date) : OnlineEvent + class End(date) : OnlineEvent + class StartReasoning(date) : OnlineEvent + class StartResponse(date, responseType: ResponseType) : OnlineEvent + class AppendText(date, body: String) : OnlineEvent + class AppendImage(date, body: ByteArray, mime: String) : OnlineEvent } ``` +Ход в терминах маркеров: онлайн `Working` → … → онлайн `End`; durable — +`UserMessage` в начале и `AssistantMessage`/`Interrupted`/`Error` в конце. + ## Чего здесь НЕТ - Никакого HTTP/SSE/JSON. Это контракт. Сериализация живёт в `:server` diff --git a/proto/build.gradle.kts b/proto/build.gradle.kts index f97247d..48b2a9f 100644 --- a/proto/build.gradle.kts +++ b/proto/build.gradle.kts @@ -28,11 +28,12 @@ kotlin { // public-сигнатуре Agent, поэтому api-висимости. // // Линейный граф зависимостей (без циклов): - // :proto ──► :outbox-api (нет обратной зависимости) - // :proto ──► :journal-api (нет обратной зависимости) - // Добились переносом AgentEvent/CommonEvent/Event из :proto в - // :outbox-api — они теперь self-contained в outbox (не нужны - // :proto-типы), а :proto использует их через :outbox-api. + // :proto ──► :content-api (Content/MessageContext/MessageOrigin/TurnTokens) + // :proto ──► :outbox-api (нет обратной зависимости) + // :proto ──► :journal-api (нет обратной зависимости) + // Общие типы содержимого живут в низкоуровневом :content-api, + // чтобы ими пользовались proto/journal/outbox без дублей и циклов. + api(project(":content-api")) api(project(":journal-api")) api(project(":outbox-api")) } diff --git a/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Content.kt b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Content.kt deleted file mode 100644 index 68ac71c..0000000 --- a/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Content.kt +++ /dev/null @@ -1,18 +0,0 @@ -package pw.binom.agentik.proto - -import kotlinx.serialization.SerialName -import kotlinx.serialization.Serializable - -@Serializable -sealed interface Content { - @Serializable - @SerialName("text") - class Text(val body: String) : Content - - /** - * Картинка. [mime] — MIME-тип, например `"image/png"`, `"image/jpeg"`. - */ - @Serializable - @SerialName("image") - class Image(val data: ByteArray, val mime: String) : Content -} diff --git a/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Conversation.kt b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Conversation.kt index edcf82e..efa2126 100644 --- a/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Conversation.kt +++ b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Conversation.kt @@ -2,6 +2,8 @@ package pw.binom.agentik.proto import kotlinx.coroutines.flow.Flow import kotlinx.coroutines.flow.flow +import pw.binom.agentik.content.Content +import pw.binom.agentik.content.MessageContext import kotlin.time.Instant /** @@ -11,18 +13,18 @@ import kotlin.time.Instant * * **Live-события** диалога НЕ часть этого интерфейса. Их два независимых * потока: - * - **durable** ([pw.binom.agentik.outbox.Event]: Working / End / Interrupted / - * Error / ToolCall / ToolResult / ToolFailed) — из + * - **durable** ([pw.binom.agentik.outbox.Event]: UserMessage / AssistantMessage / + * Interrupted / Error / ToolCall / ToolResult / ToolFailed) — из * [pw.binom.agentik.outbox.OutboxStore], перезапрашивается по курсору: * ``` * agent.outbox.conversationEvents(after = lastSeen, conversationId = id) * .map { it.event } * .collect { e -> ... } * ``` - * - **online** ([pw.binom.agentik.outbox.OnlineEvent]: StartReasoning / - * StartResponse / AppendText / AppendImage) — live-only стриминг ответа - * из [pw.binom.agentik.outbox.OnlineOutbox], без catchup/курсора: - * `agent.onlineOutbox.onlineEvents(id)`. + * - **online** ([pw.binom.agentik.outbox.OnlineEvent]: Working / End / + * StartReasoning / StartResponse / AppendText / AppendImage) — live-only + * стриминг ответа из [pw.binom.agentik.outbox.OnlineOutbox], без + * catchup/курсора: `agent.onlineOutbox.onlineEvents(id)`. * * Для cross-conversation view (admin / parent-agent / debug): * `agent.outbox.events(after)`. Для lifecycle агента (created/deleted/renamed): diff --git a/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Message.kt b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Message.kt index 19b39ef..3d40b46 100644 --- a/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Message.kt +++ b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Message.kt @@ -2,6 +2,9 @@ package pw.binom.agentik.proto import kotlinx.serialization.SerialName import kotlinx.serialization.Serializable +import pw.binom.agentik.content.Content +import pw.binom.agentik.content.MessageContext +import pw.binom.agentik.content.TurnTokens import kotlin.time.Instant @Serializable @@ -33,7 +36,17 @@ sealed interface Message { @Serializable @SerialName("assistant_message") - class AssistantMessage(override val id: String, val content: List, override val date: Instant) : Message + class AssistantMessage( + override val id: String, + val content: List, + override val date: Instant, + val tokens: TurnTokens? = null, + /** + * Текст размышлений модели (chain-of-thought / reasoning), если + * провайдер его отдаёт. Опционально. + */ + val reasoning: String? = null, + ) : Message @Serializable @SerialName("tool_call") diff --git a/proto/src/commonMain/kotlin/pw/binom/agentik/proto/MessageContext.kt b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/MessageContext.kt deleted file mode 100644 index 837635f..0000000 --- a/proto/src/commonMain/kotlin/pw/binom/agentik/proto/MessageContext.kt +++ /dev/null @@ -1,94 +0,0 @@ -package pw.binom.agentik.proto - -import kotlinx.serialization.SerialName -import kotlinx.serialization.Serializable -import kotlinx.serialization.json.JsonElement - -/** - * Кто/что инициировал данный ход сообщения. - * - * Используется в [MessageContext] — каждый ход диалога может нести - * дополнительный контекст о природе триггера: - * - [USER] — обычное сообщение от пользователя в чате (дефолт, context=null). - * - [SYSTEM] — программное системное сообщение (старт агента, режим обслуживания, - * уведомление о завершении фоновой задачи). - * - [EVENT] — внешнее событие (cron, webhook, file-changed, и т.п.). - * В этом случае [MessageContext.sourceId] и [MessageContext.description] - * позволяют модели понять, что за источник её разбудил. - * - * Семантический контракт: - * - origin != USER ⇒ [MessageContext.description] обязателен и должен быть - * человекочитаемым (короткая фраза для модели). - * - origin == USER ⇒ context может быть `null` (дефолт), и если задан — поля - * интерпретируются как «дополнительная мета» (например, ui_client). - */ -@Serializable -enum class MessageOrigin { - @SerialName("user") - USER, - - @SerialName("system") - SYSTEM, - - @SerialName("event") - EVENT, -} - -/** - * Контекст инициации сообщения: кто/что и почему вызвало этот ход. - * - * Примеры: - * ``` - * // cron-задача утренней сводки - * MessageContext( - * origin = MessageOrigin.EVENT, - * description = "scheduled cron 'morning-briefing'", - * sourceId = "cron-42", - * metadata = buildJsonObject { put("scheduledAt", "2026-09-14T08:00:00Z") }, - * ) - * - * // обычное сообщение из IRC - * MessageContext( - * origin = MessageOrigin.USER, - * sourceId = "irc-channel:agentik", - * description = "PRIVMSG from nick", - * ) - * - * // старт агента после рестарта - * MessageContext( - * origin = MessageOrigin.SYSTEM, - * description = "agent startup greeting", - * ) - * ``` - * - * Сериализация: snake_case для стабильного wire-формата ([origin] идёт как - * `user`/`system`/`event` благодаря @SerialName на enum). - * - * Forward-совместимо: добавление новых полей — non-breaking для старых - * клиентов, которые их игнорируют. - */ -@Serializable -data class MessageContext( - val origin: MessageOrigin, - /** - * Короткая человекочитаемая фраза для LLM: попадает в working memory - * как префикс `[origin] description (sourceId=…)` к user-сообщению, - * чтобы модель видела, что её разбудил не пользователь, а событие. - */ - val description: String? = null, - /** - * Идентификатор источника: id cron-job'а, webhook endpoint'а, имя канала IRC, - * id фонового события. Помогает модели и оператору при логировании понять, - * откуда пришёл ход. - */ - val sourceId: String? = null, - /** - * Произвольный структурированный payload о событии. - * Например: `{"scheduledAt": "...", "rule": "..."}` для cron, - * или `{"headers": {...}, "ip": "..."}` для webhook. - * - * Никогда не попадает в LLM-нагрузку как сырой JSON — используется - * только для логирования и пост-аналитики. - */ - val metadata: JsonElement? = null, -) diff --git a/server/src/commonMain/kotlin/pw/binom/agentik/server/Dto.kt b/server/src/commonMain/kotlin/pw/binom/agentik/server/Dto.kt index 6cfed8a..9eec9c8 100644 --- a/server/src/commonMain/kotlin/pw/binom/agentik/server/Dto.kt +++ b/server/src/commonMain/kotlin/pw/binom/agentik/server/Dto.kt @@ -1,9 +1,9 @@ package pw.binom.agentik.server import kotlinx.serialization.Serializable -import pw.binom.agentik.proto.Content +import pw.binom.agentik.content.Content import pw.binom.agentik.proto.Conversation -import pw.binom.agentik.proto.MessageContext +import pw.binom.agentik.content.MessageContext import kotlin.time.Instant /** 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 17c628d..5f8ddb8 100644 --- a/server/src/commonMain/kotlin/pw/binom/agentik/server/Routes.kt +++ b/server/src/commonMain/kotlin/pw/binom/agentik/server/Routes.kt @@ -21,7 +21,7 @@ import kotlinx.serialization.KSerializer import kotlinx.serialization.json.Json import pw.binom.agentik.proto.Agent import pw.binom.agentik.proto.Conversation -import pw.binom.agentik.proto.Content +import pw.binom.agentik.content.Content import pw.binom.agentik.outbox.AgentEvent import pw.binom.agentik.outbox.CommonEvent import pw.binom.agentik.outbox.Event @@ -189,6 +189,15 @@ internal fun Route.agentikRoutes(agent: Agent) { call.streamJsonSse(agent.onlineOutbox.onlineEvents(id), OnlineEvent.serializer()) } + /** + * `GET /online` — **live-only** SSE поток онлайн-событий всех диалогов + * агента (аналог `/events/all` для durable). Без `after`: онлайн-события + * не хранятся и не реплеятся. + */ + get("/online") { + call.streamJsonSse(agent.onlineOutbox.onlineEvents(), OnlineEvent.serializer()) + } + get("/events") { val after = call.parseAfter() ?: return@get // agent.outbox.agentEvents(after) возвращает Flow; diff --git a/server/src/commonTest/kotlin/pw/binom/agentik/server/JournalRoutesCountTest.kt b/server/src/commonTest/kotlin/pw/binom/agentik/server/JournalRoutesCountTest.kt index 42e4c8e..0bada58 100644 --- a/server/src/commonTest/kotlin/pw/binom/agentik/server/JournalRoutesCountTest.kt +++ b/server/src/commonTest/kotlin/pw/binom/agentik/server/JournalRoutesCountTest.kt @@ -19,7 +19,7 @@ import kotlinx.coroutines.flow.emptyFlow import kotlinx.coroutines.runBlocking import kotlinx.serialization.Serializable import pw.binom.agentik.journal.ConversationStore -import pw.binom.agentik.journal.Content as JContent +import pw.binom.agentik.content.Content as JContent import pw.binom.agentik.journal.JournalStore import pw.binom.agentik.journal.MessageRecord import pw.binom.agentik.journal.inmemory.InMemoryJournalStore diff --git a/server/src/commonTest/kotlin/pw/binom/agentik/server/TestOnlineOutbox.kt b/server/src/commonTest/kotlin/pw/binom/agentik/server/TestOnlineOutbox.kt index 48be547..44b725e 100644 --- a/server/src/commonTest/kotlin/pw/binom/agentik/server/TestOnlineOutbox.kt +++ b/server/src/commonTest/kotlin/pw/binom/agentik/server/TestOnlineOutbox.kt @@ -9,6 +9,7 @@ import pw.binom.agentik.outbox.OnlineOutbox * онлайн-канал тестам не нужен, но интерфейс обязывает его отдать. */ internal fun emptyOnlineOutbox(): OnlineOutbox = object : OnlineOutbox { + override fun onlineEvents() = emptyFlow() override fun onlineEvents(conversationId: String) = emptyFlow() override fun close() {} } diff --git a/settings.gradle.kts b/settings.gradle.kts index 39f8ba1..39ef57f 100644 --- a/settings.gradle.kts +++ b/settings.gradle.kts @@ -79,6 +79,12 @@ include(":reflection-api") include(":journal-api") include(":context-api") include(":outbox-api") +// Общие типы содержимого сообщения (Content/MessageContext/MessageOrigin/ +// TurnTokens). Вынесены в отдельный низкоуровневый модуль, чтобы ими +// пользовались И :proto, И :journal-api, И :outbox-api (message-события +// несут List + MessageContext) без дублирования типов и циклов +// в графе зависимостей. +include(":content-api") // Bounded-tail event log с auto-TTL. Двухуровневое хранилище: этот модуль — // короткий live tail + recent replay; полный audit log живёт в :message-store-api // (там — MessageStore + ConversationStore). EventStore сам управляет eviction, diff --git a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/A2aBridge.kt b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/A2aBridge.kt index c4d9802..892a229 100644 --- a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/A2aBridge.kt +++ b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/A2aBridge.kt @@ -13,7 +13,7 @@ import pw.binom.a2a.server.AgentHandler import pw.binom.agentik.outbox.Event import pw.binom.agentik.outbox.OnlineEvent import pw.binom.agentik.proto.Agent -import pw.binom.agentik.proto.Content +import pw.binom.agentik.content.Content import pw.binom.agentik.proto.Conversation import java.util.concurrent.ConcurrentHashMap @@ -31,7 +31,7 @@ private val log = KotlinLogging.logger {} * Ответ A2A = склеенные [OnlineEvent.AppendText] нашего хода. Подписку на онлайн-поток * ([pw.binom.agentik.outbox.OnlineOutbox]) открываем ДО [Conversation.send] (live-only, * без catchup — события начала хода иначе можно упустить), завершение хода ждём - * по durable-событиям [Event.End] / [Event.Interrupted] / [Event.Error]. + * по онлайн [OnlineEvent.End] и durable [Event.AssistantMessage] / [Event.Interrupted] / [Event.Error]. * * Ограничение v1: tool-события и картинки в A2A-ответ не транслируются; * при нескольких ходов в очереди за контекстом текст предыдущего хода @@ -53,14 +53,20 @@ class A2aBridge(private val agent: Agent) : AgentHandler { // Онлайн-поток — дельты ответа (live-only, без catchup). val onlineJob = async { agent.onlineOutbox.onlineEvents(conv.id).collect { e -> - if (e is OnlineEvent.AppendText) reply.append(e.body) + when (e) { + is OnlineEvent.AppendText -> reply.append(e.body) + is OnlineEvent.End -> turnDone.complete(Unit) + else -> {} + } } } - // Durable-поток — терминатор хода (catchup + live). + // Durable-поток — терминатор хода (catchup + live): целый ответ + // приходит [Event.AssistantMessage], обрыв — [Event.Interrupted], + // провал — [Event.Error]. val turnJob = async { agent.outbox.conversationEvents(since, conv.id).collect { ce -> when (val e = ce.event) { - is Event.End, is Event.Interrupted -> turnDone.complete(Unit) + is Event.AssistantMessage, is Event.Interrupted -> turnDone.complete(Unit) is Event.Error -> turnDone.completeExceptionally( IllegalStateException("agent turn failed: ${e.message}") diff --git a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/DebugRoutes.kt b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/DebugRoutes.kt index fc0a2a0..6c94fde 100644 --- a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/DebugRoutes.kt +++ b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/DebugRoutes.kt @@ -137,5 +137,5 @@ internal suspend fun recentTurns(workingMemoryStore: ContextStore, conversationI } /** Текстовое содержимое записей 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/commonMain/kotlin/pw/binom/agentik/standalone/agent/ChatAgent.kt b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ChatAgent.kt index 1916a23..a80c958 100644 --- a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ChatAgent.kt +++ b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ChatAgent.kt @@ -380,10 +380,10 @@ private val onlineEventStore: MutableOnlineOutbox = pw.binom.agentik.outbox.inme for (entry in filtered.takeLast(limit * 2)) { val text = when (entry) { is pw.binom.agentik.context.WorkingMemoryEntry.User -> - entry.content.filterIsInstance() + entry.content.filterIsInstance() .joinToString("\n") { it.body } is pw.binom.agentik.context.WorkingMemoryEntry.Assistant -> - entry.content.filterIsInstance() + entry.content.filterIsInstance() .joinToString("\n") { it.body } else -> continue } diff --git a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/CompactionCoordinator.kt b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/CompactionCoordinator.kt index c662d1a..8c2742d 100644 --- a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/CompactionCoordinator.kt +++ b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/CompactionCoordinator.kt @@ -6,7 +6,7 @@ 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.journal.Content +import pw.binom.agentik.content.Content import pw.binom.agentik.context.WorkingMemoryEntry import pw.binom.agentik.context.WorkingMemoryRow import pw.binom.agentik.context.ContextStore diff --git a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ContextBuilder.kt b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ContextBuilder.kt index 9c76fa5..a5af284 100644 --- a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ContextBuilder.kt +++ b/standalone/src/commonMain/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.journal.MessageContext -import pw.binom.agentik.journal.MessageOrigin +import pw.binom.agentik.content.MessageContext +import pw.binom.agentik.content.MessageOrigin internal class ContextBuilder( private val memoryPrefetcher: MemoryPrefetcher?, diff --git a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ConversationLoop.kt b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ConversationLoop.kt index 1603156..f3333ca 100644 --- a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ConversationLoop.kt +++ b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ConversationLoop.kt @@ -19,22 +19,22 @@ import mu.KotlinLogging import pw.binom.agentik.memory.MemoryPrefetcher import pw.binom.agentik.memory.MemoryReviewer import pw.binom.agentik.memory.MemoryStore -import pw.binom.agentik.proto.Content as ProtoContent +import pw.binom.agentik.content.Content as ProtoContent import pw.binom.agentik.proto.Conversation as ProtoConversation import pw.binom.agentik.outbox.Event as ProtoEvent import pw.binom.agentik.outbox.OnlineEvent import pw.binom.agentik.outbox.MutableOnlineOutbox import pw.binom.agentik.reflection.ReflectionStore import pw.binom.agentik.proto.Message as ProtoMessage -import pw.binom.agentik.proto.MessageContext as ProtoMessageContext -import pw.binom.agentik.journal.Content +import pw.binom.agentik.content.MessageContext as ProtoMessageContext +import pw.binom.agentik.content.Content import pw.binom.agentik.journal.ConversationRecord import pw.binom.agentik.journal.MutableConversationStore -import pw.binom.agentik.journal.MessageContext -import pw.binom.agentik.journal.MessageOrigin +import pw.binom.agentik.content.MessageContext +import pw.binom.agentik.content.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.content.TurnTokens import pw.binom.agentik.context.WorkingMemoryEntry import pw.binom.agentik.context.ContextStore import pw.binom.agentik.toolsets.ToolsetDispatchPolicy @@ -179,11 +179,11 @@ class ConversationLoop( check(!state.isClosed) { "Conversation closed: $id" } val turnStarted = now() // Working-маркер — самый первый event хода. Эмитим синхронно через - // tryEmit (events.tryEmit → outbox.append, без сетевого I/O), чтобы - // клиент увидел «агент работает» ещё до turnLock.withLock { launch } - // и до первого токена от LLM. Терминатор — End/Interrupted/Error - // (см. KDoc Event.Working). - emitEvent(pw.binom.agentik.outbox.Event.Working(date = turnStarted)) + // tryEmit (events.tryEmitOnline → без сетевого I/O), чтобы клиент + // увидел «агент работает» ещё до turnLock.withLock { launch } и до + // первого токена от LLM. Working/End — онлайн-маркеры (live-only), + // терминатор durable-части — AssistantMessage/Interrupted/Error. + emitOnline(OnlineEvent.Working(date = turnStarted)) val userMessageId = newId("msg") val storageContext = context?.toStorage() @@ -206,6 +206,16 @@ class ConversationLoop( ), now = turnStarted, ) + // Durable-событие user-сообщения: позволяет восстановить историю + // по курсору outbox без отдельного запроса в journal. + emitEvent( + ProtoEvent.UserMessage( + date = turnStarted, + id = userMessageId, + content = userRecord.content, + context = storageContext, + ) + ) } turnLock.withLock { @@ -448,6 +458,17 @@ class ConversationLoop( ) messageStore.append(assistantRecord) + // Durable-событие готового ответа агента. + emitEvent( + ProtoEvent.AssistantMessage( + date = assistantAt, + id = assistantId, + content = assistantContent, + reasoning = null, + tokens = turnTokens, + ) + ) + workingMemoryStore.append( conversationId = id, entry = WorkingMemoryEntry.Assistant( @@ -489,7 +510,7 @@ class ConversationLoop( if (wasInterrupted || interrupted.get()) { emitEvent(ProtoEvent.Interrupted(date = now())) } - emitEvent(ProtoEvent.End(date = now())) + emitOnline(OnlineEvent.End(date = now())) interrupted.set(false) } @@ -541,9 +562,9 @@ internal fun ProtoContent.toStorage(): Content = when (this) { internal fun ProtoMessageContext.toStorage(): MessageContext = MessageContext( origin = when (origin) { - pw.binom.agentik.proto.MessageOrigin.USER -> MessageOrigin.USER - pw.binom.agentik.proto.MessageOrigin.SYSTEM -> MessageOrigin.SYSTEM - pw.binom.agentik.proto.MessageOrigin.EVENT -> MessageOrigin.EVENT + pw.binom.agentik.content.MessageOrigin.USER -> MessageOrigin.USER + pw.binom.agentik.content.MessageOrigin.SYSTEM -> MessageOrigin.SYSTEM + pw.binom.agentik.content.MessageOrigin.EVENT -> MessageOrigin.EVENT }, description = description, sourceId = sourceId, @@ -552,9 +573,9 @@ internal fun ProtoMessageContext.toStorage(): MessageContext = MessageContext( internal fun MessageContext.toProto(): ProtoMessageContext { val protoOrigin = when (origin) { - MessageOrigin.USER -> pw.binom.agentik.proto.MessageOrigin.USER - MessageOrigin.SYSTEM -> pw.binom.agentik.proto.MessageOrigin.SYSTEM - MessageOrigin.EVENT -> pw.binom.agentik.proto.MessageOrigin.EVENT + MessageOrigin.USER -> pw.binom.agentik.content.MessageOrigin.USER + MessageOrigin.SYSTEM -> pw.binom.agentik.content.MessageOrigin.SYSTEM + MessageOrigin.EVENT -> pw.binom.agentik.content.MessageOrigin.EVENT } return ProtoMessageContext( origin = protoOrigin, diff --git a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ReflectionScheduler.kt b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ReflectionScheduler.kt index bf8f57e..02b0161 100644 --- a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ReflectionScheduler.kt +++ b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ReflectionScheduler.kt @@ -9,7 +9,7 @@ import kotlinx.coroutines.flow.onEach import kotlinx.coroutines.launch import mu.KotlinLogging import pw.binom.agentik.memory.ConversationTurn -import pw.binom.agentik.journal.Content +import pw.binom.agentik.content.Content import pw.binom.agentik.reflection.ReflectionStore import pw.binom.agentik.context.WorkingMemoryEntry import pw.binom.agentik.context.ContextStore diff --git a/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/agent/ChatAgentTest.kt b/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/agent/ChatAgentTest.kt index 54516a3..8ab3b9d 100644 --- a/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/agent/ChatAgentTest.kt +++ b/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/agent/ChatAgentTest.kt @@ -9,7 +9,7 @@ import kotlinx.coroutines.launch import kotlinx.coroutines.test.runTest import pw.binom.agentik.outbox.AgentEvent import pw.binom.agentik.outbox.OnlineEvent -import pw.binom.agentik.proto.Content +import pw.binom.agentik.content.Content import pw.binom.agentik.outbox.Event as ProtoEvent import pw.binom.agentik.skill.mining.SkillReadTool import pw.binom.agentik.skills.SkillCatalog @@ -186,7 +186,7 @@ class ChatAgentTest { pw.binom.agentik.journal.MessageRecord.UserMessage( id = "m1", conversationId = id, - content = listOf(pw.binom.agentik.journal.Content.Text("hi")), + content = listOf(pw.binom.agentik.content.Content.Text("hi")), createdAt = Instant.fromEpochMilliseconds(1_700_000_000_000), ), ) @@ -214,10 +214,10 @@ class ChatAgentTest { val msgs = sqliteStores.messages.listFlow(conv.id, Instant.DISTANT_PAST).toList() assertEquals(2, msgs.size) assertEquals("hi", (msgs[0] as pw.binom.agentik.journal.MessageRecord.UserMessage).content.let { - (it[0] as pw.binom.agentik.journal.Content.Text).body + (it[0] as pw.binom.agentik.content.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 + (it[0] as pw.binom.agentik.content.Content.Text).body }) } @@ -341,12 +341,18 @@ class ChatAgentTest { val job = launch(start = kotlinx.coroutines.CoroutineStart.UNDISPATCHED) { agent.outbox.conversationEvents(Instant.DISTANT_PAST, conv.id).collect { events.add(it.event) } } + val online = mutableListOf() + val onlineJob = launch(start = kotlinx.coroutines.CoroutineStart.UNDISPATCHED) { + agent.onlineOutbox.onlineEvents(conv.id).collect { online.add(it) } + } conv.send(listOf(Content.Text("hi"))) delay(50) job.cancel() + onlineJob.cancel() assertTrue(events.any { it is ProtoEvent.Error && it.message == "boom from llm" }, "events=$events") - assertTrue(events.any { it is ProtoEvent.End }, "events=$events") + assertTrue(events.any { it is ProtoEvent.UserMessage }, "events=$events") + assertTrue(online.any { it is OnlineEvent.End }, "online=$online") val msgs = sqliteStores.messages.listFlow(conv.id, Instant.DISTANT_PAST).toList() assertEquals(2, msgs.size) @@ -363,14 +369,11 @@ class ChatAgentTest { } @Test - fun `Working is first durable event, streaming goes to online channel`() = runTest { - // Working-маркер обязан прийти самым первым durable-событием хода, до - // End. Это позволяет UI показать спиннер сразу же при отправке, не - // дожидаясь первого токена от LLM. - // - // Стриминг (StartReasoning / StartResponse / AppendText) — теперь - // онлайн-события (live-only, не сохраняются) и приходят из - // [OnlineOutbox], а НЕ из durable-потока. + fun `message events are durable, Working-and-End and streaming go online`() = runTest { + // Durable-поток несёт целые сообщения (UserMessage/AssistantMessage). + // Working/End — маркеры фаз (онлайн, live-only), как и стриминг + // (StartReasoning / StartResponse / AppendText), который приходит + // из [OnlineOutbox], а НЕ из durable-потока. val agent = newAgent() fakeLlm.reply = "ok" val conv = agent.createConversation(temp = false) @@ -388,17 +391,20 @@ class ChatAgentTest { durableJob.cancel() onlineJob.cancel() - // Working — первый durable event хода (индекс 0). + // durable: первое событие хода — user-сообщение; терминатор успешного + // хода — готовое assistant-сообщение. assertTrue(durable.isNotEmpty(), "no events captured: $durable") val first = durable.first() - assertIs(first) - assertTrue(durable.any { it is ProtoEvent.End }, "no End: $durable") - // Ранее стриминговые маркеры — теперь в онлайн-канале. + assertIs(first) + assertTrue(durable.any { it is ProtoEvent.AssistantMessage }, "no AssistantMessage: $durable") + // Working/End и стриминговые маркеры — онлайн-канал (live-only). + assertTrue(online.any { it is OnlineEvent.Working }, "no Working: $online") assertTrue(online.any { it is OnlineEvent.StartReasoning }, "no StartReasoning: $online") assertTrue(online.any { it is OnlineEvent.StartResponse }, "no StartResponse: $online") + assertTrue(online.any { it is OnlineEvent.End }, "no End: $online") // Working не дублируется (emit'ится один раз на send()). - assertEquals(1, durable.count { it is ProtoEvent.Working }, "Working emitted >1 times: $durable") + assertEquals(1, online.count { it is OnlineEvent.Working }, "Working emitted >1 times: $online") } @Test @@ -417,6 +423,10 @@ class ChatAgentTest { val eventsJob = launch(start = kotlinx.coroutines.CoroutineStart.UNDISPATCHED) { agent.outbox.conversationEvents(Instant.DISTANT_PAST, conv.id).collect { events.add(it.event) } } + val online = mutableListOf() + val onlineJob = launch(start = kotlinx.coroutines.CoroutineStart.UNDISPATCHED) { + agent.onlineOutbox.onlineEvents(conv.id).collect { online.add(it) } + } val sendJob = launch { try { @@ -430,6 +440,7 @@ class ChatAgentTest { conv.interrupt() sendJob.join() eventsJob.cancel() + onlineJob.cancel() // audit: только user (assistant не успел сгенериться) val msgs = sqliteStores.messages.listFlow(conv.id, Instant.DISTANT_PAST).toList() @@ -441,9 +452,9 @@ class ChatAgentTest { assertEquals(1, wm.size) assertTrue(wm[0].entry is WorkingMemoryEntry.User) - // events: должны включать Interrupted + End + // durable: включая Interrupted; End — онлайн-маркер. assertTrue(events.any { it is ProtoEvent.Interrupted }, "events=$events") - assertTrue(events.any { it is ProtoEvent.End }, "events=$events") + assertTrue(online.any { it is OnlineEvent.End }, "online=$online") } @Test @@ -504,10 +515,10 @@ class ChatAgentTest { assertFalse(exchanges[0].wasCancelled, "tool реально выполнился, не был отменён") assertTrue(exchanges[0].resultText.contains("echo")) - // events должны включать ToolCall + ToolResult. End — обязательно (turn завершился). + // events должны включать ToolCall + ToolResult. AssistantMessage — обязательно (turn завершился). assertTrue(events.any { it is ProtoEvent.ToolCall }, "events=$events") assertTrue(events.any { it is ProtoEvent.ToolResult }, "events=$events") - assertTrue(events.any { it is ProtoEvent.End }, "events=$events") + assertTrue(events.any { it is ProtoEvent.AssistantMessage }, "events=$events") } @Test @@ -615,7 +626,7 @@ class ChatAgentTest { val agent = newAgent(llm = toolLlm, tools = listOf(echoTool)) val conv = agent.createConversation(temp = false) - conv.send(listOf(pw.binom.agentik.proto.Content.Text("call the tool"))) + conv.send(listOf(pw.binom.agentik.content.Content.Text("call the tool"))) // sendStreamContents вызывается дважды: первый раз с user-сообщением // (LLM отвечает tool_call), второй раз — после addToolResult — для diff --git a/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/agent/CompactionTest.kt b/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/agent/CompactionTest.kt index 38a666a..191ad98 100644 --- a/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/agent/CompactionTest.kt +++ b/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/agent/CompactionTest.kt @@ -12,7 +12,7 @@ import pw.binom.agentik.standalone.llm.LlmBackend import pw.binom.agentik.standalone.llm.LlmConfig import pw.binom.agentik.standalone.llm.OpenAiConfig import pw.binom.agentik.standalone.persistence.SqliteStores -import pw.binom.agentik.proto.Content as ProtoContent +import pw.binom.agentik.content.Content as ProtoContent import pw.binom.litert.LiteConversation import pw.binom.litert.LiteConversationConfig import pw.binom.litert.LiteLlm diff --git a/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/agent/ContextPrefixTest.kt b/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/agent/ContextPrefixTest.kt index 95ceea2..2bd1b79 100644 --- a/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/agent/ContextPrefixTest.kt +++ b/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/agent/ContextPrefixTest.kt @@ -1,9 +1,9 @@ package pw.binom.agentik.standalone.agent -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.agentik.content.MessageContext +import pw.binom.agentik.content.MessageOrigin.EVENT +import pw.binom.agentik.content.MessageOrigin.SYSTEM +import pw.binom.agentik.content.MessageOrigin.USER import pw.binom.litert.LiteContentPart import kotlin.test.Test import kotlin.test.assertEquals diff --git a/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/agent/MemoryWiringTest.kt b/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/agent/MemoryWiringTest.kt index 6c7516b..44fa64c 100644 --- a/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/agent/MemoryWiringTest.kt +++ b/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/agent/MemoryWiringTest.kt @@ -25,7 +25,7 @@ import pw.binom.agentik.memory.MemorySystemGuidance import pw.binom.agentik.memory.NewMemoryNote import pw.binom.agentik.memory.ReviewedTurn import pw.binom.agentik.memory.md.openMdMemorySystem -import pw.binom.agentik.proto.Content +import pw.binom.agentik.content.Content import pw.binom.agentik.standalone.agent.memory.MemoryToolsFactory import pw.binom.agentik.standalone.llm.LlmBackend import pw.binom.agentik.standalone.llm.LlmConfig diff --git a/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/persistence/PersistenceTest.kt b/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/persistence/PersistenceTest.kt index c10e9c9..5096849 100644 --- a/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/persistence/PersistenceTest.kt +++ b/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/persistence/PersistenceTest.kt @@ -1,9 +1,9 @@ package pw.binom.agentik.standalone.persistence -import pw.binom.agentik.journal.MessageContext -import pw.binom.agentik.journal.MessageOrigin +import pw.binom.agentik.content.MessageContext +import pw.binom.agentik.content.MessageOrigin import pw.binom.agentik.journal.ConversationRecord import pw.binom.agentik.journal.MessageRecord -import pw.binom.agentik.journal.Content +import pw.binom.agentik.content.Content import pw.binom.agentik.context.WorkingMemoryEntry import kotlinx.coroutines.flow.toList