From 8f85612665f026161df2149f3990106508388fc8 Mon Sep 17 00:00:00 2001 From: subochev Date: Mon, 21 Sep 2026 02:43:29 +0300 Subject: [PATCH] feat(client): implement `HttpJournalStore` and integrate journal endpoints - Added `HttpJournalStore` as an HTTP-backed implementation of `JournalStore` for read-only access to the audit log. - Integrated `GET /journal/conversations/{id}/messages` endpoint to fetch conversation transcripts with full payloads. - Updated `AgentClient` to expose `HttpJournalStore` as the `journal` property. - Adjusted `HttpEventStore` to align with updated endpoint structure (`/outbox/events`). --- client/build.gradle.kts | 1 + .../pw/binom/agentik/client/AgentClient.kt | 9 +-- .../pw/binom/agentik/client/HttpEventStore.kt | 23 ++++--- .../binom/agentik/client/HttpJournalStore.kt | 60 +++++++++++++++++++ 4 files changed, 74 insertions(+), 19 deletions(-) create mode 100644 client/src/commonMain/kotlin/pw/binom/agentik/client/HttpJournalStore.kt diff --git a/client/build.gradle.kts b/client/build.gradle.kts index 2a95945..00d451d 100644 --- a/client/build.gradle.kts +++ b/client/build.gradle.kts @@ -22,6 +22,7 @@ kotlin { commonMain.dependencies { api(project(":proto")) api(project(":event-store")) + api(project(":journal-api")) api(libs.ktor.client.core) implementation(libs.ktor.client.content.negotiation) 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 4617919..8b3ce5a 100644 --- a/client/src/commonMain/kotlin/pw/binom/agentik/client/AgentClient.kt +++ b/client/src/commonMain/kotlin/pw/binom/agentik/client/AgentClient.kt @@ -46,13 +46,10 @@ internal class AgentClient( override val outbox: OutboxStore = HttpEventStore(httpClient = httpClient, baseUrl = agentUrl) /** - * HTTP-фасад для journal пока не реализован: на стороне `:server` ещё - * не выставлены endpoint'ы `/journal/conversations/{id}/messages`. - * Как только появятся — заменить на `HttpJournalStore(httpClient, agentUrl)`. + * HTTP-фасад для journal: ходит в `:server`'s `GET /journal/conversations/{id}/messages`. + * См. [HttpJournalStore] и [pw.binom.agentik.server.journalRoutes]. */ - override val journal: JournalStore = error( - "HttpJournalStore ещё не реализован — дождаться :server endpoint'а /journal/...", - ) + override val journal: JournalStore = HttpJournalStore(httpClient = httpClient, baseUrl = agentUrl) override fun createConversation(temp: Boolean): Conversation = runBlocking { 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 b1ad455..2bd3719 100644 --- a/client/src/commonMain/kotlin/pw/binom/agentik/client/HttpEventStore.kt +++ b/client/src/commonMain/kotlin/pw/binom/agentik/client/HttpEventStore.kt @@ -17,21 +17,18 @@ import kotlin.time.Instant * HTTP-реализация [EventStore] (= [pw.binom.agentik.outbox.OutboxStore]), * ходящая в `:server`-фасад. * - * **Хитрый план**: вместо того, чтобы все методы шли в один общий endpoint и - * фильтровали client-side ([EventStore.events]/[filterIsInstance]), эта - * реализация бьёт запросы по URL'ам в зависимости от того, какой класс - * событий нужен: - * - [events] → `GET /events/all` (полный поток CommonEvent) - * - [agentEvents] → `GET /events` (только lifecycle диалогов) + * **Endpoint-раскладка** (новый дизайн — storage handles на [Agent]): + * - [events] → `GET {baseUrl}/outbox/events?after=` (полный поток + * [CommonEvent], bounded-tail + live SSE, см. [pw.binom.agentik.server.outboxRoutes]) + * - [agentEvents] → `GET {baseUrl}/events?after=` (legacy proto-роут: + * сервер пробрасывает [pw.binom.agentik.outbox.agentEvents] и распаковывает + * `.event` для обратной совместимости с форматом AgentEvent) * - [conversationEvents] с `conversationId != null` → `GET /conversations/{id}/events` * - * Так серверный фильтр (SQL `WHERE` или разные буферы) работает на своей стороне, - * а клиент получает ровно тот срез, который ему нужен, без лишнего трафика. - * * Для [conversationEvents] с `conversationId == null` (события всех диалогов) - * fallback на default [EventStore.conversationEvents] — общий поток `/events/all` - * + фильтр client-side. Это редкий кейс (admin-дашборды), и оптимизировать его - * отдельно нерационально. + * fallback на default [EventStore.conversationEvents] — общий поток + * `/outbox/events` + filter. Это редкий кейс (admin-дашборды), и + * оптимизировать его отдельно нерационально. * * [earliestEventDate] не имеет своего endpoint'а; возвращает `Clock.System.now()` * (см. KDoc [EventStore.earliestEventDate] — для пустого буфера это и есть @@ -53,7 +50,7 @@ internal class HttpEventStore( override fun events(after: Instant?): Flow = flow { val url = buildString { - append("$agentUrl/events/all") + append("$agentUrl/outbox/events") if (after != null) append("?after=$after") } httpClient.prepareGet(url) { noSseReadTimeout() } diff --git a/client/src/commonMain/kotlin/pw/binom/agentik/client/HttpJournalStore.kt b/client/src/commonMain/kotlin/pw/binom/agentik/client/HttpJournalStore.kt new file mode 100644 index 0000000..ea5fb06 --- /dev/null +++ b/client/src/commonMain/kotlin/pw/binom/agentik/client/HttpJournalStore.kt @@ -0,0 +1,60 @@ +package pw.binom.agentik.client + +import io.ktor.client.HttpClient +import io.ktor.client.call.body +import io.ktor.client.request.get +import io.ktor.client.request.parameter +import io.ktor.http.HttpStatusCode +import pw.binom.agentik.journal.JournalStore +import pw.binom.agentik.journal.MessageRecord +import kotlin.time.Instant + +/** + * HTTP-реализация [JournalStore] (append-only audit log сообщений диалога), + * ходящая в `:server`-фасад. + * + * **Endpoint**: `GET {baseUrl}/journal/conversations/{id}/messages?after=&offset=&limit=` + * (см. [pw.binom.agentik.server.journalRoutes]). + * + * Возвращает raw [MessageRecord] (все типы: UserMessage / AssistantMessage / + * ToolCall / ToolResult / Error). В отличие от `GET /conversations/{id}/messages` + * в `:server`'s proto-роутах (который отдаёт project'нутые + * [pw.binom.agentik.proto.Message]), здесь клиент получает полный transcript + * с tool-call/tool-result/error payload'ами, turn-tokens и context'ом. + * + * **listFlow** — default cold-flow paging через [list] (N+1 round-trip, + * дефолтная реализация из [JournalStore]). Для remote/SQL-backed store'а + * это OK: server-side paging + client-side flow compose'ится естественно. + * + * **Read-only**: [JournalStore] не имеет `append` — запись только через + * writer-референс, который ChatAgent держит внутри (тип + * `MutableJournalStore`, не выставлен наружу через [pw.binom.agentik.proto.Agent]). + */ +internal class HttpJournalStore( + private val httpClient: HttpClient, + private val baseUrl: String, +) : JournalStore { + + private val agentUrl: String = baseUrl.trimEnd('/') + + override suspend fun list( + conversationId: String, + after: Instant, + offset: Int, + limit: Int, + ): List { + val response = httpClient.get("$agentUrl/journal/conversations/$conversationId/messages") { + parameter("after", after.toString()) + parameter("offset", offset) + parameter("limit", limit) + } + check(response.status == HttpStatusCode.OK) { + "journal.list: server returned ${response.status}" + } + return response.body>() + } + + override fun close() { + // HttpClient закрывает владелец (AgentClient / AgentikAgent). + } +}