proto/server/client: Agent.info, Working-маркер хода, /record, таймауты запросов

- proto: Agent.info (AgentInfo: name/description/usefulness) — человекочитаемое
  имя отдельно от opaque id; сериализация и тесты
- server: GET {path} отдаёт Agent.info; GET /conversations/{id}/record —
  ConversationRecord без handle'а (для клиентского кэша после AgentEvent.Created)
- outbox: Event.Working — первый event хода, эмитится из Conversation.send()
  до LLM-цикла и turnLock, чтобы UI показал спиннер сразу
- standalone: AGENTIK_NAME/AGENTIK_DESCRIPTION/AGENTIK_USEFULNESS → AgentInfo
- client: AgentClient.create eagerly фетчит info (GET {baseUrl})
- client: HttpConversationStore.get читает /record (ConversationRecord,
  а не ConversationSnapshot — рассинхрон типов)
- client: noReadTimeout() на POST /conversations/{id}/messages — сервер отвечает
  по завершении всего хода агента (реально 0.5–144 с), дефолтные 15 с рвали
  живую реплику на клиенте
- journal-api: ConversationRecord @Serializable
- ksqlite 0.1.3 → 0.1.4
- .gitignore: runtime-данные standalone-агента и hs_err-дампы
This commit is contained in:
2026-09-28 10:09:34 +03:00
parent 4156e5f95f
commit 509763cf36
24 changed files with 681 additions and 16 deletions
+5
View File
@@ -30,3 +30,8 @@ agentik.db
agentik.db-shm
agentik.db-wal
memory-md/agentik-mem-*/
# Runtime-данные standalone-агента (db/memory/skills при локальном запуске)
/standalone/agentik/
hs_err_pid*.log
core.*
@@ -7,6 +7,7 @@ import pw.binom.agentik.journal.ConversationStore
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.proto.Conversation
import pw.binom.agentik.outbox.Event
@@ -23,6 +24,7 @@ internal class FakeAgent(
private val conversationFactory: () -> FakeConversation = { FakeConversation() },
) : Agent {
override val id: String = "fake"
override val info: AgentInfo = AgentInfo(name = "fake")
var createCount: Int = 0
private set
val conversations = mutableListOf<FakeConversation>()
@@ -16,6 +16,7 @@ import pw.binom.agentik.journal.ConversationStore
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.Conversation
import kotlin.time.Instant
@@ -27,9 +28,16 @@ import kotlin.time.Instant
* **Storage handles** ([journal], [outbox], [conversationStore]) — read-only
* views на серверные хранилища. Запись — только через команды
* [createConversation] / [deleteConversation] / [renameConversation].
*
* Конструируется через suspend [Companion.create], который **eagerly**
* фетчит [info] (`GET {baseUrl}`) и сохраняет снимок в поле. Это убирает
* необходимость в `lazy { runBlocking { ... } }` на горячем пути —
* `runBlocking` живёт один раз в [Companion.create], оттуда же [AgentikAgent]
* его и вызывает (там он приемлем: одноразовая инициализация агента).
*/
internal class AgentClient(
internal class AgentClient private constructor(
override val id: String,
override val info: AgentInfo,
private val baseUrl: String,
private val httpClient: HttpClient,
) : Agent {
@@ -74,4 +82,32 @@ internal class AgentClient(
override fun close() {
httpClient.close()
}
companion object {
/**
* Создаёт [AgentClient] и eagerly загружает [Agent.info]
* (`GET {baseUrl}` на серверном фасаде). Любой сбой сети на этом
* этапе пробрасывается как исключение — агент без `info` бесполезен
* (UI/A2A сразу упрутся в `agent.info`).
*
* Single-shot инициализация, `runBlocking` тут допустим (см. KDoc
* класса). Хосты, которым нужен полностью неблокирующий старт,
* могут обернуть вызов в свой `CoroutineScope`.
*/
suspend fun create(
id: String,
baseUrl: String,
httpClient: HttpClient,
): AgentClient {
val agentUrl = baseUrl.trimEnd('/')
val info: AgentInfo = httpClient.get(agentUrl).body()
return AgentClient(
id = id,
info = info,
baseUrl = agentUrl,
httpClient = httpClient,
)
}
}
}
@@ -6,6 +6,7 @@ import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.Job
import kotlinx.coroutines.SupervisorJob
import kotlinx.coroutines.cancel
import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.launch
import kotlinx.coroutines.runBlocking
import pw.binom.agentik.journal.ConversationRecord
@@ -95,7 +96,7 @@ fun AgentikAgent(
token: String? = null,
): Agent {
val httpClient = agentikHttpClient(engineFactory = engineFactory, token = token)
val client = AgentClient(id = id, baseUrl = baseUrl, httpClient = httpClient)
val client = runBlocking { AgentClient.create(id = id, baseUrl = baseUrl, httpClient = httpClient) }
return wrapWithLocalConversationCache(client, scopeClient = client)
}
@@ -59,6 +59,9 @@ internal class ConversationClient(
httpClient.post("$convUrl/messages") {
contentType(ContentType.Application.Json)
setBody(SendPayload(content, context))
// Сервер отвечает только по завершении хода агента (LLM + тулы),
// а это минуты, а не 15 секунд дефолтного request-timeout.
noReadTimeout()
}
}
@@ -33,7 +33,10 @@ internal class HttpConversationStore(
private val agentUrl: String = baseUrl.trimEnd('/')
override suspend fun get(id: String): ConversationRecord? {
val response = httpClient.get("$agentUrl/conversations/$id")
// `/record` (а НЕ `/conversations/{id}`): последний отдаёт
// ConversationSnapshot для `AgentClient.getConversation`, у которого
// другой shape (handle + isImageSupported, без createdAt/updatedAt).
val response = httpClient.get("$agentUrl/conversations/$id/record")
if (response.status == HttpStatusCode.NotFound) return null
check(response.status == HttpStatusCode.OK) {
"conversationStore.get($id): server returned ${response.status}"
@@ -53,7 +53,7 @@ internal class HttpEventStore(
append("$agentUrl/outbox/events")
if (after != null) append("?after=$after")
}
httpClient.prepareGet(url) { noSseReadTimeout() }
httpClient.prepareGet(url) { noReadTimeout() }
.execute { response ->
check(response.status == HttpStatusCode.OK) {
"events: server returned ${response.status}"
@@ -74,7 +74,7 @@ internal class HttpEventStore(
append("$agentUrl/events")
if (after != null) append("?after=$after")
}
httpClient.prepareGet(url) { noSseReadTimeout() }
httpClient.prepareGet(url) { noReadTimeout() }
.execute { response ->
check(response.status == HttpStatusCode.OK) {
"agentEvents: server returned ${response.status}"
@@ -104,7 +104,7 @@ internal class HttpEventStore(
append("$agentUrl/conversations/$conversationId/events")
if (after != null) append("?after=$after")
}
httpClient.prepareGet(url) { noSseReadTimeout() }
httpClient.prepareGet(url) { noReadTimeout() }
.execute { response ->
check(response.status == HttpStatusCode.OK) {
"conversationEvents: server returned ${response.status}"
@@ -8,18 +8,23 @@ import io.ktor.client.request.HttpRequestBuilder
* Отключает request/connect/socket-таймауты для конкретного запроса через
* [HttpTimeoutCapability] со всеми таймаутами = [HttpTimeoutConfig.INFINITE_TIMEOUT_MS].
*
* Зачем: наш SSE-ридер ([readSse]) читает `bodyAsChannel()` руками и не
* использует плагин `SSE`, поэтому движок не считает запрос SSE-шным
* (`HttpRequestBuilder.supportsRequestTimeout` проверяет
* `body is SSEClientContent`, а у нас тело — обычный GET без тела).
* Без capability встроенный `HttpTimeoutPlugin.requestTimeoutMillis` (по умолчанию
* **15000 мс**) молча убивает долгий idle-стрим через 15 секунд.
* Зачем: два вида запросов живут дольше дефолтных 15 секунд:
* - **SSE-чтение** ([readSse]) — читает `bodyAsChannel()` руками и не
* использует плагин `SSE`, поэтому движок не считает запрос SSE-шным
* ([HttpRequestBuilder.supportsRequestTimeout] проверяет
* `body is SSEClientContent`, а у нас тело — обычный GET без тела).
* - **`POST /conversations/{id}/messages`** — сервер отвечает не сразу, а
* только когда ход агента полностью завершён (LLM + тулы). Реальный ход
* легко длится минуты, и дефолтный request-timeout убивал бы его на
* 15-й секунде, обрывая ещё живой ход на сервере.
* Без capability встроенный `HttpTimeoutPlugin.requestTimeoutMillis`
* (по умолчанию **15000 мс**) молча убивает такой запрос.
*
* Конфиг создаётся заново на каждый вызов — плагин `HttpTimeout` при
* установленном capability мутирует его поля через `?:`, так что шаренный
* инстанс мог бы утечь между запросами.
*/
internal fun HttpRequestBuilder.noSseReadTimeout() {
internal fun HttpRequestBuilder.noReadTimeout() {
setCapability(
HttpTimeoutCapability,
HttpTimeoutConfig(
+1 -1
View File
@@ -16,7 +16,7 @@ logback = "1.5.18"
mosaic = "0.18.0"
clikt = "5.0.3"
kotlinx-cli = "0.3.6"
ksqlite = "0.1.3"
ksqlite = "0.1.4"
junit = "4.13.2"
[plugins]
@@ -1,10 +1,12 @@
package pw.binom.agentik.journal
import kotlinx.serialization.Serializable
import kotlin.time.Instant
/**
* Snapshot диалога. В таблице `conversation` хранится как есть.
*/
@Serializable
data class ConversationRecord(
val id: String,
val title: String?,
@@ -12,8 +12,13 @@ import kotlin.time.Instant
* порядка при равных timestamps.
*
* Базовая структура хода:
* `StartReasoning?` → `StartResponse(TEXT|IMAGE)` → ...контент... → `End` | `Interrupted` | `Error`.
* `Working` → `StartReasoning?` → `StartResponse(TEXT|IMAGE)` → ...контент... → `End` | `Interrupted` | `Error`.
* `StartReasoning` может отсутствовать, если агент не показывал рассуждения.
* `Working` — маркер «агент принял запрос и пошёл обрабатывать», эмитится
* синхронно в `Conversation.send()` ДО старта LLM-цикла (и возможной долгой
* очереди turnLock'а). Парный терминатор не нужен: `End`/`Interrupted`/`Error`
* уже закрывают ход. UI использует `Working` чтобы показать спиннер ещё до
* первого токена ответа.
*
* **История**: до 2026-09-21 жил в `:proto` как `pw.binom.agentik.proto.Event`;
* при миграции в `:outbox-api` был оставлен typealias в `:proto` для
@@ -33,6 +38,16 @@ sealed interface Event {
@SerialName("image") IMAGE
}
/**
* Маркер «агент принял запрос и пошёл обрабатывать». Эмитится **до**
* [StartReasoning]/[StartResponse], синхронно из `Conversation.send()`,
* чтобы UI мог показать спиннер ещё до первого токена ответа.
* Терминатор хода ([End]/[Interrupted]/[Error]) — парный.
*/
@Serializable
@SerialName("working")
data class Working(override val date: Instant) : Event
/** Ассистент начал рассуждение (опциональный маркер; контент рассуждения приходит через [AppendText]). */
@Serializable
@SerialName("start_reasoning")
@@ -28,6 +28,26 @@ interface Agent : AutoCloseable {
/** Идентификатор агента. */
val id: String
/**
* Публичное представление имени/назначения агента — то, что агент
* сам "рассказывает о себе" клиентам и UI.
*
* Контрактно отличается от [id]:
* - [id] — opaque stable identifier (логи, корреляция, A2A contextId);
* - [info.name] — **человекочитаемое** имя, которое пользователь
* может сменить через конфиг или UI; может меняться в течение жизни
* агента.
*
* Используется HTTP-фасадом `:server` (endpoint под `path` самого
* агента) и A2A-фасадом (для заполнения AgentCard `name` / `description`).
*
* Доступ к [info] всегда дешёвый — это snapshot `data class` без побочных
* эффектов. Конкретные реализации ([MutableAgent]) могут выставлять
* `info` конструкторным параметром; [Agent] контрактно не требует
* реактивности — клиенты получают снимок в момент обращения.
*/
val info: AgentInfo
/**
* Освобождает ресурсы агента (HTTP-клиент, сетевые handles, подписки).
* После [close] вызовы [createConversation] / [getConversation] и т.п.
@@ -0,0 +1,34 @@
package pw.binom.agentik.proto
import kotlinx.serialization.SerialName
import kotlinx.serialization.Serializable
/**
* Идентификационная информация об [Agent] — публичное представление для
* UI/A2A-агентов/IRC/визуальных дашбордов: имя ("как зовут ассистента"),
* краткое описание ("что он делает") и зачем он может быть полезен
* ("в каких задачах помогает").
*
* Контрактное отличие от [Agent.id]:
* - [Agent.id] — стабильный opaque identifier для корреляции, логов и
* транспорта (A2A `contextId`, IRC CTCP-запросы, etc). Менять нельзя.
* - [AgentInfo.name] — **человекочитаемое имя**, которое показывается
* клиенту и которое пользователь может настроить (например,
* `AGENTIK_NAME` в standalone-конфиге). Может меняться.
*
* Wire-формат (через `:server` JSON-фасад): `name` — обязательное поле,
* `description` и `usefulness` — опциональные, могут отсутствовать.
*
* @property name Человекочитаемое имя ассистента. Обязательное.
* @property description Краткое описание (что агент делает/кто он).
* Не сериализуется, если `null`.
* @property usefulness Зачем агент может быть полезен. Не сериализуется,
* если `null`.
*/
@Serializable
data class AgentInfo(
val name: String,
val description: String? = null,
@SerialName("usefulness")
val usefulness: String? = null,
)
@@ -0,0 +1,70 @@
package pw.binom.agentik.proto
import kotlinx.serialization.encodeToString
import kotlinx.serialization.json.Json
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertNull
/**
* Сериализация [AgentInfo] — контракт, по которому
* `GET {path}` отдаёт JSON клиенту/A2A.
*
* JSON-форма `description` и `usefulness` исключаются, когда null
* (благодаря `explicitNulls = false` в `agentikJson`).
*/
class AgentInfoSerializationTest {
private val json = Json {
explicitNulls = false
}
@Test
fun onlyName_whenOptionalsAreNull() {
val info = AgentInfo(name = "agentik")
val s = json.encodeToString(info)
assertEquals("""{"name":"agentik"}""", s)
}
@Test
fun allFields_whenAllSet() {
val info = AgentInfo(
name = "Hermes",
description = "Ассистент-разработчик на Kotlin",
usefulness = "Помогает писать код, искать баги, объяснять legacy",
)
val s = json.encodeToString(info)
assertEquals(
"""{"name":"Hermes","description":"Ассистент-разработчик на Kotlin","usefulness":"Помогает писать код, искать баги, объяснять legacy"}""",
s,
)
}
@Test
fun descriptionOnly_usefulnessOmitted() {
val info = AgentInfo(name = "x", description = "d")
val s = json.encodeToString(info)
assertEquals("""{"name":"x","description":"d"}""", s)
}
@Test
fun usefulnessOnly_descriptionOmitted() {
val info = AgentInfo(name = "x", usefulness = "u")
val s = json.encodeToString(info)
assertEquals("""{"name":"x","usefulness":"u"}""", s)
}
@Test
fun roundTrip_preservesFields() {
val original = AgentInfo(name = "Hermes", description = "d", usefulness = "u")
val decoded = json.decodeFromString<AgentInfo>(json.encodeToString(original))
assertEquals(original, decoded)
}
@Test
fun defaultConstructor_optionalsAreNull() {
val info = AgentInfo(name = "x")
assertNull(info.description)
assertNull(info.usefulness)
}
}
@@ -33,6 +33,19 @@ internal fun Route.agentikRoutes(agent: Agent) {
call.respondText("ok")
}
/**
* `GET {path}` — публичная информация об агенте (`name` +
* опциональные `description` / `usefulness`). Контрактно отличается
* от `/health` тем, что отдаёт человекочитаемые данные об агенте,
* а не просто "alive"-сигнал.
*
* Используется UI-дашбордами / A2A-fetch'ерами AgentCard для
* отображения карточки агента без подписки на события.
*/
get("") {
call.respond(agent.info)
}
// ---- Agent: множество диалогов ----
get("/conversations") {
@@ -54,6 +67,25 @@ internal fun Route.agentikRoutes(agent: Agent) {
else call.respond(c.snapshot())
}
/**
* `GET /conversations/{id}/record` — лёгкая метадата записи
* ([pw.binom.agentik.journal.ConversationRecord]: id/title/isTemporal/
* createdAt/updatedAt) БЕЗ handle'а и image-support флагов.
*
* Отдельно от `GET /conversations/{id}` (который отдаёт
* [pw.binom.agentik.proto.ConversationSnapshot] — его ждёт
* `AgentClient.getConversation`): это тот же shape, что и элемент
* списка `GET /conversations`, его читает клиентский кэш
* (`HttpConversationStore.get`), чтобы дотянуть запись после
* outbox-события `AgentEvent.Created`.
*/
get("/conversations/{id}/record") {
val id = call.parameters["id"]!!
val rec = agent.conversationStore.get(id)
if (rec == null) call.respond(HttpStatusCode.NotFound)
else call.respond(rec)
}
delete("/conversations/{id}") {
val id = call.parameters["id"]!!
call.respond(if (agent.deleteConversation(id)) HttpStatusCode.NoContent else HttpStatusCode.NotFound)
@@ -0,0 +1,169 @@
package pw.binom.agentik.server
import io.ktor.client.HttpClient
import io.ktor.client.call.body
import io.ktor.client.engine.cio.CIO
import io.ktor.client.plugins.contentnegotiation.ContentNegotiation
import io.ktor.client.request.get
import io.ktor.client.statement.bodyAsText
import io.ktor.http.HttpHeaders
import io.ktor.http.HttpStatusCode
import io.ktor.serialization.kotlinx.json.json
import io.ktor.server.cio.CIO as ServerCIO
import io.ktor.server.engine.EmbeddedServer
import io.ktor.server.engine.embeddedServer
import io.ktor.server.routing.routing
import kotlinx.coroutines.flow.emptyFlow
import kotlinx.coroutines.runBlocking
import kotlinx.serialization.Serializable
import kotlinx.serialization.json.Json
import pw.binom.agentik.journal.ConversationStore
import pw.binom.agentik.journal.JournalStore
import pw.binom.agentik.journal.MessageRecord
import pw.binom.agentik.outbox.CommonEvent
import pw.binom.agentik.outbox.OutboxStore
import pw.binom.agentik.proto.Agent
import pw.binom.agentik.proto.AgentInfo
import pw.binom.agentik.proto.Conversation
import kotlin.test.AfterTest
import kotlin.test.BeforeTest
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.time.Instant
/**
* Интеграционные тесты endpoint'а `GET {path}` (без суффикса) — отдаёт
* публичную информацию об агенте ([AgentInfo]) для UI/A2A/IRC.
*
* Проверяет:
* - happy path с полным [AgentInfo] (name + description + usefulness);
* - минимальный [AgentInfo] (только name, optional поля null/отсутствуют);
* - JSON-shape: только присутствующие в [AgentInfo] поля, `null`-опциональные
* выкинуты (благодаря `explicitNulls = false` в `agentikJson`);
* - endpoint доступен без авторизации (auth отключён в этом тесте;
* bearer-token тест уже в [BearerTokenTest]).
*/
class AgentInfoRouteTest {
private lateinit var server: EmbeddedServer<*, *>
private var port: Int = 0
@Serializable
private data class InfoDto(val name: String, val description: String? = null, val usefulness: String? = null)
@BeforeTest
fun setup() {
val fake = object : Agent {
override val id: String = "test-agent"
override val info: AgentInfo = AgentInfo(
name = "Hermes",
description = "Kotlin-ассистент для agentik",
usefulness = "Помогает писать код и искать баги",
)
override val journal: JournalStore = object : JournalStore {
override suspend fun list(conversationId: String, after: Instant, offset: Int, limit: Int) = emptyList<MessageRecord>()
override suspend fun count(conversationId: String): Long = 0L
override suspend fun count(conversationId: String, after: Instant): Long = 0L
override fun listFlow(conversationId: String, after: Instant, pageSize: Int) = emptyFlow<MessageRecord>()
override fun close() {}
}
override val outbox: OutboxStore = object : OutboxStore {
override fun events(after: Instant?) = emptyFlow<CommonEvent>()
override fun agentEvents(after: Instant?) = emptyFlow<CommonEvent.Agent>()
override fun conversationEvents(after: Instant?, conversationId: String?) = emptyFlow<CommonEvent.Conversation>()
override suspend fun earliestEventDate(): Instant = Instant.DISTANT_PAST
override fun close() {}
}
override val conversationStore: ConversationStore = object : ConversationStore {
override suspend fun get(id: String) = null
override suspend fun list(offset: Int, limit: Int) = emptyList<pw.binom.agentik.journal.ConversationRecord>()
override fun close() {}
}
override fun createConversation(temp: Boolean): Conversation = TODO("not used")
override suspend fun getConversation(id: String): Conversation? = null
override suspend fun deleteConversation(id: String): Boolean = false
override suspend fun renameConversation(id: String, title: String?): Instant? = null
override fun close() {}
}
server = embeddedServer(ServerCIO, port = 0, host = "127.0.0.1") {
routing { agentikAgent(fake, path = "/agentik", token = null) }
}.start(wait = false)
port = runBlocking { server.engine.resolvedConnectors()[0].port }
}
@AfterTest
fun tearDown() {
server.stop(100, 200)
}
private fun client(): HttpClient = HttpClient(CIO) {
install(ContentNegotiation) {
json(Json { ignoreUnknownKeys = true })
}
}
@Test
fun get_returns_full_info_as_json() = runBlocking {
val response = client().get("http://127.0.0.1:$port/agentik")
assertEquals(HttpStatusCode.OK, response.status)
val dto = response.body<InfoDto>()
assertEquals(InfoDto("Hermes", "Kotlin-ассистент для agentik", "Помогает писать код и искать баги"), dto)
}
@Test
fun get_body_shape_has_all_three_fields() = runBlocking {
val response = client().get("http://127.0.0.1:$port/agentik")
val raw = response.bodyAsText()
// explicitNulls=false → все три поля всегда присутствуют у full-info agent'а.
assertEquals(
"""{"name":"Hermes","description":"Kotlin-ассистент для agentik","usefulness":"Помогает писать код и искать баги"}""",
raw,
)
}
@Test
fun get_returns_minimal_info_when_only_name_set() {
// override agent: minimal info (только name).
// Storage handles — noop-stub'ы (см. BearerTokenTest: agentikAgent
// вычитывает их eagerly при старте).
val minimal = object : Agent {
override val id: String = "minimal"
override val info: AgentInfo = AgentInfo(name = "agentik")
override val journal: JournalStore = object : JournalStore {
override suspend fun list(conversationId: String, after: Instant, offset: Int, limit: Int) = emptyList<MessageRecord>()
override suspend fun count(conversationId: String): Long = 0L
override suspend fun count(conversationId: String, after: Instant): Long = 0L
override fun listFlow(conversationId: String, after: Instant, pageSize: Int) = emptyFlow<MessageRecord>()
override fun close() {}
}
override val outbox: OutboxStore = object : OutboxStore {
override fun events(after: Instant?) = emptyFlow<CommonEvent>()
override fun agentEvents(after: Instant?) = emptyFlow<CommonEvent.Agent>()
override fun conversationEvents(after: Instant?, conversationId: String?) = emptyFlow<CommonEvent.Conversation>()
override suspend fun earliestEventDate(): Instant = Instant.DISTANT_PAST
override fun close() {}
}
override val conversationStore: ConversationStore = object : ConversationStore {
override suspend fun get(id: String) = null
override suspend fun list(offset: Int, limit: Int) = emptyList<pw.binom.agentik.journal.ConversationRecord>()
override fun close() {}
}
override fun createConversation(temp: Boolean): Conversation = TODO("not used")
override suspend fun getConversation(id: String): Conversation? = null
override suspend fun deleteConversation(id: String): Boolean = false
override suspend fun renameConversation(id: String, title: String?): Instant? = null
override fun close() {}
}
val server2 = embeddedServer(ServerCIO, port = 0, host = "127.0.0.1") {
routing { agentikAgent(minimal, path = "/agentik", token = null) }
}.start(wait = false)
try {
val p = runBlocking { server2.engine.resolvedConnectors()[0].port }
val raw = runBlocking { client().get("http://127.0.0.1:$p/agentik").bodyAsText() }
// explicitNulls=false выкидывает null-опциональные поля → только {"name":"agentik"}.
assertEquals("""{"name":"agentik"}""", raw)
} finally {
server2.stop(100, 200)
}
}
}
@@ -19,6 +19,7 @@ import pw.binom.agentik.journal.MessageRecord
import pw.binom.agentik.outbox.CommonEvent
import pw.binom.agentik.outbox.OutboxStore
import pw.binom.agentik.proto.Agent
import pw.binom.agentik.proto.AgentInfo
import pw.binom.agentik.proto.Conversation
import kotlin.time.Instant
import kotlin.test.Test
@@ -41,6 +42,7 @@ class BearerTokenTest {
*/
private class FakeAgent(
override val id: String = "test",
override val info: AgentInfo = AgentInfo(name = "test"),
) : Agent {
override val journal: JournalStore = object : JournalStore {
override suspend fun list(conversationId: String, after: Instant, offset: Int, limit: Int) = emptyList<MessageRecord>()
@@ -0,0 +1,181 @@
package pw.binom.agentik.server
import io.ktor.client.HttpClient
import io.ktor.client.call.body
import io.ktor.client.engine.cio.CIO
import io.ktor.client.plugins.contentnegotiation.ContentNegotiation
import io.ktor.client.request.get
import io.ktor.client.statement.bodyAsText
import io.ktor.http.HttpStatusCode
import io.ktor.serialization.kotlinx.json.json
import io.ktor.server.cio.CIO as ServerCIO
import io.ktor.server.engine.EmbeddedServer
import io.ktor.server.engine.embeddedServer
import io.ktor.server.routing.routing
import kotlinx.coroutines.flow.emptyFlow
import kotlinx.coroutines.runBlocking
import pw.binom.agentik.journal.ConversationRecord
import pw.binom.agentik.journal.ConversationStore
import pw.binom.agentik.journal.JournalStore
import pw.binom.agentik.journal.inmemory.InMemoryJournalStore
import pw.binom.agentik.journal.inmemory.InMemoryMutableConversationStore
import pw.binom.agentik.outbox.CommonEvent
import pw.binom.agentik.outbox.OutboxStore
import pw.binom.agentik.proto.Agent
import pw.binom.agentik.proto.AgentInfo
import kotlin.test.AfterTest
import kotlin.test.BeforeTest
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertNotNull
import kotlin.time.Duration.Companion.seconds
import kotlin.time.Instant
/**
* Тесты `GET /conversations` и `GET /conversations/{id}/record` — обе отдают
* [ConversationRecord] (лёгкую метадату элемента списка).
*
* Это **регрессия** на два бага, из-за которых клиентский
* `wrapWithLocalConversationCache` не наполнялся:
* 1. [ConversationRecord] не был `@Serializable` → `GET /conversations`
* падал 500 `Serializer for class ... not found`;
* 2. клиентский `HttpConversationStore.get` читал `ConversationRecord` из
* `GET /conversations/{id}`, который отдаёт `ConversationSnapshot`
* (другой shape: handle + isImageSupported, без createdAt/updatedAt)
* → `MissingFieldException: Field 'createdAt' is required`. Для этого
* заведён отдельный `/record`-endpoint.
*/
class ConversationRoutesTest {
private lateinit var server: EmbeddedServer<*, *>
private lateinit var journal: InMemoryJournalStore
private lateinit var store: InMemoryMutableConversationStore
private var port: Int = 0
private val t0 = Instant.parse("2026-09-22T10:00:00Z")
@BeforeTest
fun setup() {
journal = InMemoryJournalStore()
store = InMemoryMutableConversationStore()
val fakeAgent = FakeAgent(journal, store)
server = embeddedServer(ServerCIO, port = 0, host = "127.0.0.1") {
routing { agentikAgent(fakeAgent, path = "/agentik", token = null) }
}.start(wait = false)
port = runBlocking { server.engine.resolvedConnectors()[0].port }
}
@AfterTest
fun tearDown() {
server.stop(100, 200)
journal.close()
store.close()
}
private fun client(): HttpClient = HttpClient(CIO) {
install(ContentNegotiation) { json(agentikJson) }
}
private fun seed(id: String, title: String?, isTemporal: Boolean = false, updatedAt: Instant = t0) =
runBlocking {
store.upsert(ConversationRecord(id, title, isTemporal, t0, updatedAt))
}
@Test
fun `list returns serializable conversation records`() = runBlocking {
seed("c1", "Первый", updatedAt = t0)
seed("c2", null, isTemporal = true, updatedAt = t0 + 1.seconds)
val client = client()
try {
val resp = client.get("http://127.0.0.1:$port/agentik/conversations")
assertEquals(HttpStatusCode.OK, resp.status)
// До фикса `@Serializable` здесь был 500.
val records = resp.body<List<ConversationRecord>>()
// Сортировка хранилища — updatedAt DESC.
assertEquals(listOf("c2", "c1"), records.map { it.id })
assertEquals(listOf(null, "Первый"), records.map { it.title })
assertEquals(Instant.parse("2026-09-22T10:00:01Z"), records[0].updatedAt)
} finally {
client.close()
}
}
@Test
fun `list body carries timestamps that the client cache needs`() = runBlocking {
seed("c1", "Первый")
val client = client()
try {
val raw = client.get("http://127.0.0.1:$port/agentik/conversations").bodyAsText()
assertEquals(
"""[{"id":"c1","title":"Первый","isTemporal":false,""" +
""""createdAt":"2026-09-22T10:00:00Z","updatedAt":"2026-09-22T10:00:00Z"}]""",
raw,
)
} finally {
client.close()
}
}
@Test
fun `record returns conversation record by id`() = runBlocking {
seed("c1", "Первый", isTemporal = true)
val client = client()
try {
val resp = client.get("http://127.0.0.1:$port/agentik/conversations/c1/record")
assertEquals(HttpStatusCode.OK, resp.status)
val rec = assertNotNull(resp.body<ConversationRecord>())
assertEquals("c1", rec.id)
assertEquals("Первый", rec.title)
assertEquals(true, rec.isTemporal)
assertEquals(t0, rec.createdAt)
} finally {
client.close()
}
}
@Test
fun `record for unknown conversation returns not found`() = runBlocking {
val client = client()
try {
val resp = client.get("http://127.0.0.1:$port/agentik/conversations/never-existed/record")
assertEquals(HttpStatusCode.NotFound, resp.status)
} finally {
client.close()
}
}
/**
* Минимальный stub-агент: реальные [JournalStore] и
* [InMemoryMutableConversationStore]; outbox — noop, conversation-хендлы
* (`getConversation`/`createConversation`) тесту не нужны.
*/
private class FakeAgent(
private val js: JournalStore,
private val cs: ConversationStore,
) : Agent {
override val id: String = "test"
override val info: AgentInfo = AgentInfo(name = "test")
override val journal: JournalStore = js
override val outbox: OutboxStore = object : OutboxStore {
override fun events(after: Instant?) = emptyFlow<CommonEvent>()
override fun agentEvents(after: Instant?) = emptyFlow<CommonEvent.Agent>()
override fun conversationEvents(after: Instant?, conversationId: String?) =
emptyFlow<CommonEvent.Conversation>()
override suspend fun earliestEventDate(): Instant = Instant.DISTANT_PAST
override fun close() {}
}
override val conversationStore: ConversationStore = cs
override fun createConversation(temp: Boolean): pw.binom.agentik.proto.Conversation =
TODO("not used")
override suspend fun getConversation(id: String): pw.binom.agentik.proto.Conversation? = null
override suspend fun deleteConversation(id: String): Boolean = false
override suspend fun renameConversation(id: String, title: String?): Instant? = null
override fun close() {}
}
}
@@ -26,6 +26,7 @@ import pw.binom.agentik.journal.inmemory.InMemoryJournalStore
import pw.binom.agentik.outbox.CommonEvent
import pw.binom.agentik.outbox.OutboxStore
import pw.binom.agentik.proto.Agent
import pw.binom.agentik.proto.AgentInfo
import kotlin.test.AfterTest
import kotlin.test.BeforeTest
import kotlin.test.Test
@@ -177,6 +178,7 @@ class JournalRoutesCountTest {
*/
private class FakeAgent(private val js: JournalStore) : Agent {
override val id: String = "test"
override val info: AgentInfo = AgentInfo(name = "test")
override val journal: JournalStore = js
override val outbox: OutboxStore = object : OutboxStore {
override fun events(after: Instant?) = emptyFlow<CommonEvent>()
@@ -374,6 +374,7 @@ private fun runServer() {
val agent = ChatAgent(
id = "agentik",
info = config.agent.toAgentInfo(),
mutableConversationStore = sqliteStores.conversations,
messageStore = sqliteStores.messages,
workingMemoryStore = sqliteStores.workingMemory,
@@ -24,6 +24,7 @@ import pw.binom.agentik.memory.MemoryPrefetcher
import pw.binom.agentik.memory.MemoryReviewer
import pw.binom.agentik.memory.MemorySystemGuidance
import pw.binom.agentik.proto.Agent as ProtoAgent
import pw.binom.agentik.proto.AgentInfo
import pw.binom.agentik.outbox.AgentEvent
import pw.binom.agentik.outbox.Event as OutboxEvent
import pw.binom.agentik.outbox.CommonEvent
@@ -83,6 +84,15 @@ import pw.binom.litert.LiteLlm
*/
internal class ChatAgent(
override val id: String,
/**
* Публичное представление об агенте (имя + описание + чем полезен).
* Доступно через [ProtoAgent.info] — `:server` HTTP-фасад отдаёт его
* как `GET {path}`, A2A — использует для AgentCard. По умолчанию
* имя "agentik", описание/usefulness — отсутствуют (host может
* переопределить через standalone-конфиг `AGENTIK_NAME`,
* `AGENTIK_DESCRIPTION`, `AGENTIK_USEFULNESS`).
*/
override val info: AgentInfo = AgentInfo(name = "agentik"),
private val mutableConversationStore: MutableConversationStore,
private val messageStore: MutableJournalStore,
private val workingMemoryStore: ContextStore,
@@ -170,6 +170,12 @@ class ConversationLoop(
override suspend fun send(content: List<ProtoContent>, context: ProtoMessageContext?) {
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))
val userMessageId = newId("msg")
val storageContext = context?.toStorage()
@@ -7,6 +7,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.mcp.bridge.McpConfig
import pw.binom.agentik.proto.AgentInfo
// Inlined here (AppLimits.kt удалён параллельным рефакторингом): сетевые капы
// и env-parser limits живут рядом с тем, кто их использует, чтобы :config не
@@ -62,7 +63,38 @@ data class AppConfig(
* `null` — авторизация выключена. Источник: env `AGENTIK_A2A_TOKEN`.
*/
val a2aToken: String? = null,
)
/**
* Публичное имя ассистента — то, что видит пользователь/UI/A2A.
* Источник: env `AGENTIK_NAME`. Default — `"agentik"`.
*
* Контрактное отличие от [Agent.id] — id opaque-stable для корреляции,
* `name` — человекочитаемый и может меняться. В терминологии
* [pw.binom.agentik.proto.AgentInfo] это `info.name`.
*/
val name: String = DEFAULT_AGENT_NAME,
/**
* Краткое описание ассистента (что он делает / кто он).
* Источник: env `AGENTIK_DESCRIPTION`. `null` — секция не отдаётся
* клиенту (см. [pw.binom.agentik.proto.AgentInfo.description]).
*/
val description: String? = null,
/**
* Зачем ассистент может быть полезен (в каких задачах помогает).
* Источник: env `AGENTIK_USEFULNESS`. `null` — секция не отдаётся.
*/
val usefulness: String? = null,
) {
/**
* Снапшот [pw.binom.agentik.proto.AgentInfo] из полей секции.
* Имя обязательно ([name] non-blank), остальные опциональные.
*/
fun toAgentInfo(): AgentInfo = AgentInfo(
name = name,
description = description,
usefulness = usefulness,
)
}
/** Долговременная память. */
@Serializable
@@ -146,6 +178,7 @@ data class AppConfig(
companion object {
const val DEFAULT_PORT: Int = 8080
const val DEFAULT_DB_PATH: String = "./agentik.db"
const val DEFAULT_AGENT_NAME: String = "agentik"
const val DEFAULT_COMPRESSION_THRESHOLD: Double = 0.8
const val DEFAULT_EMBEDDING_MODEL: String = "text-embedding-3-small"
const val DEFAULT_EMBEDDING_DIMENSION: Int = 1536
@@ -197,6 +230,9 @@ data class AppConfig(
soulPath = env("AGENTIK_SOUL")?.takeIf { it.isNotBlank() },
authToken = env("AGENTIK_TOKEN")?.takeIf { it.isNotBlank() },
a2aToken = env("AGENTIK_A2A_TOKEN")?.takeIf { it.isNotBlank() },
name = env("AGENTIK_NAME")?.takeIf { it.isNotBlank() } ?: DEFAULT_AGENT_NAME,
description = env("AGENTIK_DESCRIPTION")?.takeIf { it.isNotBlank() },
usefulness = env("AGENTIK_USEFULNESS")?.takeIf { it.isNotBlank() },
),
llm = llm,
mcp = McpConfig.fromEnv(env),
@@ -361,6 +361,36 @@ class ChatAgentTest {
)
}
@Test
fun `Working event is the first event of a turn`() = runTest {
// Working-маркер обязан прийти самым первым событием хода, до
// StartReasoning/StartResponse/AppendText/End. Это позволяет UI
// показать спиннер сразу же при отправке, не дожидаясь первого
// токена от LLM.
val agent = newAgent()
fakeLlm.reply = "ok"
val conv = agent.createConversation(temp = false)
val events = mutableListOf<ProtoEvent>()
val job = launch(start = kotlinx.coroutines.CoroutineStart.UNDISPATCHED) {
agent.outbox.conversationEvents(Instant.DISTANT_PAST, conv.id).collect { events.add(it.event) }
}
conv.send(listOf(Content.Text("hi")))
delay(50)
job.cancel()
// Working — первый event хода (индекс 0).
assertTrue(events.isNotEmpty(), "no events captured: $events")
val first = events.first()
assertIs<ProtoEvent.Working>(first)
assertTrue(events.any { it is ProtoEvent.StartReasoning }, "no StartReasoning: $events")
assertTrue(events.any { it is ProtoEvent.StartResponse }, "no StartResponse: $events")
assertTrue(events.any { it is ProtoEvent.End }, "no End: $events")
// Working не дублируется (emit'ится один раз на send()).
assertEquals(1, events.count { it is ProtoEvent.Working }, "Working emitted >1 times: $events")
}
@Test
fun `interrupt mid-slow-stream preserves user message and no assistant`() = runTest {
// Новая семантика interrupt (commit 7): ставится флаг, LiteConv.cancel()