feat(client): refactor AgentikAgent to manage its own HttpClient
ci / JVM build + tests (push) Failing after 54s
release / Publish KMP libraries → caffeine Nexus (release) Failing after 5s

- `AgentikAgent` now accepts `engineFactory` and an optional `token` to create an internal `HttpClient`, handling all configuration (JSON, Bearer).
- Removed `applyAgentikDefaults` and replaced it with `agentikHttpClient` for `HttpClient` creation with consistent settings.
- Updated `Agent` to implement `AutoCloseable`, ensuring proper resource closure with `agent.close()`.
- Adjusted tests, docs, and examples to align with the new `AgentikAgent` API.
This commit is contained in:
2026-09-21 03:53:35 +03:00
parent fb963bfb6b
commit f946186ef5
7 changed files with 95 additions and 102 deletions
+32 -29
View File
@@ -6,64 +6,63 @@
## Что есть ## Что есть
- `AgentikAgent(id, baseUrl, httpClient)` — entry-point. Возвращает `Agent` - `AgentikAgent(id, baseUrl, engineFactory, token?)` — entry-point. Возвращает
(тот же интерфейс, что в `:proto`). `Agent` (тот же интерфейс, что в `:proto`). HttpClient создаётся внутри
из переданной `engineFactory` (`CIO`, `OkHttp`, `Darwin`).
- `Agent`: `createConversation` / `getConversation` / `getConversations` / - `Agent`: `createConversation` / `getConversation` / `getConversations` /
`deleteConversation` / `journal` / `outbox`. `deleteConversation` / `journal` / `outbox` / `close`.
- `Conversation`: `send(content, context?)` / `events(after)` (SSE `Flow<Event>`) - `Conversation`: `send(content, context?)` / `events(after)` (SSE `Flow<Event>`)
/ `getMessages(after, offset, limit)` / `rename` / `interrupt` / `close`. / `getMessages(after, offset, limit)` / `rename` / `interrupt` / `close`.
- `HttpJournalStore` — `list(convId, after, offset, limit)` → `List<MessageRecord>` - `HttpJournalStore` — `list(convId, after, offset, limit)` → `List<MessageRecord>`
со всеми типами записей (User/Assistant/ToolCall/ToolResult/Error + tokens). со всеми типами записей (User/Assistant/ToolCall/ToolResult/Error + tokens).
- `HttpEventStore` — `events` / `agentEvents` / `conversationEvents` (SSE). - `HttpEventStore` — `events` / `agentEvents` / `conversationEvents` (SSE).
`HttpClient` создаётся снаружи (выбор движка — на тебе: CIO, OkHttp, `Agent` — `AutoCloseable`; `agent.close()` закрывает HttpClient. Не нужно
Darwin). Конфигурация (JSON + Bearer-токен) — через `applyAgentikDefaults`. вручную создавать `HttpClient` и накатывать на него JSON/Bearer-плагины.
## Подключение ## Подключение
```kotlin ```kotlin
// build.gradle.kts // build.gradle.kts
kotlin { dependencies {
sourceSets.commonMain.dependencies {
api("pw.binom.agentik:client:0.1.0") api("pw.binom.agentik:client:0.1.0")
// Движок — на твой выбор (один из):
implementation("io.ktor:ktor-client-cio:3.x") // JVM/Native
implementation("io.ktor:ktor-client-okhttp:3.x") // JVM
implementation("io.ktor:ktor-client-darwin:3.x") // iOS/macOS
// Опционально — только если будешь использовать `InMemoryJournalStore` // Опционально — только если будешь использовать `InMemoryJournalStore`
// как клиентский кэш. Свой `MutableJournalStore` — не нужен. // как клиентский кэш. Свой `MutableJournalStore` — не нужен.
api("pw.binom.agentik:journal-inmemory:0.1.0") api("pw.binom.agentik:journal-inmemory:0.1.0")
}
} }
``` ```
## Быстрый старт: свой клиент за 5 минут ## Быстрый старт: свой клиент за 5 минут
Один self-contained пример: создаём HTTP-клиент, открываем диалог, Один self-contained пример: создаём агента, открываем диалог,
отправляем сообщение, печатаем streaming-ответ. отправляем сообщение, печатаем streaming-ответ.
```kotlin ```kotlin
import pw.binom.agentik.client.AgentikAgent import pw.binom.agentik.client.AgentikAgent
import pw.binom.agentik.client.applyAgentikDefaults
import pw.binom.agentik.proto.Content import pw.binom.agentik.proto.Content
import pw.binom.agentik.proto.Event import pw.binom.agentik.proto.Event
import io.ktor.client.HttpClient
import io.ktor.client.engine.cio.CIO import io.ktor.client.engine.cio.CIO
import kotlinx.coroutines.runBlocking import kotlinx.coroutines.runBlocking
import kotlin.time.Clock import kotlin.time.Clock
fun main() = runBlocking { fun main() = runBlocking {
// 1. HTTP-клиент. Движок выбираешь сам (CIO/OkHttp/Darwin). // 1. Agent — обёртка над :server фасадом. HttpClient создаётся внутри.
val http = HttpClient(CIO) { applyAgentikDefaults(token = "s3cret") }
// 2. Agent — обёртка над :server фасадом.
val agent = AgentikAgent( val agent = AgentikAgent(
id = "my-client", id = "my-client",
baseUrl = "http://localhost:8080/agentik", baseUrl = "http://localhost:8080/agentik",
httpClient = http, engineFactory = CIO,
token = "s3cret", // или null, если не нужен
) )
// 3. Открыть диалог, отправить сообщение. // 2. Открыть диалог, отправить сообщение.
val conv = agent.createConversation(temp = false) val conv = agent.createConversation(temp = false)
conv.send(listOf(Content.Text("Привет"))) conv.send(listOf(Content.Text("Привет")))
// 4. Собирать streaming-ответ. // 3. Собирать streaming-ответ.
conv.events(after = Clock.System.now()).collect { ev -> conv.events(after = Clock.System.now()).collect { ev ->
when (ev) { when (ev) {
is Event.StartResponse -> println("[start]") is Event.StartResponse -> println("[start]")
@@ -74,15 +73,18 @@ fun main() = runBlocking {
} }
} }
// 5. Чистый shutdown. // 4. Чистый shutdown.
conv.close() conv.close()
http.close() agent.close()
} }
``` ```
**Это весь клиент.** `:server` сам хранит историю, контекст, события. **Это весь клиент.** `:server` сам хранит историю, контекст, события.
Ты только получаешь типизированный `Flow<Event>` и рендеришь как хочешь. Ты только получаешь типизированный `Flow<Event>` и рендеришь как хочешь.
`HttpClient`, `applyAgentikDefaults`, выбор engine'а — всё скрыто
внутри `AgentikAgent`. Один вызов — один готовый `Agent`.
### Добавить локальный кэш истории (ещё 4 строки) ### Добавить локальный кэш истории (ещё 4 строки)
```kotlin ```kotlin
@@ -117,14 +119,17 @@ history.forEach { rec ->
### Что вообще не нужно писать самому ### Что вообще не нужно писать самому
- HTTP-сериализация `Event`/`Message` — `applyAgentikDefaults` регистрирует - HTTP-сериализация `Event`/`Message` — `agentikHttpClient` регистрирует
`agentikJson` и `InstantSerializer`. `agentikJson` и `InstantSerializer`.
- SSE-парсер — `readSse()` внутри `:client`. - SSE-парсер — `readSse()` внутри `:client`.
- Cursor-менеджмент для `listFlow` — дефолтная имплементация в - Cursor-менеджмент для `listFlow` — дефолтная имплементация в
`JournalStore.listFlow` сама пагинирует. `JournalStore.listFlow` сама пагинирует.
- Lifecycle подписок на `events()` — `Conversation.close()` отменяет SSE-job. - Lifecycle подписок на `events()` — `Conversation.close()` отменяет SSE-job.
- Bearer-токен в каждом запросе — `applyAgentikDefaults(token = ...)` инжектит - HTTP-клиент и Bearer — `AgentikAgent` создаёт `HttpClient(engineFactory)`
один раз на весь `HttpClient`. с Bearer'ом из `token=` под капотом; `agent.close()` его закрывает.
- Движковые настройки (requestTimeout и пр.) — `HttpClient(engineFactory) { ... }`
создаётся здесь; для нестандартных движковых настроек используй
`agentikHttpClient(engineFactory, token)` напрямую (он экспортирован).
### Что нужно написать самому ### Что нужно написать самому
@@ -138,17 +143,14 @@ history.forEach { rec ->
```kotlin ```kotlin
import pw.binom.agentik.client.AgentikAgent import pw.binom.agentik.client.AgentikAgent
import pw.binom.agentik.client.applyAgentikDefaults
import pw.binom.agentik.proto.Content import pw.binom.agentik.proto.Content
import pw.binom.agentik.proto.Event import pw.binom.agentik.proto.Event
import io.ktor.client.HttpClient
import io.ktor.client.engine.cio.CIO import io.ktor.client.engine.cio.CIO
val http = HttpClient(CIO) { applyAgentikDefaults(token = "s3cret") }
val agent = AgentikAgent( val agent = AgentikAgent(
id = "agentik", id = "agentik",
baseUrl = "http://localhost:8080/agentik", baseUrl = "http://localhost:8080/agentik",
httpClient = http, engineFactory = CIO,
) )
val conv = agent.createConversation(temp = false) val conv = agent.createConversation(temp = false)
@@ -309,8 +311,9 @@ UI-обновление списка — отдельная задача, реш
- **Персистентность кэша** — `InMemoryJournalStore` хранит в RAM. Для - **Персистентность кэша** — `InMemoryJournalStore` хранит в RAM. Для
диска пиши свой `MutableJournalStore` (см. `KsqliteJournalStore` в диска пиши свой `MutableJournalStore` (см. `KsqliteJournalStore` в
`:journal-ksqlite` как образец). `:journal-ksqlite` как образец).
- **Авторизация** — `applyAgentikDefaults(token = "...")` для Bearer; - **Нестандартные движковые настройки** — для `requestTimeout`,
для OAuth/что-то ещё — конфигурируй `HttpClient` сам. прокси и т.п. используй `agentikHttpClient(engineFactory, token)`
напрямую.
## Тесты ## Тесты
@@ -7,6 +7,7 @@ import io.ktor.client.request.get
import io.ktor.client.request.parameter import io.ktor.client.request.parameter
import io.ktor.client.request.post import io.ktor.client.request.post
import io.ktor.client.request.setBody import io.ktor.client.request.setBody
import io.ktor.client.statement.HttpResponse
import io.ktor.http.ContentType import io.ktor.http.ContentType
import io.ktor.http.HttpStatusCode import io.ktor.http.HttpStatusCode
import io.ktor.http.contentType import io.ktor.http.contentType
@@ -19,36 +20,20 @@ import pw.binom.agentik.proto.Conversation
/** /**
* HTTP-реализация [Agent]. Ходит в `:server`-фасад, см. `agentikAgent(...)`. * HTTP-реализация [Agent]. Ходит в `:server`-фасад, см. `agentikAgent(...)`.
* *
* Замечание по [createConversation]: интерфейс [Agent] объявлен не-suspend * HttpClient создаётся внутри из переданного engine и закрывается в [close].
* (in-process кейс этого не требует), но HTTP-вариант обязан ждать ответа
* POST `/conversations`. Используем `runBlocking` — это одноразовая
* операция (открытие чата), не горячий путь. В UI-контексте вызывающий сам
* решает, что делать.
* *
* **Storage handles** ([journal], [outbox]) — read-only views на серверные * **Storage handles** ([journal], [outbox]) — read-only views на серверные
* хранилища. [outbox] уже реализован ([HttpEventStore]); [journal] — * хранилища.
* заглушка, потому что соответствующий HTTP endpoint'ы (`/journal/...`)
* ещё не выставлены на стороне `:server`. После их добавления подменить
* `error(...)` на `HttpJournalStore(...)`.
*/ */
internal class AgentClient( internal class AgentClient(
private val httpClient: HttpClient,
private val baseUrl: String,
override val id: String, override val id: String,
private val baseUrl: String,
private val httpClient: HttpClient,
) : Agent { ) : Agent {
private val agentUrl: String = baseUrl.trimEnd('/') private val agentUrl: String = baseUrl.trimEnd('/')
/**
* Единый канал событий (lifecycle + per-conversation). Под капотом —
* [HttpEventStore]: каждый метод бьёт свой URL (см. KDoc).
*/
override val outbox: OutboxStore = HttpEventStore(httpClient = httpClient, baseUrl = agentUrl) override val outbox: OutboxStore = HttpEventStore(httpClient = httpClient, baseUrl = agentUrl)
/**
* HTTP-фасад для journal: ходит в `:server`'s `GET /journal/conversations/{id}/messages`.
* См. [HttpJournalStore] и [pw.binom.agentik.server.journalRoutes].
*/
override val journal: JournalStore = HttpJournalStore(httpClient = httpClient, baseUrl = agentUrl) override val journal: JournalStore = HttpJournalStore(httpClient = httpClient, baseUrl = agentUrl)
override fun createConversation(temp: Boolean): Conversation = override fun createConversation(temp: Boolean): Conversation =
@@ -68,7 +53,7 @@ internal class AgentClient(
} }
override suspend fun deleteConversation(id: String): Boolean { override suspend fun deleteConversation(id: String): Boolean {
val response = httpClient.delete("$agentUrl/conversations/$id") val response: HttpResponse = httpClient.delete("$agentUrl/conversations/$id")
return response.status == HttpStatusCode.NoContent return response.status == HttpStatusCode.NoContent
} }
@@ -79,4 +64,8 @@ internal class AgentClient(
}.body<List<ConversationSnapshot>>() }.body<List<ConversationSnapshot>>()
return snapshots.map { ConversationClient(httpClient, agentUrl, it) } return snapshots.map { ConversationClient(httpClient, agentUrl, it) }
} }
override fun close() {
httpClient.close()
}
} }
@@ -1,38 +1,43 @@
package pw.binom.agentik.client package pw.binom.agentik.client
import io.ktor.client.HttpClient import io.ktor.client.engine.HttpClientEngineFactory
import pw.binom.agentik.proto.Agent import pw.binom.agentik.proto.Agent
/** /**
* Создаёт [Agent], который под капотом ходит в HTTP-фасад `agentikAgent` * Создаёт [Agent], который ходит в HTTP-фасад `agentikAgent` (модуль `:server`).
* (модуль `:server`). *
* Принимает [engineFactory] — `HttpClientEngineFactory<*>` (`CIO`, `OkHttp`,
* `Darwin`, ...). Внутри сам создаёт `HttpClient`, накатывает JSON-конфиг
* [agentikJson] и опциональный Bearer [token]. Никакого `applyAgentikDefaults`
* снаружи — всё под капотом.
* *
* ``` * ```
* val http = HttpClient(engine) { * val agent = AgentikAgent(
* applyAgentikDefaults(token = "s3cret") * id = "my-client",
* // engine — на выбор потребителя (CIO, OkHttp, Darwin, ...);
* // движковые настройки (requestTimeout и пр.) — там же
* }
* val client = AgentikAgent(
* id = "my-agent",
* baseUrl = "http://localhost:8080/agentik", * baseUrl = "http://localhost:8080/agentik",
* httpClient = http, * engineFactory = CIO,
* token = "s3cret",
* ) * )
* val conv = client.createConversation(temp = false) * val conv = agent.createConversation(temp = false)
* conv.send(listOf(Content.Text("hi"))) * conv.send(listOf(Content.Text("hi")))
* conv.events(Instant.DISTANT_PAST).collect { ev -> ... } * conv.events(Instant.DISTANT_PAST).collect { ... }
* agent.close() // закрывает HttpClient
* ``` * ```
* *
* [id] пробрасывается в реализацию [Agent.id] — сервер про идентичность * [id] пробрасывается в `Agent.id` — сервер про идентичность агента не знает,
* агента не знает, поэтому клиент должен её знать сам (или взять из * поэтому клиент должен её знать сам (или взять из конфига).
* конфига).
* *
* `:client` НЕ выбирает движок: [HttpClient] (с уже созданным engine'ом) * **Lifecycle**: [Agent] — `AutoCloseable`. `agent.close()` закрывает
* приходит снаружи. К блоку конфигурации применяется [applyAgentikDefaults] * HttpClient (идемпотентно). После этого `createConversation` /
* (JSON + опциональный Bearer-токен). * `getConversation` etc. не определены.
*/ */
fun AgentikAgent( fun AgentikAgent(
id: String, id: String,
baseUrl: String, baseUrl: String,
httpClient: HttpClient, engineFactory: HttpClientEngineFactory<*>,
): Agent = AgentClient(httpClient = httpClient, baseUrl = baseUrl, id = id) token: String? = null,
): Agent = AgentClient(
id = id,
baseUrl = baseUrl,
httpClient = agentikHttpClient(engineFactory = engineFactory, token = token),
)
@@ -1,7 +1,7 @@
package pw.binom.agentik.client package pw.binom.agentik.client
import io.ktor.client.HttpClient import io.ktor.client.HttpClient
import io.ktor.client.HttpClientConfig import io.ktor.client.engine.HttpClientEngineFactory
import io.ktor.client.plugins.DefaultRequest import io.ktor.client.plugins.DefaultRequest
import io.ktor.client.plugins.contentnegotiation.ContentNegotiation import io.ktor.client.plugins.contentnegotiation.ContentNegotiation
import io.ktor.client.request.header import io.ktor.client.request.header
@@ -9,32 +9,20 @@ import io.ktor.http.HttpHeaders
import io.ktor.serialization.kotlinx.json.json import io.ktor.serialization.kotlinx.json.json
/** /**
* Общая конфигурация HTTP-клиента agentik — платформо-независимая часть. * Создаёт [HttpClient] поверх [engineFactory] с конфигурацией agentik.
* *
* `:client` НЕ выбирает движок: его приносит потребитель. Здесь живёт только то, * Внутренний helper для [AgentikAgent]. Потребителю `:client` обычно
* без чего клиент несовместим с `/agentik`: * не нужен — он передаёт engine в [AgentikAgent] и получает готовый
* - JSON-конфиг [agentikJson] (обязан совпадать с серверным); * [pw.binom.agentik.proto.Agent] с уже закрытым HttpClient'ом
* - при заданном [token] — `Authorization: Bearer <token>` на ВСЕ запросы * на [pw.binom.agentik.proto.Agent.close].
* через [DefaultRequest] (накрывает 10 REST-вызовов и оба SSE-потока;
* заголовок живёт на клиенте, а не в отдельных запросах).
* *
* `null` — авторизация выключена, заголовок не отправляется. * Экспортируется для случаев, когда нужен прямой доступ к `HttpClient`
* * (например, дополнительные нестандартные запросы в обход `Agent` API).
* Применяется к уже сконструированному [HttpClient]:
* ```
* val engine = HttpClientEngineFactory().create() // потребитель выбирает движок
* val http = HttpClient(engine) {
* applyAgentikDefaults(token)
* // движковые настройки (requestTimeout и пр.) — потребитель знает свой движок
* }
* ```
*
* Раньше `:client` экспортировал фабрику `agentikHttpClient(engineFactory, ...)`,
* но она навязывала тип-параметр `HttpClientEngineFactory<T>` и неявно тянула за
* собой конкретный движок в виде примера в KDoc. Сейчас фабрики нет — потребитель
* сам создаёт [HttpClient] и накатывает на блок конфигурации [applyAgentikDefaults].
*/ */
fun HttpClientConfig<*>.applyAgentikDefaults(token: String? = null) { fun agentikHttpClient(
engineFactory: HttpClientEngineFactory<*>,
token: String? = null,
): HttpClient = HttpClient(engineFactory) {
install(ContentNegotiation) { json(agentikJson) } install(ContentNegotiation) { json(agentikJson) }
if (token != null) { if (token != null) {
install(DefaultRequest) { install(DefaultRequest) {
@@ -21,7 +21,7 @@ import kotlin.test.Test
import kotlin.test.assertEquals import kotlin.test.assertEquals
/** /**
* Тесты клиентской части: [applyAgentikDefaults] с заданным `token` прикладывает * Тесты клиентской части: [agentikHttpClient] с заданным `token` прикладывает
* `Authorization: Bearer <token>` ко всем запросам через плагин `DefaultRequest`, * `Authorization: Bearer <token>` ко всем запросам через плагин `DefaultRequest`,
* без токена — заголовок не отправляется. * без токена — заголовок не отправляется.
* *
@@ -50,7 +50,7 @@ class BearerHeaderTest {
} }
private fun clientWith(token: String?): HttpClient = private fun clientWith(token: String?): HttpClient =
HttpClient(CIO) { applyAgentikDefaults(token) } agentikHttpClient(engineFactory = CIO, token = token)
private suspend fun startServer(): Pair<EmbeddedServer<*, *>, Int> { private suspend fun startServer(): Pair<EmbeddedServer<*, *>, Int> {
val server = embeddedServer(ServerCIO, port = 0) { val server = embeddedServer(ServerCIO, port = 0) {
@@ -19,11 +19,18 @@ import kotlin.time.Instant
* `outbox.agentEvents(after)`. Это даёт единый путь для всех read-операций * `outbox.agentEvents(after)`. Это даёт единый путь для всех read-операций
* по хранилищу и убирает дублирование между протоколом и хранилищем. * по хранилищу и убирает дублирование между протоколом и хранилищем.
*/ */
public interface Agent { interface Agent : AutoCloseable {
/** Идентификатор агента. */ /** Идентификатор агента. */
val id: String val id: String
/**
* Освобождает ресурсы агента (HTTP-клиент, сетевые handles, подписки).
* После [close] вызовы [createConversation] / [getConversation] и т.п.
* не определены. Idempotent.
*/
override fun close()
/** /**
* Append-only audit log всех сообщений диалогов (read-only view). * Append-only audit log всех сообщений диалогов (read-only view).
* *
@@ -57,6 +57,7 @@ class BearerTokenTest {
override suspend fun getConversation(id: String): Conversation? = null override suspend fun getConversation(id: String): Conversation? = null
override suspend fun deleteConversation(id: String): Boolean = false override suspend fun deleteConversation(id: String): Boolean = false
override suspend fun getConversations(offset: Int, limit: Int): List<Conversation> = emptyList() override suspend fun getConversations(offset: Int, limit: Int): List<Conversation> = emptyList()
override fun close() {}
} }
private suspend fun startServer(token: String?): Pair<EmbeddedServer<*, *>, Int> { private suspend fun startServer(token: String?): Pair<EmbeddedServer<*, *>, Int> {