Compare commits
2 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 4ad59d5f5d | |||
| c140d0b758 |
@@ -41,6 +41,7 @@ kotlin {
|
||||
implementation(libs.kotlinx.cli)
|
||||
|
||||
implementation(libs.kotlinx.coroutines.core)
|
||||
implementation(libs.ktor.client.cio)
|
||||
}
|
||||
// :agentik-cli — commonMain-only (нет jvmMain/nativeMain разделения):
|
||||
// весь код, включая platformEnv, лежит в commonMain.
|
||||
|
||||
@@ -0,0 +1,23 @@
|
||||
package pw.binom.agentik.cli
|
||||
|
||||
import io.ktor.client.HttpClient
|
||||
import io.ktor.client.engine.cio.CIO
|
||||
import pw.binom.agentik.client.agentikHttpClient
|
||||
|
||||
/**
|
||||
* HTTP-клиент CLI: движок CIO + конфигурация agentik.
|
||||
*
|
||||
* Движок живёт здесь, а не в `:client`: библиотека не выбирает транспорт за
|
||||
* потребителя. Таргеты `:agentik-cli` (jvm + linuxX64/macosX64/macosArm64/mingwX64)
|
||||
* покрываются CIO.
|
||||
*
|
||||
* `requestTimeout = 0` — отключение встроенного request-таймаута CIO;
|
||||
* defense-in-depth против обрыва долгих SSE-idle (основная защита —
|
||||
* `noSseReadTimeout` в `:client`).
|
||||
*
|
||||
* [token] = `null` — авторизация выключена.
|
||||
*/
|
||||
internal fun defaultCliHttpClient(token: String? = null): HttpClient =
|
||||
agentikHttpClient(engineFactory = CIO, token = token) {
|
||||
engine { requestTimeout = 0 }
|
||||
}
|
||||
+2
-1
@@ -2,13 +2,14 @@ package pw.binom.agentik.cli.commands
|
||||
|
||||
import kotlinx.cli.ArgType
|
||||
import pw.binom.agentik.cli.AgentikSubcommand
|
||||
import pw.binom.agentik.cli.defaultCliHttpClient
|
||||
import pw.binom.agentik.client.AgentikAgent
|
||||
|
||||
class ConvDeleteSubcommand : ConvSubcommand("delete", "Удалить диалог") {
|
||||
val id by argument(ArgType.String, description = "ID диалога")
|
||||
|
||||
override fun execute() = kotlinx.coroutines.runBlocking {
|
||||
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl)
|
||||
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl, httpClient = defaultCliHttpClient())
|
||||
val ok = agent.deleteConversation(id)
|
||||
if (ok) println("deleted: $id") else println("conversation not found: $id")
|
||||
}
|
||||
|
||||
@@ -3,6 +3,7 @@ package pw.binom.agentik.cli.commands
|
||||
import kotlinx.cli.ArgType
|
||||
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.Agent
|
||||
|
||||
@@ -10,7 +11,7 @@ class ConvLsSubcommand : ConvSubcommand("ls", "Список диалогов а
|
||||
val limit by option(ArgType.Int, fullName = "limit", description = "Максимум диалогов").default(Agent.PAGE_SIZE)
|
||||
|
||||
override fun execute() = kotlinx.coroutines.runBlocking {
|
||||
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl)
|
||||
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl, httpClient = defaultCliHttpClient())
|
||||
val convs = agent.getConversations(offset = 0, limit = limit.coerceAtMost(Agent.PAGE_SIZE))
|
||||
if (convs.isEmpty()) {
|
||||
println("(no conversations)")
|
||||
|
||||
+2
-1
@@ -3,13 +3,14 @@ package pw.binom.agentik.cli.commands
|
||||
import kotlinx.cli.ArgType
|
||||
import kotlinx.cli.default
|
||||
import pw.binom.agentik.cli.AgentikSubcommand
|
||||
import pw.binom.agentik.cli.defaultCliHttpClient
|
||||
import pw.binom.agentik.client.AgentikAgent
|
||||
|
||||
class ConvNewSubcommand : ConvSubcommand("new", "Создать диалог; печатает id") {
|
||||
val temp by option(ArgType.Boolean, fullName = "temp", description = "Временный диалог").default(false)
|
||||
|
||||
override fun execute() = kotlinx.coroutines.runBlocking {
|
||||
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl)
|
||||
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl, httpClient = defaultCliHttpClient())
|
||||
val conv = agent.createConversation(temp = temp)
|
||||
println(conv.id)
|
||||
}
|
||||
|
||||
+2
-1
@@ -2,6 +2,7 @@ package pw.binom.agentik.cli.commands
|
||||
|
||||
import kotlinx.cli.ArgType
|
||||
import pw.binom.agentik.cli.AgentikSubcommand
|
||||
import pw.binom.agentik.cli.defaultCliHttpClient
|
||||
import pw.binom.agentik.client.AgentikAgent
|
||||
|
||||
class ConvRenameSubcommand : ConvSubcommand("rename", "Переименовать диалог") {
|
||||
@@ -9,7 +10,7 @@ class ConvRenameSubcommand : ConvSubcommand("rename", "Переименоват
|
||||
val title by argument(ArgType.String, description = "Новое название")
|
||||
|
||||
override fun execute() = kotlinx.coroutines.runBlocking {
|
||||
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl)
|
||||
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl, httpClient = defaultCliHttpClient())
|
||||
val conv = agent.getConversation(id) ?: run {
|
||||
println("conversation not found: $id")
|
||||
return@runBlocking
|
||||
|
||||
+2
-1
@@ -2,13 +2,14 @@ package pw.binom.agentik.cli.commands
|
||||
|
||||
import kotlinx.cli.ArgType
|
||||
import pw.binom.agentik.cli.AgentikSubcommand
|
||||
import pw.binom.agentik.cli.defaultCliHttpClient
|
||||
import pw.binom.agentik.client.AgentikAgent
|
||||
|
||||
class ConvShowSubcommand : ConvSubcommand("show", "Метаданные диалога") {
|
||||
val id by argument(ArgType.String, description = "ID диалога")
|
||||
|
||||
override fun execute() = kotlinx.coroutines.runBlocking {
|
||||
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl)
|
||||
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl, httpClient = defaultCliHttpClient())
|
||||
val conv = agent.getConversation(id) ?: run {
|
||||
println("conversation not found: $id")
|
||||
return@runBlocking
|
||||
|
||||
+2
-1
@@ -2,13 +2,14 @@ package pw.binom.agentik.cli.commands
|
||||
|
||||
import kotlinx.cli.ArgType
|
||||
import pw.binom.agentik.cli.AgentikSubcommand
|
||||
import pw.binom.agentik.cli.defaultCliHttpClient
|
||||
import pw.binom.agentik.client.AgentikAgent
|
||||
|
||||
class InterruptSubcommand : AgentikSubcommand("interrupt", "Прервать текущий ход диалога") {
|
||||
val id by argument(ArgType.String, description = "ID диалога")
|
||||
|
||||
override fun execute() = kotlinx.coroutines.runBlocking {
|
||||
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl)
|
||||
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl, httpClient = defaultCliHttpClient())
|
||||
val conv = agent.getConversation(id) ?: run {
|
||||
println("conversation not found: $id")
|
||||
return@runBlocking
|
||||
|
||||
@@ -3,6 +3,7 @@ package pw.binom.agentik.cli.commands
|
||||
import kotlinx.cli.ArgType
|
||||
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.proto.Message
|
||||
@@ -13,7 +14,7 @@ class MsgsSubcommand : AgentikSubcommand("msgs", "Показать сообще
|
||||
val limit by option(ArgType.Int, fullName = "limit", description = "Максимум сообщений").default(100)
|
||||
|
||||
override fun execute() = kotlinx.coroutines.runBlocking {
|
||||
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl)
|
||||
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl, httpClient = defaultCliHttpClient())
|
||||
val conv = agent.getConversation(id) ?: run {
|
||||
println("conversation not found: $id")
|
||||
return@runBlocking
|
||||
|
||||
@@ -7,6 +7,7 @@ import kotlinx.coroutines.flow.onEach
|
||||
import kotlinx.coroutines.flow.takeWhile
|
||||
import kotlinx.coroutines.launch
|
||||
import pw.binom.agentik.cli.AgentikSubcommand
|
||||
import pw.binom.agentik.cli.defaultCliHttpClient
|
||||
import pw.binom.agentik.client.AgentikAgent
|
||||
import pw.binom.agentik.proto.Content
|
||||
import pw.binom.agentik.proto.Event
|
||||
@@ -17,7 +18,7 @@ class SendSubcommand : AgentikSubcommand("send", "Отправить user-ход
|
||||
val text by argument(ArgType.String, description = "Текст хода (все позиционные после <id> склеиваются пробелом)").vararg()
|
||||
|
||||
override fun execute() = kotlinx.coroutines.runBlocking {
|
||||
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl)
|
||||
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl, httpClient = defaultCliHttpClient())
|
||||
val conv = agent.getConversation(id) ?: run {
|
||||
println("conversation not found: $id")
|
||||
return@runBlocking
|
||||
|
||||
@@ -22,8 +22,7 @@ kotlin {
|
||||
commonMain.dependencies {
|
||||
api(project(":proto"))
|
||||
|
||||
implementation(libs.ktor.client.core)
|
||||
implementation(libs.ktor.client.cio)
|
||||
api(libs.ktor.client.core)
|
||||
implementation(libs.ktor.client.content.negotiation)
|
||||
implementation(libs.ktor.serialization.kotlinx.json)
|
||||
|
||||
@@ -37,6 +36,7 @@ kotlin {
|
||||
implementation(libs.ktor.server.core)
|
||||
implementation(libs.ktor.server.test.host)
|
||||
implementation(libs.ktor.client.content.negotiation)
|
||||
implementation(libs.ktor.client.cio)
|
||||
implementation(libs.ktor.server.cio)
|
||||
implementation(libs.ktor.server.sse)
|
||||
}
|
||||
|
||||
@@ -17,6 +17,7 @@ import kotlinx.coroutines.flow.flow
|
||||
import kotlinx.coroutines.runBlocking
|
||||
import pw.binom.agentik.proto.Agent
|
||||
import pw.binom.agentik.proto.AgentEvent
|
||||
import pw.binom.agentik.proto.AllEvent
|
||||
import pw.binom.agentik.proto.Conversation
|
||||
import kotlin.time.Instant
|
||||
|
||||
@@ -78,4 +79,74 @@ internal class AgentClient(
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
/**
|
||||
* Подписка на ВСЕ события: agent lifecycle + все conversation events.
|
||||
* Использует SSE endpoint /events/all.
|
||||
*/
|
||||
override fun allEvents(after: Instant): Flow<AllEvent> = flow {
|
||||
httpClient.prepareGet("$agentUrl/events/all?after=$after") { noSseReadTimeout() }
|
||||
.execute { response ->
|
||||
check(response.status == HttpStatusCode.OK) {
|
||||
"allEvents: server returned ${response.status}"
|
||||
}
|
||||
readSse(response.bodyAsChannel())
|
||||
.collect { payload ->
|
||||
emit(agentikJson.decodeFromString(AllEvent.serializer(), payload))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Catchup для /events/replay — пагинированно читает events после [afterId].
|
||||
* Caller делает несколько вызовов пока `result.size < limit` (= конец).
|
||||
*
|
||||
* @param afterId exclusive cursor. `null` = с начала.
|
||||
* @param limit max per-request (default 100, max 1000 на сервере).
|
||||
*/
|
||||
internal suspend fun replayAllEvents(
|
||||
afterId: String? = null,
|
||||
limit: Int = 100,
|
||||
): List<EventRecordDto> {
|
||||
val response = httpClient.get("$agentUrl/events/replay") {
|
||||
afterId?.let { parameter("after_id", it) }
|
||||
parameter("limit", limit)
|
||||
}
|
||||
return response.body()
|
||||
}
|
||||
|
||||
/**
|
||||
* Catchup для /conversations/{id}/events/replay — пагинированно.
|
||||
*/
|
||||
internal suspend fun replayConversationEvents(
|
||||
conversationId: String,
|
||||
afterId: String? = null,
|
||||
limit: Int = 100,
|
||||
): List<EventRecordDto> {
|
||||
val response = httpClient.get("$agentUrl/conversations/$conversationId/events/replay") {
|
||||
afterId?.let { parameter("after_id", it) }
|
||||
parameter("limit", limit)
|
||||
}
|
||||
return response.body()
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
/**
|
||||
* Запись event'а, которую возвращает /events/replay endpoint.
|
||||
*
|
||||
* Это **мини-зеркало** `pw.binom.agentik.storage.events.EventRecord` — клиент
|
||||
* не зависит от `:storage-core`, поэтому определяет свою модель (wire-only).
|
||||
*
|
||||
* Формат полностью совместим с сервером — `agentikJson.encodeToString(...)` там
|
||||
* и `agentikJson.decodeFromString(...)` здесь.
|
||||
*/
|
||||
@kotlinx.serialization.Serializable
|
||||
internal data class EventRecordDto(
|
||||
val id: String,
|
||||
val conversationId: String? = null,
|
||||
val createdAt: kotlin.time.Instant,
|
||||
val type: String,
|
||||
val payload: String,
|
||||
)
|
||||
|
||||
@@ -8,9 +8,11 @@ import pw.binom.agentik.proto.Agent
|
||||
* (модуль `:server`).
|
||||
*
|
||||
* ```
|
||||
* val http = HttpClient(CIO) { applyAgentikDefaults(token = "s3cret") }
|
||||
* val client = AgentikAgent(
|
||||
* id = "my-agent",
|
||||
* baseUrl = "http://localhost:8080/agentik",
|
||||
* httpClient = http,
|
||||
* )
|
||||
* val conv = client.createConversation(temp = false)
|
||||
* conv.send(listOf(Content.Text("hi")))
|
||||
@@ -21,27 +23,13 @@ import pw.binom.agentik.proto.Agent
|
||||
* агента не знает, поэтому клиент должен её знать сам (или взять из
|
||||
* конфига).
|
||||
*
|
||||
* [httpClient] по умолчанию — [defaultAgentikHttpClient] (платформо-зависимый
|
||||
* движок: CIO на JVM, libcurl на desktop-native). Можно передать свой.
|
||||
* Клиент приходит снаружи: `:client` не выбирает движок. Собрать [HttpClient]
|
||||
* можно через [agentikHttpClient] (фабрика движка + опциональные движковые
|
||||
* настройки) или вручную, применив к блоку конфигурации [applyAgentikDefaults]
|
||||
* (JSON + опциональный Bearer-токен).
|
||||
*/
|
||||
fun AgentikAgent(
|
||||
id: String,
|
||||
baseUrl: String,
|
||||
token: String? = null,
|
||||
httpClient: HttpClient = defaultAgentikHttpClient(token),
|
||||
httpClient: HttpClient,
|
||||
): Agent = AgentClient(httpClient = httpClient, baseUrl = baseUrl, id = id)
|
||||
|
||||
/**
|
||||
* Дефолтный [HttpClient] для общения с `agentikAgent`. SSE-парсер ([readSse])
|
||||
* живёт в общем коде и плагина `SSEClientContent` не требует.
|
||||
*
|
||||
* **Платформы:**
|
||||
* - JVM: движок CIO. `engine { requestTimeout = 0 }` отключает встроенный
|
||||
* 15-секундный request-таймаут движка (наш кастомный SSE-ридер не маркирует
|
||||
* для долгих idle-стримов). Defense-in-depth: SSE-запросы в
|
||||
* `ConversationClient.events`/`AgentClient.events` уже ставят
|
||||
* `HttpTimeoutCapability` = INFINITE (см. [noSseReadTimeout]).
|
||||
*
|
||||
* Один движок CIO работает и на JVM, и на всех desktop-native (linux/macos/mingw).
|
||||
* Реализация — в [HttpClientFactory.kt].
|
||||
*/
|
||||
|
||||
@@ -1,7 +1,9 @@
|
||||
package pw.binom.agentik.client
|
||||
|
||||
import io.ktor.client.HttpClient
|
||||
import io.ktor.client.engine.cio.CIO
|
||||
import io.ktor.client.HttpClientConfig
|
||||
import io.ktor.client.engine.HttpClientEngineConfig
|
||||
import io.ktor.client.engine.HttpClientEngineFactory
|
||||
import io.ktor.client.plugins.DefaultRequest
|
||||
import io.ktor.client.plugins.contentnegotiation.ContentNegotiation
|
||||
import io.ktor.client.request.header
|
||||
@@ -9,19 +11,26 @@ import io.ktor.http.HttpHeaders
|
||||
import io.ktor.serialization.kotlinx.json.json
|
||||
|
||||
/**
|
||||
* Единый HTTP-клиент для JVM и всех 5 native-таргетов (:agentik-cli).
|
||||
* CIO в ktor 3.x — KMP, поддерживает linuxX64/Arm64, macosX64/Arm64, mingwX64.
|
||||
* Общая конфигурация HTTP-клиента agentik — платформо-независимая часть.
|
||||
*
|
||||
* `requestTimeout = 0` — defense-in-depth против read-таймаута на SSE:
|
||||
* основная защита в `HttpRequestBuilder.noSseReadTimeout()` ([SseTimeout]).
|
||||
* `:client` НЕ выбирает движок: его приносит потребитель. Здесь живёт только то,
|
||||
* без чего клиент несовместим с `/agentik`:
|
||||
* - JSON-конфиг [agentikJson] (обязан совпадать с серверным);
|
||||
* - при заданном [token] — `Authorization: Bearer <token>` на ВСЕ запросы
|
||||
* через [DefaultRequest] (накрывает 10 REST-вызовов и оба SSE-потока;
|
||||
* заголовок живёт на клиенте, а не в отдельных запросах).
|
||||
*
|
||||
* При заданном [token] на ВСЕ запросы клиента навешивается
|
||||
* `Authorization: Bearer <token>` через плагин [DefaultRequest]. Это накрывает
|
||||
* все 10 REST-вызовов и оба SSE-потока сразу — заголовок живёт на HTTP-клиенте,
|
||||
* а не в отдельных запросах.
|
||||
* `null` — авторизация выключена, заголовок не отправляется.
|
||||
*
|
||||
* Потребитель, знающий свой движок, добавляет к этому движковые настройки, напр.:
|
||||
* ```
|
||||
* val http = HttpClient(CIO) {
|
||||
* engine { requestTimeout = 0 } // CIO-специфика, живёт у потребителя
|
||||
* applyAgentikDefaults(token)
|
||||
* }
|
||||
* ```
|
||||
*/
|
||||
fun defaultAgentikHttpClient(token: String? = null): HttpClient = HttpClient(CIO) {
|
||||
engine { requestTimeout = 0 }
|
||||
fun HttpClientConfig<*>.applyAgentikDefaults(token: String? = null) {
|
||||
install(ContentNegotiation) { json(agentikJson) }
|
||||
if (token != null) {
|
||||
install(DefaultRequest) {
|
||||
@@ -29,3 +38,23 @@ fun defaultAgentikHttpClient(token: String? = null): HttpClient = HttpClient(CIO
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Создаёт [HttpClient] из фабрики движка потребителя и сразу применяет к нему
|
||||
* конфигурацию agentik ([applyAgentikDefaults]).
|
||||
*
|
||||
* Это точка, где `:client` НЕ привязан к реализации транспорта: [engineFactory]
|
||||
* выбирает потребитель (CIO, OkHttp, Darwin, …), а `:client` только конфигурирует
|
||||
* созданный клиент.
|
||||
*
|
||||
* [configure] — опциональный последний штрих потребителя (движковые настройки:
|
||||
* таймауты, прокси, логирование). Вызывается ПОСЛЕ [applyAgentikDefaults].
|
||||
*/
|
||||
fun <T : HttpClientEngineConfig> agentikHttpClient(
|
||||
engineFactory: HttpClientEngineFactory<T>,
|
||||
token: String? = null,
|
||||
configure: (HttpClientConfig<T>.() -> Unit)? = null,
|
||||
): HttpClient = HttpClient(engineFactory) {
|
||||
applyAgentikDefaults(token)
|
||||
configure?.invoke(this)
|
||||
}
|
||||
@@ -1,5 +1,7 @@
|
||||
package pw.binom.agentik.client
|
||||
|
||||
import io.ktor.client.HttpClient
|
||||
import io.ktor.client.engine.cio.CIO
|
||||
import io.ktor.client.request.get
|
||||
import io.ktor.client.statement.bodyAsText
|
||||
import io.ktor.http.ContentType
|
||||
@@ -19,7 +21,7 @@ import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
|
||||
/**
|
||||
* Тесты клиентской части: [defaultAgentikHttpClient] с заданным `token` прикладывает
|
||||
* Тесты клиентской части: [applyAgentikDefaults] с заданным `token` прикладывает
|
||||
* `Authorization: Bearer <token>` ко всем запросам через плагин `DefaultRequest`,
|
||||
* без токена — заголовок не отправляется.
|
||||
*
|
||||
@@ -47,6 +49,9 @@ class BearerHeaderTest {
|
||||
var token: String? = null
|
||||
}
|
||||
|
||||
private fun clientWith(token: String?): HttpClient =
|
||||
HttpClient(CIO) { applyAgentikDefaults(token) }
|
||||
|
||||
private suspend fun startServer(): Pair<EmbeddedServer<*, *>, Int> {
|
||||
val server = embeddedServer(ServerCIO, port = 0) {
|
||||
routing {
|
||||
@@ -66,7 +71,7 @@ class BearerHeaderTest {
|
||||
fun clientWithTokenAttachesBearerHeader() = runBlocking {
|
||||
val (server, port) = startServer()
|
||||
try {
|
||||
val client = defaultAgentikHttpClient("secret")
|
||||
val client = clientWith("secret")
|
||||
val resp = client.get("http://127.0.0.1:$port/agentik/conversations")
|
||||
assertEquals(HttpStatusCode.OK, resp.status)
|
||||
assertEquals("[]", resp.bodyAsText())
|
||||
@@ -79,7 +84,7 @@ class BearerHeaderTest {
|
||||
fun clientWithoutTokenGets401(): Unit = runBlocking {
|
||||
val (server, port) = startServer()
|
||||
try {
|
||||
val client = defaultAgentikHttpClient(null)
|
||||
val client = clientWith(null)
|
||||
val resp = client.get("http://127.0.0.1:$port/agentik/conversations")
|
||||
assertEquals(HttpStatusCode.Unauthorized, resp.status)
|
||||
} finally {
|
||||
@@ -91,7 +96,7 @@ class BearerHeaderTest {
|
||||
fun clientWithWrongTokenGets401(): Unit = runBlocking {
|
||||
val (server, port) = startServer()
|
||||
try {
|
||||
val client = defaultAgentikHttpClient("wrong")
|
||||
val client = clientWith("wrong")
|
||||
val resp = client.get("http://127.0.0.1:$port/agentik/conversations")
|
||||
assertEquals(HttpStatusCode.Unauthorized, resp.status)
|
||||
} finally {
|
||||
|
||||
@@ -49,6 +49,18 @@ public interface Agent {
|
||||
*/
|
||||
fun events(after: Instant): Flow<AgentEvent>
|
||||
|
||||
/**
|
||||
* All events in one stream: agent lifecycle (Created/Deleted/Renamed) +
|
||||
* all conversation turns. Useful for admin dashboards, debug tools,
|
||||
* parent agents.
|
||||
*
|
||||
* For UI use [events] + [Conversation.events]. This one-feed variant is
|
||||
* for cases where everything-in-one is preferred.
|
||||
*
|
||||
* Cold (no replay). For catchup use EventStore.
|
||||
*/
|
||||
fun allEvents(after: Instant): Flow<AllEvent>
|
||||
|
||||
companion object {
|
||||
|
||||
const val PAGE_SIZE: Int = 100
|
||||
|
||||
@@ -0,0 +1,39 @@
|
||||
package pw.binom.agentik.proto
|
||||
|
||||
import kotlin.time.Instant
|
||||
import kotlinx.serialization.SerialName
|
||||
import kotlinx.serialization.Serializable
|
||||
|
||||
/**
|
||||
* Unified wrapper for all agent events in a single stream.
|
||||
*
|
||||
* Useful for admin dashboards, debug tools, parent agents: one subscription
|
||||
* instead of N+1. For regular UI use two separate SSE feeds
|
||||
* ([AgentEvent] via /events and [Event] via /conversations/{id}/events);
|
||||
* [AllEvent] is for those who need everything in one place.
|
||||
*
|
||||
* Server endpoint: GET /events/all (SSE), or replay via EventStore.
|
||||
*
|
||||
* Not used for persistence payload: EventStore stores AgentEvent and
|
||||
* Conversation.Event natively (compact form); this wrapper is wire-format
|
||||
* only.
|
||||
*/
|
||||
@Serializable
|
||||
sealed interface AllEvent {
|
||||
val date: Instant
|
||||
|
||||
@Serializable
|
||||
@SerialName("agent")
|
||||
data class Agent(
|
||||
override val date: Instant,
|
||||
val event: AgentEvent,
|
||||
) : AllEvent
|
||||
|
||||
@Serializable
|
||||
@SerialName("conversation")
|
||||
data class Conversation(
|
||||
override val date: Instant,
|
||||
val conversationId: String,
|
||||
val event: Event,
|
||||
) : AllEvent
|
||||
}
|
||||
@@ -23,6 +23,7 @@ kotlin {
|
||||
sourceSets {
|
||||
commonMain.dependencies {
|
||||
implementation(project(":proto"))
|
||||
implementation(project(":storage-core"))
|
||||
|
||||
// Ktor (без engine — engine подключает потребитель, см. :standalone).
|
||||
implementation(libs.ktor.server.core)
|
||||
|
||||
@@ -6,6 +6,7 @@ import io.ktor.server.plugins.contentnegotiation.ContentNegotiation
|
||||
import io.ktor.server.routing.Route
|
||||
import io.ktor.server.routing.route
|
||||
import pw.binom.agentik.proto.Agent
|
||||
import pw.binom.agentik.storage.events.EventStore
|
||||
|
||||
/**
|
||||
* Встраивает HTTP/SSE-фасад протокола agentik в твой Ktor-роутинг.
|
||||
@@ -30,10 +31,20 @@ import pw.binom.agentik.proto.Agent
|
||||
* - `POST /conversations/{id}/interrupt` — `interrupt`
|
||||
* - `GET /conversations/{id}/messages` — история
|
||||
* - `GET /conversations/{id}/events` — SSE: события хода
|
||||
* - `GET /conversations/{id}/events/replay` — replay-after-disconnect (events after ?after_id=X)
|
||||
* - `GET /events` — SSE: события агента
|
||||
* - `GET /events/replay` — replay-after-disconnect (global)
|
||||
* - `GET /health` — `"ok"`
|
||||
*
|
||||
* @param eventStore если null — `/events/replay` endpoints возвращают 503 (event
|
||||
* persistence не настроен). Live-streaming работает как обычно.
|
||||
*/
|
||||
fun Route.agentikAgent(agent: Agent, path: String = "/agentik", token: String? = null) {
|
||||
fun Route.agentikAgent(
|
||||
agent: Agent,
|
||||
eventStore: EventStore? = null,
|
||||
path: String = "/agentik",
|
||||
token: String? = null,
|
||||
) {
|
||||
route(path) {
|
||||
install(ContentNegotiation) {
|
||||
json(agentikJson)
|
||||
@@ -43,6 +54,6 @@ fun Route.agentikAgent(agent: Agent, path: String = "/agentik", token: String? =
|
||||
this.token = token
|
||||
}
|
||||
}
|
||||
agentikRoutes(agent)
|
||||
agentikRoutes(agent, eventStore)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -21,12 +21,14 @@ import kotlinx.serialization.KSerializer
|
||||
import kotlinx.serialization.json.Json
|
||||
import pw.binom.agentik.proto.Agent
|
||||
import pw.binom.agentik.proto.AgentEvent
|
||||
import pw.binom.agentik.proto.AllEvent
|
||||
import pw.binom.agentik.proto.Conversation
|
||||
import pw.binom.agentik.proto.Content
|
||||
import pw.binom.agentik.proto.Event
|
||||
import pw.binom.agentik.storage.events.EventStore
|
||||
import kotlin.time.Instant
|
||||
|
||||
internal fun Route.agentikRoutes(agent: Agent) {
|
||||
internal fun Route.agentikRoutes(agent: Agent, eventStore: EventStore? = null) {
|
||||
|
||||
get("/health") {
|
||||
call.respondText("ok")
|
||||
@@ -134,6 +136,48 @@ internal fun Route.agentikRoutes(agent: Agent) {
|
||||
val after = call.parseAfter() ?: return@get
|
||||
call.streamJsonSse(agent.events(after), AgentEvent.serializer())
|
||||
}
|
||||
|
||||
/**
|
||||
* Все события в одном потоке: agent lifecycle + все conversation events.
|
||||
* Для admin-дашборда, debug-инструментов, parent-агента.
|
||||
* Для UI достаточно `/events` + `/conversations/{id}/events`.
|
||||
*/
|
||||
get("/events/all") {
|
||||
val after = call.parseAfter() ?: return@get
|
||||
call.streamJsonSse(agent.allEvents(after), AllEvent.serializer())
|
||||
}
|
||||
|
||||
// ---- Replay-after-disconnect (event log) ----
|
||||
//
|
||||
// Stream `/events` и `/conversations/{id}/events` — cold (no replay).
|
||||
// Клиент, отвалившийся от SSE, при reconnect делает GET на replay-endpoint
|
||||
// с `?after_id=X` чтобы получить events, которые произошли во время разрыва.
|
||||
// paginated: делает несколько запросов пока `result.size < limit`.
|
||||
//
|
||||
// Если [eventStore] == null (не сконфигурирован), эти endpoints возвращают 503.
|
||||
|
||||
get("/events/replay") {
|
||||
if (eventStore == null) {
|
||||
call.respond(HttpStatusCode.ServiceUnavailable, "EventStore not configured on this agent")
|
||||
return@get
|
||||
}
|
||||
val afterId = call.request.queryParameters["after_id"]?.takeIf { it.isNotBlank() }
|
||||
val limit = (call.request.queryParameters["limit"]?.toIntOrNull() ?: 100).coerceIn(1, 1000)
|
||||
val records = eventStore.query(conversationId = null, afterId = afterId, limit = limit)
|
||||
call.respond(records)
|
||||
}
|
||||
|
||||
get("/conversations/{id}/events/replay") {
|
||||
if (eventStore == null) {
|
||||
call.respond(HttpStatusCode.ServiceUnavailable, "EventStore not configured on this agent")
|
||||
return@get
|
||||
}
|
||||
val id = call.parameters["id"]!!
|
||||
val afterId = call.request.queryParameters["after_id"]?.takeIf { it.isNotBlank() }
|
||||
val limit = (call.request.queryParameters["limit"]?.toIntOrNull() ?: 100).coerceIn(1, 1000)
|
||||
val records = eventStore.query(conversationId = id, afterId = afterId, limit = limit)
|
||||
call.respond(records)
|
||||
}
|
||||
}
|
||||
|
||||
// ---------- helpers ----------
|
||||
|
||||
@@ -1,38 +1,50 @@
|
||||
package pw.binom.agentik.standalone.agent
|
||||
|
||||
import kotlin.time.Instant
|
||||
import kotlinx.coroutines.flow.Flow
|
||||
import kotlinx.coroutines.flow.MutableSharedFlow
|
||||
import kotlinx.coroutines.flow.asSharedFlow
|
||||
import kotlinx.coroutines.flow.emitAll
|
||||
import kotlinx.coroutines.flow.filterNotNull
|
||||
import kotlinx.coroutines.flow.flow
|
||||
import kotlinx.coroutines.flow.map
|
||||
import kotlinx.coroutines.flow.merge
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlinx.coroutines.runBlocking
|
||||
import kotlinx.coroutines.sync.Mutex
|
||||
import kotlinx.coroutines.sync.withLock
|
||||
import kotlinx.serialization.json.Json
|
||||
import pw.binom.agentik.llm.tools.ContextCompactor
|
||||
import pw.binom.agentik.llm.tools.LlmMemoryReviewer
|
||||
import pw.binom.agentik.llm.tools.LlmReflector
|
||||
import pw.binom.agentik.llm.tools.SkillMiner
|
||||
import pw.binom.agentik.memory.MemoryPrefetcher
|
||||
import pw.binom.agentik.memory.MemoryReviewer
|
||||
import pw.binom.agentik.memory.MemorySystemGuidance
|
||||
import pw.binom.agentik.proto.Agent as ProtoAgent
|
||||
import pw.binom.agentik.proto.AgentEvent
|
||||
import pw.binom.agentik.proto.AllEvent
|
||||
import pw.binom.agentik.proto.Conversation as ProtoConversation
|
||||
import pw.binom.agentik.proto.Event as ProtoEvent
|
||||
import pw.binom.agentik.skills.SkillCatalog
|
||||
import pw.binom.agentik.skills.renderSystemPromptSection
|
||||
import pw.binom.agentik.standalone.agent.memory.MemoryToolsFactory
|
||||
import pw.binom.agentik.standalone.llm.LlmConfig
|
||||
import pw.binom.agentik.storage.ConversationRecord
|
||||
import pw.binom.agentik.storage.Ids
|
||||
import pw.binom.agentik.storage.Reflection
|
||||
import pw.binom.agentik.storage.WorkingMemoryEntry
|
||||
import pw.binom.agentik.storage.StorageBundle
|
||||
import pw.binom.agentik.toolsets.EnableToolsetTool
|
||||
import pw.binom.agentik.storage.WorkingMemoryEntry
|
||||
import pw.binom.agentik.storage.events.EventRecord
|
||||
import pw.binom.agentik.storage.events.EventType
|
||||
import pw.binom.agentik.toolsets.DisableToolsetTool
|
||||
import pw.binom.agentik.toolsets.EnableToolsetTool
|
||||
import pw.binom.agentik.toolsets.NamedTool
|
||||
import pw.binom.agentik.toolsets.SystemPromptToolsetSection
|
||||
import pw.binom.agentik.toolsets.ToolsetContribution
|
||||
import pw.binom.agentik.toolsets.ToolsetDispatchPolicy
|
||||
import pw.binom.agentik.toolsets.ToolsetRegistry
|
||||
import pw.binom.litert.LiteLlm
|
||||
import kotlin.time.Instant
|
||||
import pw.binom.agentik.llm.tools.LlmReflector
|
||||
import pw.binom.agentik.llm.tools.SkillMiner
|
||||
import pw.binom.agentik.llm.tools.LlmMemoryReviewer
|
||||
import pw.binom.agentik.llm.tools.ContextCompactor
|
||||
import pw.binom.agentik.toolsets.NamedTool
|
||||
|
||||
/**
|
||||
* Stateful [ProtoAgent] на базе SQLite (история + working memory) и
|
||||
@@ -192,6 +204,54 @@ class ChatAgent(
|
||||
extraBufferCapacity = 64,
|
||||
)
|
||||
|
||||
/** Json-encoder для payload в EventStore. Один на весь agent. */
|
||||
private val eventJson = Json {
|
||||
ignoreUnknownKeys = true
|
||||
encodeDefaults = true
|
||||
}
|
||||
|
||||
/**
|
||||
* Fire-and-forget persist в EventStore. Если EventStore настроен (production
|
||||
* deployment с `:storage-sqlite`), каждый AgentEvent также уходит в SQLite
|
||||
* с уникальным id, чтобы клиенты могли сделать /events/replay после
|
||||
* disconnect. Если EventStore == null (например, dev mode с in-memory
|
||||
* storage или Android без persistent event log) — no-op.
|
||||
*
|
||||
* Background launch + runCatching: ошибки БД не должны ронять agent loop.
|
||||
* Логирование — если persistence падает, видно в логах, но live-stream
|
||||
* продолжает работать.
|
||||
*/
|
||||
private fun persistAgentEvent(event: AgentEvent) {
|
||||
val store = storage.eventStore ?: return
|
||||
kotlinx.coroutines.CoroutineScope(
|
||||
kotlinx.coroutines.SupervisorJob() + kotlinx.coroutines.Dispatchers.IO
|
||||
).launch {
|
||||
runCatching {
|
||||
store.append(
|
||||
EventRecord(
|
||||
id = "ev-${Ids.new("agent")}",
|
||||
conversationId = when (event) {
|
||||
is AgentEvent.Created -> event.conversationId
|
||||
is AgentEvent.Deleted -> event.id
|
||||
is AgentEvent.Renamed -> event.id
|
||||
},
|
||||
createdAt = event.date,
|
||||
type = when (event) {
|
||||
is AgentEvent.Created -> EventType.AGENT_CREATED
|
||||
is AgentEvent.Deleted -> EventType.AGENT_DELETED
|
||||
is AgentEvent.Renamed -> EventType.AGENT_RENAMED
|
||||
},
|
||||
payload = eventJson.encodeToString(AgentEvent.serializer(), event),
|
||||
)
|
||||
)
|
||||
}.onFailure {
|
||||
mu.KotlinLogging.logger("ChatAgent").warn(it) {
|
||||
"failed to persist agent event ${event::class.simpleName}: ${it.message}"
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/** Защищает карту живых диалогов. */
|
||||
private val liveLock = Mutex()
|
||||
private val live: MutableMap<String, ChatConversation> = HashMap()
|
||||
@@ -204,6 +264,36 @@ class ChatAgent(
|
||||
return agentEvents.asSharedFlow()
|
||||
}
|
||||
|
||||
/**
|
||||
* Все события в одном потоке: agent lifecycle + events всех диалогов.
|
||||
* Реализация — merge двух cold-flow'ов. Snapshot-список живых диалогов
|
||||
* берётся на момент подписки; новые Created-Event'ы НЕ переподписывают
|
||||
* (это ответственность caller'а: если хочет всё — он может
|
||||
* переподписаться или следить за AgentEvent.Created сам).
|
||||
*
|
||||
* Для admin/debug — допустимое упрощение. Для long-running мониторинга
|
||||
* (admin-дашборд часами) надо добавить reactive re-subscribe (см.
|
||||
* notes 04-sub-agents.md, Variant B).
|
||||
*/
|
||||
override fun allEvents(after: Instant): Flow<AllEvent> = flow {
|
||||
// Agent lifecycle events
|
||||
emitAll(agentEvents.asSharedFlow().map { e: AgentEvent ->
|
||||
AllEvent.Agent(date = e.date, event = e)
|
||||
})
|
||||
// Snapshot живых диалогов на момент подписки.
|
||||
// ВАЖНО: не подписываемся на новые Created — это ответственность caller'а
|
||||
// (см. KDoc выше).
|
||||
live.values
|
||||
.asSequence()
|
||||
.filterNot { it.isClosed }
|
||||
.forEach { conv: ProtoConversation ->
|
||||
val cid: String = conv.id
|
||||
emitAll(conv.events(after).map { ev: ProtoEvent ->
|
||||
AllEvent.Conversation(date = ev.date, conversationId = cid, event = ev)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
override fun createConversation(temp: Boolean): ProtoConversation {
|
||||
val now = now()
|
||||
val id = pw.binom.agentik.storage.Ids.new("conv")
|
||||
@@ -243,6 +333,7 @@ class ChatAgent(
|
||||
liveLock.withLock { live[conv.id] = conv }
|
||||
}
|
||||
agentEvents.tryEmit(AgentEvent.Created(date = now(), conversationId = conv.id))
|
||||
persistAgentEvent(AgentEvent.Created(date = now(), conversationId = conv.id))
|
||||
return conv
|
||||
}
|
||||
|
||||
@@ -258,7 +349,11 @@ class ChatAgent(
|
||||
val conv = liveLock.withLock { live.remove(id) }
|
||||
conv?.close()
|
||||
val ok = storage.conversationStore.delete(id)
|
||||
if (ok) agentEvents.tryEmit(AgentEvent.Deleted(date = now(), id = id))
|
||||
if (ok) {
|
||||
val event = AgentEvent.Deleted(date = now(), id = id)
|
||||
agentEvents.tryEmit(event)
|
||||
persistAgentEvent(event)
|
||||
}
|
||||
return ok
|
||||
}
|
||||
|
||||
|
||||
+103
-4
@@ -4,9 +4,45 @@ import kotlinx.coroutines.channels.BufferOverflow
|
||||
import kotlinx.coroutines.flow.MutableSharedFlow
|
||||
import kotlinx.coroutines.flow.SharedFlow
|
||||
import kotlinx.coroutines.flow.asSharedFlow
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlinx.serialization.json.Json
|
||||
import mu.KotlinLogging
|
||||
import pw.binom.agentik.proto.Event as ProtoEvent
|
||||
import pw.binom.agentik.storage.Ids
|
||||
import pw.binom.agentik.storage.events.EventRecord
|
||||
import pw.binom.agentik.storage.events.EventStore
|
||||
import pw.binom.agentik.storage.events.EventType
|
||||
|
||||
internal class ConversationEvents {
|
||||
private val log = KotlinLogging.logger {}
|
||||
|
||||
/**
|
||||
* SSE-события диалога + (опционально) persistence в [EventStore].
|
||||
*
|
||||
* Двойная ответственность:
|
||||
* 1. Live-streaming через [flow] — клиенты подписываются на long-lived SSE.
|
||||
* 2. Durable storage через [eventStore] — для replay после disconnect
|
||||
* через /conversations/{id}/events/replay?after_id=X.
|
||||
*
|
||||
* **Persistence strategy**: при каждом [tryEmit]/[emit] параллельно пишем в
|
||||
* EventStore (fire-and-forget в IO scope). Ошибка БД НЕ должна ронять live-stream
|
||||
* — оборачиваем в runCatching и логируем.
|
||||
*
|
||||
* **Idempotency**: каждое event имеет детерминированный id (из messageStore при
|
||||
* создании), append с тем же id в EventStore — no-op. Это критично для retry
|
||||
* между producer и БД.
|
||||
*
|
||||
* **Dual-write cost**: на каждый event одно INSERT в SQLite. SQLite на локальном
|
||||
* диске выдерживает ~50K events/sec; для hot pathов можно вынести persist в
|
||||
* отдельную batched очередь. Для v1 — синхронный launch — OK.
|
||||
*/
|
||||
internal class ConversationEvents(
|
||||
private val eventStore: EventStore? = null,
|
||||
private val conversationId: String? = null,
|
||||
private val eventJson: Json = Json {
|
||||
ignoreUnknownKeys = true
|
||||
encodeDefaults = true
|
||||
},
|
||||
) {
|
||||
private val _flow = MutableSharedFlow<ProtoEvent>(
|
||||
replay = 0,
|
||||
extraBufferCapacity = 4096,
|
||||
@@ -15,7 +51,70 @@ internal class ConversationEvents {
|
||||
|
||||
val flow: SharedFlow<ProtoEvent> get() = _flow.asSharedFlow()
|
||||
|
||||
fun tryEmit(event: ProtoEvent): Boolean = _flow.tryEmit(event)
|
||||
|
||||
suspend fun emit(event: ProtoEvent) = _flow.emit(event)
|
||||
/**
|
||||
* Emit event в live-stream + persist в EventStore (если настроен).
|
||||
*
|
||||
* @return true если event попал в live-stream (false если buffer overflow
|
||||
* и event был дропнут — DROP_OLDEST policy).
|
||||
*/
|
||||
fun tryEmit(event: ProtoEvent): Boolean {
|
||||
val ok = _flow.tryEmit(event)
|
||||
if (ok) persistAsync(event)
|
||||
return ok
|
||||
}
|
||||
|
||||
/**
|
||||
* Same as [tryEmit] но suspend — ждёт места в buffer'е (а не дропает).
|
||||
* Используется реже — там где мы хотим гарантировать доставку подписчикам.
|
||||
*/
|
||||
suspend fun emit(event: ProtoEvent) {
|
||||
_flow.emit(event)
|
||||
persistAsync(event)
|
||||
}
|
||||
|
||||
/**
|
||||
* Persist в EventStore в fire-and-forget. Если [eventStore] == null — no-op
|
||||
* (in-memory dev или Android без persistence).
|
||||
*
|
||||
* **Не использует agentScope** — мы не знаем о нём здесь (ConversationEvents
|
||||
* не владеет lifecycle). Если нужна более аккуратная lifecycle management —
|
||||
* передавать scope параметром или держать свой CoroutineScope.
|
||||
*
|
||||
* **Сейчас**: создаём transient `GlobalScope`-like через `MainScope()`-style —
|
||||
* НЕТ, лучше через `CoroutineScope(SupervisorJob + Dispatchers.IO).launch`.
|
||||
* Это сделано лениво, чтобы не плодить треды при hot path.
|
||||
*/
|
||||
private fun persistAsync(event: ProtoEvent) {
|
||||
val store = eventStore ?: return
|
||||
val convId = conversationId ?: return // не знаем к чему привязать
|
||||
kotlinx.coroutines.CoroutineScope(kotlinx.coroutines.SupervisorJob() + kotlinx.coroutines.Dispatchers.IO).launch {
|
||||
runCatching {
|
||||
store.append(
|
||||
EventRecord(
|
||||
id = "ev-${Ids.new("conv")}",
|
||||
conversationId = convId,
|
||||
createdAt = event.date,
|
||||
type = mapEventType(event),
|
||||
payload = eventJson.encodeToString(ProtoEvent.serializer(), event),
|
||||
)
|
||||
)
|
||||
}.onFailure {
|
||||
log.warn(it) {
|
||||
"failed to persist conversation event ${event::class.simpleName}: ${it.message}"
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private fun mapEventType(event: ProtoEvent): EventType = when (event) {
|
||||
is ProtoEvent.StartReasoning -> EventType.CONVERSATION_START_REASONING
|
||||
is ProtoEvent.StartResponse -> EventType.CONVERSATION_START_RESPONSE
|
||||
is ProtoEvent.AppendText -> EventType.CONVERSATION_APPEND_TEXT
|
||||
is ProtoEvent.AppendImage -> EventType.CONVERSATION_APPEND_IMAGE
|
||||
is ProtoEvent.ToolCall -> EventType.CONVERSATION_TOOL_CALL
|
||||
is ProtoEvent.ToolResult -> EventType.CONVERSATION_TOOL_RESULT
|
||||
is ProtoEvent.End -> EventType.CONVERSATION_END
|
||||
is ProtoEvent.Interrupted -> EventType.CONVERSATION_INTERRUPTED
|
||||
is ProtoEvent.Error -> EventType.CONVERSATION_ERROR
|
||||
}
|
||||
}
|
||||
|
||||
@@ -84,7 +84,10 @@ class ConversationLoop(
|
||||
agentScope = agentScope,
|
||||
)
|
||||
|
||||
private val events = ConversationEvents()
|
||||
private val events = ConversationEvents(
|
||||
eventStore = storage.eventStore,
|
||||
conversationId = state.id,
|
||||
)
|
||||
|
||||
/** Per-conversation background event bus. Lifecycle scoped к этому ConversationLoop. */
|
||||
private val backgroundEvents = BackgroundEventBus()
|
||||
|
||||
@@ -1,5 +1,7 @@
|
||||
package pw.binom.agentik.storage
|
||||
|
||||
import pw.binom.agentik.storage.events.EventStore
|
||||
|
||||
/**
|
||||
* Агрегатор всех storage-интерфейсов, нужных агенту для работы с историей диалога.
|
||||
*
|
||||
@@ -11,7 +13,10 @@ package pw.binom.agentik.storage
|
||||
* `SkillStore` НЕ входит сюда — он живёт в модуле `:skills` (другая ответственность:
|
||||
* не сообщения/рефлексии, а контент-файлы навыков) и принимается отдельно в `ChatAgent`.
|
||||
*
|
||||
* AutoCloseable: один `close()` закрывает все четыре store'а. В реализациях,
|
||||
* [EventStore] входит начиная с commit "event-store" — для replay после
|
||||
* disconnect (см. `/events/replay` endpoint в `:server`).
|
||||
*
|
||||
* AutoCloseable: один `close()` закрывает все store'ы. В реализациях,
|
||||
* которые не владеют ресурсами (in-memory), close — no-op.
|
||||
*/
|
||||
data class StorageBundle(
|
||||
@@ -19,11 +24,13 @@ data class StorageBundle(
|
||||
val messageStore: MessageStore,
|
||||
val workingMemoryStore: WorkingMemoryStore,
|
||||
val reflectionStore: ReflectionStore,
|
||||
val eventStore: EventStore? = null,
|
||||
) : AutoCloseable {
|
||||
override fun close() {
|
||||
conversationStore.close()
|
||||
messageStore.close()
|
||||
workingMemoryStore.close()
|
||||
reflectionStore.close()
|
||||
eventStore?.close()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,109 @@
|
||||
package pw.binom.agentik.storage.events
|
||||
|
||||
import kotlinx.serialization.Serializable
|
||||
import kotlin.time.Instant
|
||||
|
||||
/**
|
||||
* Persistent event log для replay после disconnect.
|
||||
*
|
||||
* Зачем: SSE-подписка на `/events` и `/conversations/{id}/events` — cold (no replay).
|
||||
* Если клиент отвалился на час, он пропустил всё. [EventStore] даёт:
|
||||
* - append() — producer (ChatAgent) пишет при каждом event
|
||||
* - query() — consumer (server SSE replay endpoint) читает по cursor
|
||||
* - prune() — maintenance: удалить старые events по TTL
|
||||
*
|
||||
* Не заменяет live-подписку на [MutableSharedFlow] — это для долговременного
|
||||
* хранения, а live-streaming идёт через in-memory channel.
|
||||
*
|
||||
* Платформо-агностичный interface (KMP): impl в `:storage-sqlite` (JVM-only),
|
||||
* `:storage-inmemory` (KMP, для тестов и dev), и в будущем `:storage-sqlite-android`
|
||||
* для Android-агента.
|
||||
*
|
||||
* Payload — opaque JSON string. [storage-core] не должен знать про
|
||||
* kotlinx.serialization или [AgentEvent]/[Conversation.Event] типы (это `:proto`-шный
|
||||
* слой). Конвертация — на стороне producer'а (:standalone ChatAgent).
|
||||
*/
|
||||
interface EventStore : AutoCloseable {
|
||||
/**
|
||||
* Записать event. Идемпотентен по [EventRecord.id] — повторный append с тем же
|
||||
* id это no-op (важно для retry при network failure между producer'ом и БД).
|
||||
*/
|
||||
suspend fun append(record: EventRecord)
|
||||
|
||||
/**
|
||||
* Catchup query для reconnect.
|
||||
*
|
||||
* @param conversationId если `null` — глобальный catchup (для `/events/replay`).
|
||||
* если задан — только этот диалог (для `/conversations/{id}/events/replay`).
|
||||
* @param afterId exclusive cursor: вернуть events СТРОГО после этого id.
|
||||
* Если `null` — с начала.
|
||||
* @param limit max количество records (default 100). Caller делает пагинацию
|
||||
* пока `result.size == limit`.
|
||||
*
|
||||
* Сортировка: по [EventRecord.createdAt] ASC, ties broken по [EventRecord.id] ASC
|
||||
* (т.к. id содержит timestamp-like prefix в нашей схеме, это даёт стабильный порядок).
|
||||
*/
|
||||
suspend fun query(
|
||||
conversationId: String? = null,
|
||||
afterId: String? = null,
|
||||
limit: Int = 100,
|
||||
): List<EventRecord>
|
||||
|
||||
/**
|
||||
* Maintenance: удалить events старше [olderThan]. Возвращает количество удалённых.
|
||||
* Default вызывается из background scope раз в час (TTL = 24h типично).
|
||||
*/
|
||||
suspend fun pruneOlderThan(olderThan: Instant): Int
|
||||
|
||||
/** Сколько events всего хранится (для observability). */
|
||||
suspend fun count(): Int
|
||||
|
||||
override fun close()
|
||||
}
|
||||
|
||||
/**
|
||||
* Платформо-агностичная запись event'а.
|
||||
*
|
||||
* @param id уникальный в пределах EventStore. Convention: `"ev-<uuid>"`.
|
||||
* Используется как cursor для [EventStore.query].
|
||||
* @param conversationId `null` для agent-level events (Created/Deleted/Renamed).
|
||||
* Задан для conversation events.
|
||||
* @param createdAt UTC timestamp. Используется для сортировки в query() и для TTL в prune().
|
||||
* @param type kind of event (для индексирования/фильтрации; payload всё равно opaque).
|
||||
* @param payload opaque JSON string. Producer (:standalone ChatAgent) сериализует
|
||||
* [pw.binom.agentik.proto.AgentEvent] или [pw.binom.agentik.proto.Event]
|
||||
* в JSON перед append. Consumer (:server Routes) парсит обратно.
|
||||
*
|
||||
* Note: payload хранится as String, не ByteArray, чтобы не зависеть от kotlinx
|
||||
* serialization и platform-specific binary encoding в [storage-core].
|
||||
*/
|
||||
@Serializable
|
||||
data class EventRecord(
|
||||
val id: String,
|
||||
val conversationId: String?,
|
||||
val createdAt: Instant,
|
||||
val type: EventType,
|
||||
val payload: String,
|
||||
)
|
||||
|
||||
/**
|
||||
* Категория event'а — для индексирования и для фильтрации в query().
|
||||
*
|
||||
* Naming: AGENT_* — agent-level, CONVERSATION_* — turn-level.
|
||||
*/
|
||||
enum class EventType {
|
||||
AGENT_CREATED,
|
||||
AGENT_DELETED,
|
||||
AGENT_RENAMED,
|
||||
|
||||
CONVERSATION_START_REASONING,
|
||||
CONVERSATION_START_RESPONSE,
|
||||
CONVERSATION_APPEND_TEXT,
|
||||
CONVERSATION_APPEND_IMAGE,
|
||||
CONVERSATION_TOOL_CALL,
|
||||
CONVERSATION_TOOL_RESULT,
|
||||
CONVERSATION_END,
|
||||
CONVERSATION_INTERRUPTED,
|
||||
CONVERSATION_ERROR,
|
||||
// reserved for future — adding new variants doesn't break older consumers
|
||||
}
|
||||
+84
@@ -0,0 +1,84 @@
|
||||
package pw.binom.agentik.storage.inmemory
|
||||
|
||||
import pw.binom.agentik.storage.events.EventRecord
|
||||
import pw.binom.agentik.storage.events.EventStore
|
||||
import pw.binom.agentik.storage.events.EventType
|
||||
import kotlin.time.Instant
|
||||
import kotlinx.coroutines.sync.Mutex
|
||||
import kotlinx.coroutines.sync.withLock
|
||||
|
||||
/**
|
||||
* Thread-safe in-memory [EventStore]. Используется в тестах и в dev-режиме
|
||||
* `:standalone` (когда persistent events не нужны — например, для agentik-cli).
|
||||
*
|
||||
* Хранит все events в одном sorted-list (по createdAt). Поиск по afterId —
|
||||
* бинарный (list отсортирован). На 10K events работает за доли ms — для тестов
|
||||
* хватает. На production нужен Sqlite-impl.
|
||||
*
|
||||
* Thread-safety: один [Mutex] на все операции. Не concurrent-write-optimised —
|
||||
* для высоких нагрузок заменить на concurrent skip-list.
|
||||
*/
|
||||
class InMemoryEventStore : EventStore {
|
||||
|
||||
private val all: MutableList<EventRecord> = mutableListOf()
|
||||
private val byId: MutableMap<String, EventRecord> = mutableMapOf()
|
||||
private val mutex = Mutex()
|
||||
|
||||
override suspend fun append(record: EventRecord) {
|
||||
mutex.withLock {
|
||||
// Идемпотентность по id
|
||||
if (byId.containsKey(record.id)) return
|
||||
byId[record.id] = record
|
||||
// Insert maintaining ASC order by createdAt, ties broken by id.
|
||||
// binarySearch returns negative (-insertionPoint - 1) if not found,
|
||||
// or non-negative index if equal element found.
|
||||
val idx = all.binarySearch {
|
||||
val cmp = it.createdAt.compareTo(record.createdAt)
|
||||
if (cmp != 0) cmp else it.id.compareTo(record.id)
|
||||
}
|
||||
if (idx < 0) {
|
||||
all.add(-idx - 1, record)
|
||||
} else {
|
||||
// Found equal element — insert AFTER it to keep insertion order.
|
||||
all.add(idx + 1, record)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
override suspend fun query(
|
||||
conversationId: String?,
|
||||
afterId: String?,
|
||||
limit: Int,
|
||||
): List<EventRecord> {
|
||||
mutex.withLock {
|
||||
val startIdx = if (afterId == null) 0 else {
|
||||
val afterIdx = all.indexOfFirst { it.id == afterId }
|
||||
if (afterIdx < 0) return emptyList()
|
||||
afterIdx + 1
|
||||
}
|
||||
val filtered = if (conversationId == null) {
|
||||
all.subList(startIdx.coerceAtMost(all.size), all.size)
|
||||
} else {
|
||||
all.subList(startIdx.coerceAtMost(all.size), all.size)
|
||||
.filter { it.conversationId == conversationId }
|
||||
}
|
||||
return filtered.take(limit)
|
||||
}
|
||||
}
|
||||
|
||||
override suspend fun pruneOlderThan(olderThan: Instant): Int {
|
||||
mutex.withLock {
|
||||
val toRemove = all.filter { it.createdAt < olderThan }.map { it.id }
|
||||
if (toRemove.isEmpty()) return 0
|
||||
all.removeAll { it.id in toRemove }
|
||||
toRemove.forEach { byId.remove(it) }
|
||||
return toRemove.size
|
||||
}
|
||||
}
|
||||
|
||||
override suspend fun count(): Int = mutex.withLock { all.size }
|
||||
|
||||
override fun close() {
|
||||
// no-op: nothing to release
|
||||
}
|
||||
}
|
||||
+1
@@ -19,5 +19,6 @@ object InMemoryStorage {
|
||||
messageStore = InMemoryMessageStore(),
|
||||
workingMemoryStore = InMemoryWorkingMemoryStore(),
|
||||
reflectionStore = InMemoryReflectionStore(),
|
||||
eventStore = InMemoryEventStore(),
|
||||
)
|
||||
}
|
||||
|
||||
+165
@@ -0,0 +1,165 @@
|
||||
package pw.binom.agentik.storage.inmemory
|
||||
|
||||
import kotlinx.coroutines.coroutineScope
|
||||
import kotlinx.coroutines.test.runTest
|
||||
import pw.binom.agentik.storage.events.EventRecord
|
||||
import pw.binom.agentik.storage.events.EventType
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertNull
|
||||
import kotlin.test.assertTrue
|
||||
import kotlin.time.Clock
|
||||
import kotlin.time.Instant
|
||||
|
||||
class InMemoryEventStoreTest {
|
||||
|
||||
private fun rec(
|
||||
id: String,
|
||||
ts: Long,
|
||||
conv: String? = null,
|
||||
type: EventType = EventType.CONVERSATION_APPEND_TEXT,
|
||||
): EventRecord = EventRecord(
|
||||
id = id,
|
||||
conversationId = conv,
|
||||
createdAt = Instant.fromEpochMilliseconds(ts),
|
||||
type = type,
|
||||
payload = "{\"i\":$id}",
|
||||
)
|
||||
|
||||
@Test
|
||||
fun `append then query returns the record`() = runTest {
|
||||
val store = InMemoryEventStore()
|
||||
val r = rec("ev-1", ts = 1000)
|
||||
store.append(r)
|
||||
val result = store.query()
|
||||
assertEquals(listOf(r), result)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `query with afterId returns events strictly after the cursor`() = runTest {
|
||||
val store = InMemoryEventStore()
|
||||
store.append(rec("ev-1", ts = 1000))
|
||||
store.append(rec("ev-2", ts = 2000))
|
||||
store.append(rec("ev-3", ts = 3000))
|
||||
|
||||
// afterId = "ev-1" → only ev-2, ev-3 (exclusive)
|
||||
assertEquals(listOf("ev-2", "ev-3"), store.query(afterId = "ev-1").map { it.id })
|
||||
// afterId = "ev-2" → only ev-3
|
||||
assertEquals(listOf("ev-3"), store.query(afterId = "ev-2").map { it.id })
|
||||
// afterId = null → all
|
||||
assertEquals(listOf("ev-1", "ev-2", "ev-3"), store.query(afterId = null).map { it.id })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `query with afterId pointing at unknown id returns empty`() = runTest {
|
||||
val store = InMemoryEventStore()
|
||||
store.append(rec("ev-1", ts = 1000))
|
||||
assertEquals(emptyList(), store.query(afterId = "ev-unknown"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `query with conversationId filters to that conversation only`() = runTest {
|
||||
val store = InMemoryEventStore()
|
||||
store.append(rec("ev-1", ts = 1000, conv = "c-1"))
|
||||
store.append(rec("ev-2", ts = 2000, conv = "c-2"))
|
||||
store.append(rec("ev-3", ts = 3000, conv = "c-1"))
|
||||
|
||||
val c1 = store.query(conversationId = "c-1")
|
||||
assertEquals(listOf("ev-1", "ev-3"), c1.map { it.id })
|
||||
|
||||
val c2 = store.query(conversationId = "c-2")
|
||||
assertEquals(listOf("ev-2"), c2.map { it.id })
|
||||
|
||||
val all = store.query(conversationId = null)
|
||||
assertEquals(listOf("ev-1", "ev-2", "ev-3"), all.map { it.id })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `query with limit caps the result size`() = runTest {
|
||||
val store = InMemoryEventStore()
|
||||
repeat(10) { i -> store.append(rec("ev-$i", ts = (i * 1000).toLong())) }
|
||||
val first5 = store.query(limit = 5)
|
||||
assertEquals(5, first5.size)
|
||||
assertEquals(listOf("ev-0", "ev-1", "ev-2", "ev-3", "ev-4"), first5.map { it.id })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `append is idempotent on id`() = runTest {
|
||||
val store = InMemoryEventStore()
|
||||
val r1 = rec("ev-1", ts = 1000, type = EventType.CONVERSATION_APPEND_TEXT)
|
||||
val r2 = rec("ev-1", ts = 1000, type = EventType.CONVERSATION_TOOL_CALL) // тот же id, разный тип
|
||||
store.append(r1)
|
||||
store.append(r2)
|
||||
// Second append no-op (idempotent by id) — первая запись побеждает
|
||||
val result = store.query()
|
||||
assertEquals(1, result.size)
|
||||
assertEquals(EventType.CONVERSATION_APPEND_TEXT, result[0].type)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `records returned in createdAt ascending order`() = runTest {
|
||||
val store = InMemoryEventStore()
|
||||
// Insert out of order
|
||||
store.append(rec("ev-b", ts = 2000))
|
||||
store.append(rec("ev-a", ts = 1000))
|
||||
store.append(rec("ev-c", ts = 3000))
|
||||
|
||||
val result = store.query()
|
||||
assertEquals(listOf("ev-a", "ev-b", "ev-c"), result.map { it.id })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `records with same timestamp ordered by id ascending (stable sort)`() = runTest {
|
||||
val store = InMemoryEventStore()
|
||||
store.append(rec("ev-c", ts = 1000))
|
||||
store.append(rec("ev-a", ts = 1000))
|
||||
store.append(rec("ev-b", ts = 1000))
|
||||
|
||||
val result = store.query()
|
||||
// id lexicographic order: a < b < c
|
||||
assertEquals(listOf("ev-a", "ev-b", "ev-c"), result.map { it.id })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `pruneOlderThan removes records before cutoff`() = runTest {
|
||||
val store = InMemoryEventStore()
|
||||
store.append(rec("ev-1", ts = 1000))
|
||||
store.append(rec("ev-2", ts = 2000))
|
||||
store.append(rec("ev-3", ts = 3000))
|
||||
|
||||
val removed = store.pruneOlderThan(Instant.fromEpochMilliseconds(2500))
|
||||
assertEquals(2, removed)
|
||||
assertEquals(listOf("ev-3"), store.query().map { it.id })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `pruneOlderThan returns 0 when nothing to remove`() = runTest {
|
||||
val store = InMemoryEventStore()
|
||||
store.append(rec("ev-1", ts = 5000))
|
||||
val removed = store.pruneOlderThan(Instant.fromEpochMilliseconds(1000))
|
||||
assertEquals(0, removed)
|
||||
assertEquals(1, store.count())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `count returns number of stored records`() = runTest {
|
||||
val store = InMemoryEventStore()
|
||||
assertEquals(0, store.count())
|
||||
store.append(rec("ev-1", ts = 1000))
|
||||
store.append(rec("ev-2", ts = 2000))
|
||||
assertEquals(2, store.count())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `concurrent append from multiple coroutines all succeed`() = runTest {
|
||||
// Не strict — InMemoryEventStore использует Mutex, поэтому concurrent calls
|
||||
// сериализуются. Тест проверяет, что при последовательных append все ids
|
||||
// попадают в store (для concurrent test нужна отдельная TestScope — это
|
||||
// покрыто integration-тестами в :standalone).
|
||||
val store = InMemoryEventStore()
|
||||
repeat(51) { i ->
|
||||
store.append(rec("ev-$i", ts = i.toLong()))
|
||||
}
|
||||
assertEquals(51, store.count())
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,86 @@
|
||||
package pw.binom.agentik.storage.sqlite
|
||||
|
||||
import kotlin.time.Instant
|
||||
import mu.KotlinLogging
|
||||
import pw.binom.agentik.storage.events.EventRecord
|
||||
import pw.binom.agentik.storage.events.EventStore
|
||||
import pw.binom.agentik.storage.events.EventType
|
||||
import pw.binom.agentik.storage.sqlite.Agent_event as DbAgentEvent
|
||||
|
||||
private val log = KotlinLogging.logger {}
|
||||
|
||||
/**
|
||||
* SQLDelight-реализация [EventStore] поверх таблицы `agent_event`.
|
||||
*
|
||||
* Использует [EventStoreQueries] (генерируется SQLDelight из EventStore.sq).
|
||||
* Все запросы готовы — мы только маппим `Agent_event` (DB) ↔ `EventRecord` (domain).
|
||||
*
|
||||
* **Idempotency**: `append()` использует `INSERT OR IGNORE` — повторный append
|
||||
* с тем же id (network retry) — no-op. Это критично для producer'а, который
|
||||
* может retry при transient failure.
|
||||
*
|
||||
* **Pruning**: вызывай [pruneOlderThan] раз в час из background scope. Типичный
|
||||
* TTL = 24h. Если eventStore разрастётся (миллионы записей), индексы
|
||||
* (created_at, conversation_id+created_at) обеспечат O(log N) для query.
|
||||
*/
|
||||
class SqliteEventStore(
|
||||
private val db: AgentikDatabase,
|
||||
) : EventStore {
|
||||
|
||||
private val queries: EventStoreQueries get() = db.eventStoreQueries
|
||||
|
||||
override suspend fun append(record: EventRecord) {
|
||||
// payload — opaque JSON, хранится как UTF-8 bytes. Используем ByteArray,
|
||||
// потому что BLOB-колонка эффективнее TEXT для >100KB строк, и
|
||||
// API sqldelight нативно работает с ByteArray.
|
||||
val payloadBytes = record.payload.encodeToByteArray()
|
||||
queries.insert(
|
||||
id = record.id,
|
||||
conversation_id = record.conversationId,
|
||||
created_at = record.createdAt.toEpochMilliseconds(),
|
||||
type = record.type.name,
|
||||
payload = payloadBytes,
|
||||
)
|
||||
log.debug { "appended event id=${record.id} type=${record.type} conv=${record.conversationId}" }
|
||||
}
|
||||
|
||||
override suspend fun query(
|
||||
conversationId: String?,
|
||||
afterId: String?,
|
||||
limit: Int,
|
||||
): List<EventRecord> {
|
||||
// Если conversationId == null — используем queryAfter без фильтра
|
||||
// (он сам обрабатывает :convId IS NULL внутри SQL).
|
||||
// Если задан — queryAfterByConv (тогда SQL имеет WHERE conversation_id = :convId).
|
||||
val rows: List<DbAgentEvent> = if (conversationId == null) {
|
||||
queries.queryAfter(convId = null, afterId = afterId, limit = limit.toLong()).executeAsList()
|
||||
} else {
|
||||
queries.queryAfterByConv(convId = conversationId, afterId = afterId, limit = limit.toLong())
|
||||
.executeAsList()
|
||||
}
|
||||
return rows.map { it.toDomain() }
|
||||
}
|
||||
|
||||
override suspend fun pruneOlderThan(olderThan: Instant): Int {
|
||||
val deleted = queries.pruneOlderThan(olderThan.toEpochMilliseconds()).value
|
||||
if (deleted > 0) log.info { "pruned $deleted events older than $olderThan" }
|
||||
return deleted.toInt()
|
||||
}
|
||||
|
||||
override suspend fun count(): Int = queries.countAll().executeAsOne().toInt()
|
||||
|
||||
override fun close() {
|
||||
// no-op: lifecycle owned by AgentikDatabase / SqliteStores
|
||||
}
|
||||
|
||||
private fun DbAgentEvent.toDomain(): EventRecord = EventRecord(
|
||||
id = id,
|
||||
conversationId = conversation_id,
|
||||
createdAt = Instant.fromEpochMilliseconds(created_at),
|
||||
// type name → enum. Если в БД оказался неизвестный тип (новая версия,
|
||||
// unknown старому коду) — fallback на AGENT_CREATED (нейтральное значение).
|
||||
// Это безопаснее чем throw: клиент просто получит event с минимальным payload.
|
||||
type = runCatching { EventType.valueOf(type) }.getOrDefault(EventType.AGENT_CREATED),
|
||||
payload = payload.decodeToString(),
|
||||
)
|
||||
}
|
||||
@@ -7,12 +7,13 @@ import pw.binom.agentik.storage.ConversationStore
|
||||
import pw.binom.agentik.storage.MessageStore
|
||||
import pw.binom.agentik.storage.ReflectionStore
|
||||
import pw.binom.agentik.storage.StorageBundle
|
||||
import pw.binom.agentik.storage.events.EventStore
|
||||
import pw.binom.agentik.storage.sqlite.SqliteReflectionStore
|
||||
import pw.binom.agentik.storage.WorkingMemoryStore
|
||||
|
||||
/**
|
||||
* Корневой объект SQLite-слоя: держит [SqlDriver] и три [WorkingMemoryStore]/[MessageStore]/[ConversationStore].
|
||||
* Закрывается вместе с приложением.
|
||||
* Корневой объект SQLite-слоя: держит [SqlDriver] и пять store'ов (включая
|
||||
* [EventStore] — replay-after-disconnect).
|
||||
*/
|
||||
class SqliteStores private constructor(
|
||||
val driver: SqlDriver,
|
||||
@@ -20,6 +21,7 @@ class SqliteStores private constructor(
|
||||
val messages: MessageStore,
|
||||
val workingMemory: WorkingMemoryStore,
|
||||
val reflections: ReflectionStore,
|
||||
val events: EventStore,
|
||||
) : AutoCloseable {
|
||||
|
||||
/**
|
||||
@@ -32,9 +34,11 @@ class SqliteStores private constructor(
|
||||
messageStore = messages,
|
||||
workingMemoryStore = workingMemory,
|
||||
reflectionStore = reflections,
|
||||
eventStore = events,
|
||||
)
|
||||
|
||||
override fun close() {
|
||||
events.close()
|
||||
conversations.close()
|
||||
messages.close()
|
||||
workingMemory.close()
|
||||
@@ -55,6 +59,7 @@ class SqliteStores private constructor(
|
||||
messages = SqliteMessageStore(db),
|
||||
workingMemory = SqliteWorkingMemoryStore(db),
|
||||
reflections = SqliteReflectionStore(db),
|
||||
events = SqliteEventStore(db),
|
||||
)
|
||||
}
|
||||
|
||||
@@ -69,6 +74,7 @@ class SqliteStores private constructor(
|
||||
messages = SqliteMessageStore(db),
|
||||
workingMemory = SqliteWorkingMemoryStore(db),
|
||||
reflections = SqliteReflectionStore(db),
|
||||
events = SqliteEventStore(db),
|
||||
)
|
||||
}
|
||||
|
||||
@@ -118,6 +124,18 @@ class SqliteStores private constructor(
|
||||
CREATE INDEX IF NOT EXISTS idx_reflection_created ON reflection(created_at DESC);
|
||||
CREATE INDEX IF NOT EXISTS idx_reflection_conv ON reflection(conversation_id, created_at DESC);
|
||||
""".trimIndent(),
|
||||
// v3: event log для replay-after-disconnect (см. EventStore.sq)
|
||||
"""
|
||||
CREATE TABLE IF NOT EXISTS agent_event (
|
||||
id TEXT NOT NULL PRIMARY KEY,
|
||||
conversation_id TEXT,
|
||||
created_at INTEGER NOT NULL,
|
||||
type TEXT NOT NULL,
|
||||
payload BLOB NOT NULL
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS idx_agent_event_conv_time ON agent_event(conversation_id, created_at);
|
||||
CREATE INDEX IF NOT EXISTS idx_agent_event_time ON agent_event(created_at);
|
||||
""".trimIndent(),
|
||||
)
|
||||
for (sql in migrations) {
|
||||
driver.execute(null, sql, 0)
|
||||
|
||||
@@ -0,0 +1,49 @@
|
||||
-- Event log для replay после disconnect (см. EventStore.kt в :storage-core).
|
||||
-- Каждая запись — один event из ChatAgent (AgentEvent) или ConversationLoop
|
||||
-- (Conversation.Event). payload — opaque JSON, сериализуется в :standalone перед append.
|
||||
--
|
||||
-- Cursor для pagination: query() принимает afterId, возвращает events строго
|
||||
-- после него (exclusive). Используется /events/replay endpoint в :server.
|
||||
|
||||
CREATE TABLE agent_event (
|
||||
id TEXT NOT NULL PRIMARY KEY,
|
||||
conversation_id TEXT, -- NULL для agent-level events (Created/Deleted/Renamed)
|
||||
created_at INTEGER NOT NULL, -- epoch millis (UTC)
|
||||
type TEXT NOT NULL, -- EventType.name (см. :storage-core/events/EventStore.kt)
|
||||
payload BLOB NOT NULL -- serialized JSON
|
||||
);
|
||||
|
||||
CREATE INDEX idx_agent_event_conv_time ON agent_event(conversation_id, created_at);
|
||||
CREATE INDEX idx_agent_event_time ON agent_event(created_at);
|
||||
|
||||
insert:
|
||||
INSERT OR IGNORE INTO agent_event (id, conversation_id, created_at, type, payload)
|
||||
VALUES (?, ?, ?, ?, ?);
|
||||
|
||||
queryById:
|
||||
SELECT * FROM agent_event WHERE id = ?;
|
||||
|
||||
queryAfter:
|
||||
-- Catchup по conversationId (или все если null). afterId exclusive.
|
||||
-- Сортировка: created_at ASC, id ASC (стабильный tie-break для events с одинаковым timestamp).
|
||||
SELECT * FROM agent_event
|
||||
WHERE (:convId IS NULL OR conversation_id = :convId)
|
||||
AND id > COALESCE(:afterId, '')
|
||||
ORDER BY created_at ASC, id ASC
|
||||
LIMIT :limit;
|
||||
|
||||
queryAfterByConv:
|
||||
SELECT * FROM agent_event
|
||||
WHERE conversation_id = :convId
|
||||
AND id > COALESCE(:afterId, '')
|
||||
ORDER BY created_at ASC, id ASC
|
||||
LIMIT :limit;
|
||||
|
||||
countAll:
|
||||
SELECT COUNT(*) FROM agent_event;
|
||||
|
||||
countByConv:
|
||||
SELECT COUNT(*) FROM agent_event WHERE conversation_id = :convId;
|
||||
|
||||
pruneOlderThan:
|
||||
DELETE FROM agent_event WHERE created_at < :cutoffEpochMillis;
|
||||
+152
@@ -0,0 +1,152 @@
|
||||
package pw.binom.agentik.storage.sqlite
|
||||
|
||||
import kotlinx.coroutines.test.runTest
|
||||
import pw.binom.agentik.storage.events.EventRecord
|
||||
import pw.binom.agentik.storage.events.EventType
|
||||
import kotlin.test.AfterTest
|
||||
import kotlin.test.BeforeTest
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertTrue
|
||||
import kotlin.time.Instant
|
||||
|
||||
class SqliteEventStoreTest {
|
||||
|
||||
private lateinit var driver: app.cash.sqldelight.driver.jdbc.sqlite.JdbcSqliteDriver
|
||||
private lateinit var db: AgentikDatabase
|
||||
private lateinit var store: SqliteEventStore
|
||||
|
||||
@BeforeTest
|
||||
fun setup() {
|
||||
driver = app.cash.sqldelight.driver.jdbc.sqlite.JdbcSqliteDriver(
|
||||
app.cash.sqldelight.driver.jdbc.sqlite.JdbcSqliteDriver.IN_MEMORY,
|
||||
)
|
||||
AgentikDatabase.Schema.create(driver)
|
||||
db = AgentikDatabase(driver)
|
||||
store = SqliteEventStore(db)
|
||||
}
|
||||
|
||||
@AfterTest
|
||||
fun tearDown() {
|
||||
driver.close()
|
||||
}
|
||||
|
||||
private fun rec(
|
||||
id: String,
|
||||
ts: Long,
|
||||
conv: String? = null,
|
||||
type: EventType = EventType.CONVERSATION_APPEND_TEXT,
|
||||
): EventRecord = EventRecord(
|
||||
id = id,
|
||||
conversationId = conv,
|
||||
createdAt = Instant.fromEpochMilliseconds(ts),
|
||||
type = type,
|
||||
payload = """{"i":"$id"}""",
|
||||
)
|
||||
|
||||
@Test
|
||||
fun `append then query returns the record`() = runTest {
|
||||
val r = rec("ev-1", ts = 1000)
|
||||
store.append(r)
|
||||
assertEquals(listOf(r), store.query())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `query with conversationId filters to that conversation only`() = runTest {
|
||||
store.append(rec("ev-1", ts = 1000, conv = "c-1"))
|
||||
store.append(rec("ev-2", ts = 2000, conv = "c-2"))
|
||||
store.append(rec("ev-3", ts = 3000, conv = "c-1"))
|
||||
|
||||
val c1 = store.query(conversationId = "c-1")
|
||||
assertEquals(listOf("ev-1", "ev-3"), c1.map { it.id })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `query with afterId returns events strictly after cursor`() = runTest {
|
||||
store.append(rec("ev-1", ts = 1000))
|
||||
store.append(rec("ev-2", ts = 2000))
|
||||
store.append(rec("ev-3", ts = 3000))
|
||||
|
||||
assertEquals(listOf("ev-2", "ev-3"), store.query(afterId = "ev-1").map { it.id })
|
||||
assertEquals(listOf("ev-3"), store.query(afterId = "ev-2").map { it.id })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `append is idempotent (INSERT OR IGNORE)`() = runTest {
|
||||
val r1 = rec("ev-1", ts = 1000, type = EventType.CONVERSATION_APPEND_TEXT)
|
||||
val r2 = rec("ev-1", ts = 1000, type = EventType.CONVERSATION_TOOL_CALL)
|
||||
store.append(r1)
|
||||
store.append(r2)
|
||||
// INSERT OR IGNORE — вторая попытка no-op
|
||||
val result = store.query()
|
||||
assertEquals(1, result.size)
|
||||
assertEquals(EventType.CONVERSATION_APPEND_TEXT, result[0].type)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `records ordered by createdAt then id ascending`() = runTest {
|
||||
store.append(rec("ev-c", ts = 2000))
|
||||
store.append(rec("ev-a", ts = 1000))
|
||||
store.append(rec("ev-d", ts = 1000)) // same ts as a, different id
|
||||
store.append(rec("ev-b", ts = 1500))
|
||||
|
||||
val result = store.query()
|
||||
assertEquals(listOf("ev-a", "ev-d", "ev-b", "ev-c"), result.map { it.id })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `query with limit caps the result`() = runTest {
|
||||
repeat(10) { i -> store.append(rec("ev-$i", ts = (i * 100).toLong())) }
|
||||
val first3 = store.query(limit = 3)
|
||||
assertEquals(3, first3.size)
|
||||
assertEquals(listOf("ev-0", "ev-1", "ev-2"), first3.map { it.id })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `pruneOlderThan removes records before cutoff`() = runTest {
|
||||
store.append(rec("ev-1", ts = 1000))
|
||||
store.append(rec("ev-2", ts = 2000))
|
||||
store.append(rec("ev-3", ts = 3000))
|
||||
|
||||
val removed = store.pruneOlderThan(Instant.fromEpochMilliseconds(2500))
|
||||
assertEquals(2, removed)
|
||||
assertEquals(listOf("ev-3"), store.query().map { it.id })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `count returns total record count`() = runTest {
|
||||
assertEquals(0, store.count())
|
||||
store.append(rec("ev-1", ts = 1000))
|
||||
store.append(rec("ev-2", ts = 2000))
|
||||
assertEquals(2, store.count())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `unknown EventType in DB is loaded as AGENT_CREATED fallback`() = runTest {
|
||||
// Insert record with type that doesn't exist in current enum (simulating
|
||||
// a future enum value that older code doesn't know about).
|
||||
store.append(rec("ev-future", ts = 1000, type = EventType.CONVERSATION_END))
|
||||
// Manually mutate DB to use a fake type name (simulating forward-compat)
|
||||
db.eventStoreQueries.insert(
|
||||
id = "ev-fake",
|
||||
conversation_id = null,
|
||||
created_at = 2000,
|
||||
type = "FUTURE_TYPE_NOT_IN_ENUM",
|
||||
payload = "{}".encodeToByteArray(),
|
||||
)
|
||||
// Оба должны загрузиться (первый с правильным type, второй — с fallback)
|
||||
val result = store.query()
|
||||
assertEquals(2, result.size)
|
||||
assertEquals(EventType.CONVERSATION_END, result[0].type)
|
||||
assertEquals(EventType.AGENT_CREATED, result[1].type) // fallback для неизвестного type
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `close is no-op (does not close shared driver)`() = runTest {
|
||||
// store.close() НЕ должен закрывать driver — driver shared с другими store'ами.
|
||||
store.close()
|
||||
// Если бы close закрыл driver — следующий запрос упал бы. Проверяем что работает.
|
||||
store.append(rec("ev-1", ts = 1000))
|
||||
assertEquals(1, store.count())
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user