2 Commits
13 ... 16

Author SHA1 Message Date
subochev f191387e99 chore(deps): bump ksqlite 0.1.2 → 0.1.3 + journal count routes
release / Publish KMP libraries → caffeine Nexus (release) Successful in 1m18s
ksqlite 0.1.3 опубликован в Maven Central — POM/module-metadata/jar все
на месте. Поднят в шести модулях (vector-index-ksqlite, journal-ksqlite,
context-ksqlite, reflection-ksqlite, standalone, memory-md-vector);
комментарии про 0.1.2 в build.gradle.kts обновлены.

Расширение :journal-api: добавлены два count-метода (total + count after
cursor) в интерфейс JournalStore — реализованы в :journal-ksqlite /
:journal-inmemory. Над ними добавлены HTTP-endpoint'ы
  GET /journal/conversations/{id}/count
  GET /journal/conversations/{id}/count?after=
(объединены в один маршрут с опциональным параметром) и HTTP-клиент
HttpJournalStore.count/concount. Сервер-фасад расширен тестом
JournalRoutesCountTest (5 кейсов через embedded CIO + реальный
InMemoryJournalStore).

:server:jvmTest 10/0, :standalone:jvmTest 129/0, :journal-ksqlite:jvmTest 25/0,
:journal-inmemory:jvmTest 19/0. jvmTest агрегат 425/0/0.
2026-09-23 15:59:25 +03:00
subochev d1b4f897b7 Add GET /conversations/{id}/count endpoint for total/filtered message counts, update JournalStore API, and implement client/server support with tests.
release / Publish KMP libraries → caffeine Nexus (release) Successful in 1m4s
2026-09-23 05:43:55 +03:00
13 changed files with 265 additions and 18 deletions
@@ -5,6 +5,7 @@ import io.ktor.client.call.body
import io.ktor.client.request.get import io.ktor.client.request.get
import io.ktor.client.request.parameter import io.ktor.client.request.parameter
import io.ktor.http.HttpStatusCode import io.ktor.http.HttpStatusCode
import kotlinx.serialization.Serializable
import pw.binom.agentik.journal.JournalStore import pw.binom.agentik.journal.JournalStore
import pw.binom.agentik.journal.MessageRecord import pw.binom.agentik.journal.MessageRecord
import kotlin.time.Instant import kotlin.time.Instant
@@ -13,8 +14,11 @@ import kotlin.time.Instant
* HTTP-реализация [JournalStore] (append-only audit log сообщений диалога), * HTTP-реализация [JournalStore] (append-only audit log сообщений диалога),
* ходящая в `:server`-фасад. * ходящая в `:server`-фасад.
* *
* **Endpoint**: `GET {baseUrl}/journal/conversations/{id}/messages?after=&offset=&limit=` * **Endpoints** (см. [pw.binom.agentik.server.journalRoutes]):
* (см. [pw.binom.agentik.server.journalRoutes]). * - `GET {baseUrl}/journal/conversations/{id}/messages?after=&offset=&limit=`
* → [list]
* - `GET {baseUrl}/journal/conversations/{id}/count` → [count] (total)
* - `GET {baseUrl}/journal/conversations/{id}/count?after=` → [count] (after cursor)
* *
* Возвращает raw [MessageRecord] (все типы: UserMessage / AssistantMessage / * Возвращает raw [MessageRecord] (все типы: UserMessage / AssistantMessage /
* ToolCall / ToolResult / Error). В отличие от `GET /conversations/{id}/messages` * ToolCall / ToolResult / Error). В отличие от `GET /conversations/{id}/messages`
@@ -54,7 +58,28 @@ internal class HttpJournalStore(
return response.body<List<MessageRecord>>() return response.body<List<MessageRecord>>()
} }
override suspend fun count(conversationId: String): Long {
val response = httpClient.get("$agentUrl/journal/conversations/$conversationId/count")
check(response.status == HttpStatusCode.OK) {
"journal.count: server returned ${response.status}"
}
return response.body<CountResponse>().count
}
override suspend fun count(conversationId: String, after: Instant): Long {
val response = httpClient.get("$agentUrl/journal/conversations/$conversationId/count") {
parameter("after", after.toString())
}
check(response.status == HttpStatusCode.OK) {
"journal.count(after): server returned ${response.status}"
}
return response.body<CountResponse>().count
}
override fun close() { override fun close() {
// HttpClient закрывает владелец (AgentClient / AgentikAgent). // HttpClient закрывает владелец (AgentClient / AgentikAgent).
} }
} }
@Serializable
private data class CountResponse(val count: Long)
+2 -2
View File
@@ -21,9 +21,9 @@ kotlin {
sourceSets { sourceSets {
commonMain.dependencies { commonMain.dependencies {
// ksqlite 0.1.2 опубликован в Maven Central — обычный // ksqlite 0.1.3 опубликован в Maven Central — обычный
// `mavenCentral()` в settings.gradle.kts его подтянет. // `mavenCentral()` в settings.gradle.kts его подтянет.
implementation("pw.binom.db:ksqlite:0.1.2") implementation("pw.binom.db:ksqlite:0.1.3")
implementation(libs.kotlinx.serialization.json) implementation(libs.kotlinx.serialization.json)
api(project(":context-api")) api(project(":context-api"))
+2 -2
View File
@@ -20,9 +20,9 @@ kotlin {
sourceSets { sourceSets {
commonMain.dependencies { commonMain.dependencies {
// ksqlite 0.1.2 опубликован в Maven Central — обычный // ksqlite 0.1.3 опубликован в Maven Central — обычный
// `mavenCentral()` в settings.gradle.kts его подтянет. // `mavenCentral()` в settings.gradle.kts его подтянет.
implementation("pw.binom.db:ksqlite:0.1.2") implementation("pw.binom.db:ksqlite:0.1.3")
implementation(libs.kotlinx.serialization.json) implementation(libs.kotlinx.serialization.json)
api(project(":journal-api")) api(project(":journal-api"))
+2 -2
View File
@@ -30,9 +30,9 @@ kotlin {
sourceSets { sourceSets {
commonMain.dependencies { commonMain.dependencies {
// ksqlite 0.1.2 опубликован в Maven Central — обычный // ksqlite 0.1.3 опубликован в Maven Central — обычный
// `mavenCentral()` в settings.gradle.kts его подтянет. // `mavenCentral()` в settings.gradle.kts его подтянет.
implementation("pw.binom.db:ksqlite:0.1.2") implementation("pw.binom.db:ksqlite:0.1.3")
implementation(libs.kotlinx.coroutines.core) implementation(libs.kotlinx.coroutines.core)
implementation(libs.kotlinx.io.core) implementation(libs.kotlinx.io.core)
+3 -3
View File
@@ -9,7 +9,7 @@ plugins {
// других ksqlite-модулях; этот модуль не претендует на полную схему агента. // других ksqlite-модулях; этот модуль не претендует на полную схему агента.
// //
// Цели сборки — jvm() + linuxX64() + mingwX64(); Apple targets auto-disabled // Цели сборки — jvm() + linuxX64() + mingwX64(); Apple targets auto-disabled
// на Linux (ksqlite 0.1.2 не публикует native артефакты для Apple). // на Linux (ksqlite 0.1.3 не публикует native артефакты для Apple).
kotlin { kotlin {
jvmToolchain(21) jvmToolchain(21)
@@ -20,9 +20,9 @@ kotlin {
sourceSets { sourceSets {
commonMain.dependencies { commonMain.dependencies {
// ksqlite 0.1.2 опубликован в Maven Central — обычный // ksqlite 0.1.3 опубликован в Maven Central — обычный
// `mavenCentral()` в settings.gradle.kts его подтянет. // `mavenCentral()` в settings.gradle.kts его подтянет.
implementation("pw.binom.db:ksqlite:0.1.2") implementation("pw.binom.db:ksqlite:0.1.3")
api(project(":reflection-api")) api(project(":reflection-api"))
// :reflection-api ссылается на :journal-api типы в сигнатурах // :reflection-api ссылается на :journal-api типы в сигнатурах
+3
View File
@@ -49,6 +49,9 @@ kotlin {
implementation(kotlin("test")) implementation(kotlin("test"))
implementation(libs.ktor.server.test.host) implementation(libs.ktor.server.test.host)
implementation(libs.ktor.server.cio) implementation(libs.ktor.server.cio)
implementation(libs.ktor.client.content.negotiation)
implementation(libs.ktor.serialization.kotlinx.json)
implementation(project(":journal-inmemory"))
} }
} }
} }
@@ -23,6 +23,10 @@ import pw.binom.agentik.journal.JournalStore
* project'нутые proto-[pw.binom.agentik.proto.Message]), здесь клиент * project'нутые proto-[pw.binom.agentik.proto.Message]), здесь клиент
* получает полный transcript с tool-call/tool-result/error payload-ами, * получает полный transcript с tool-call/tool-result/error payload-ами,
* turn-tokens и context-метаданными. * turn-tokens и context-метаданными.
* - `GET /conversations/{id}/count?after=` — сколько сообщений в диалоге
* всего (без `after`) или строго позже `after` (с `after`). Лёгкий
* endpoint для UI-бейджей "N новых сообщений" и compaction-метрик;
* тело ответа — JSON `{"count": <Long>}`.
* *
* **Read-only:** [JournalStore] не имеет `append` — запись только через * **Read-only:** [JournalStore] не имеет `append` — запись только через
* writer-референс, который ChatAgent держит внутри (тип `MutableJournalStore`, * writer-референс, который ChatAgent держит внутри (тип `MutableJournalStore`,
@@ -40,5 +44,19 @@ fun Route.journalRoutes(
val limit = call.request.queryParameters["limit"]?.toIntOrNull() ?: JournalStore.PAGE_SIZE val limit = call.request.queryParameters["limit"]?.toIntOrNull() ?: JournalStore.PAGE_SIZE
call.respond(journal.list(id, after, offset, limit)) call.respond(journal.list(id, after, offset, limit))
} }
get("/conversations/{id}/count") {
val id = call.parameters["id"]!!
val after = call.parseAfter()
val count = if (after == null) {
journal.count(id)
} else {
journal.count(id, after)
}
call.respond(CountResponse(count = count))
} }
} }
}
/** Тело ответа `GET /conversations/{id}/count`. */
@kotlinx.serialization.Serializable
private data class CountResponse(val count: Long)
@@ -44,6 +44,8 @@ class BearerTokenTest {
) : Agent { ) : Agent {
override val journal: JournalStore = object : JournalStore { override val journal: JournalStore = object : JournalStore {
override suspend fun list(conversationId: String, after: Instant, offset: Int, limit: Int) = emptyList<MessageRecord>() 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 listFlow(conversationId: String, after: Instant, pageSize: Int) = emptyFlow<MessageRecord>()
override fun close() {} override fun close() {}
} }
@@ -0,0 +1,201 @@
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.request.parameter
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.serialization.json.Json
import kotlinx.coroutines.flow.emptyFlow
import kotlinx.coroutines.runBlocking
import kotlinx.serialization.Serializable
import pw.binom.agentik.journal.ConversationStore
import pw.binom.agentik.journal.Content as JContent
import pw.binom.agentik.journal.JournalStore
import pw.binom.agentik.journal.MessageRecord
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 kotlin.test.AfterTest
import kotlin.test.BeforeTest
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.time.Duration.Companion.seconds
import kotlin.time.Instant
/**
* Интеграционные тесты нового `GET /journal/conversations/{id}/count`-endpoint'а
* в [journalRoutes].
*
* Поднимается embedded CIO-сервер с реальным [InMemoryJournalStore], пишутся
* `UserMessage` через [pw.binom.agentik.journal.MutableJournalStore.append],
* затем дёргаются оба варианта (с/без `after`) и проверяются JSON-shape +
* семантика.
*/
class JournalRoutesCountTest {
private lateinit var server: EmbeddedServer<*, *>
private lateinit var journal: InMemoryJournalStore
private var port: Int = 0
@Serializable
private data class CountResponse(val count: Long)
@BeforeTest
fun setup() {
journal = InMemoryJournalStore()
val fakeAgent = FakeAgent(journal)
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()
}
private fun append(id: String, convId: String, text: String, at: Instant) {
kotlinx.coroutines.runBlocking {
journal.append(
MessageRecord.UserMessage(
id = id,
conversationId = convId,
content = listOf(JContent.Text(text)),
createdAt = at,
),
)
}
}
private fun client(): HttpClient = HttpClient(CIO) {
install(ContentNegotiation) {
json(Json { ignoreUnknownKeys = true })
}
}
@Test
fun `count without after returns total in conversation`() = runBlocking {
val t0 = Instant.parse("2026-09-22T10:00:00Z")
append("m1", "c1", "a", t0)
append("m2", "c1", "b", t0 + 1.seconds)
append("m3", "c2", "c", t0 + 2.seconds)
val client = client()
try {
val resp = client.get("http://127.0.0.1:$port/agentik/journal/conversations/c1/count")
assertEquals(HttpStatusCode.OK, resp.status)
assertEquals(CountResponse(count = 2L), resp.body<CountResponse>())
} finally {
client.close()
}
}
@Test
fun `count with after returns strictly-after count`() = runBlocking {
val t0 = Instant.parse("2026-09-22T10:00:00Z")
append("m1", "c1", "a", t0)
append("m2", "c1", "b", t0 + 10.seconds)
append("m3", "c1", "c", t0 + 20.seconds)
val client = client()
try {
val resp = client.get("http://127.0.0.1:$port/agentik/journal/conversations/c1/count") {
parameter("after", t0.toString())
}
assertEquals(HttpStatusCode.OK, resp.status)
// strictly > t0 → m2, m3
assertEquals(CountResponse(count = 2L), resp.body<CountResponse>())
} finally {
client.close()
}
}
@Test
fun `count for unknown conversation returns zero`() = runBlocking {
val client = client()
try {
val resp = client.get("http://127.0.0.1:$port/agentik/journal/conversations/never-existed/count")
assertEquals(HttpStatusCode.OK, resp.status)
assertEquals(CountResponse(count = 0L), resp.body<CountResponse>())
} finally {
client.close()
}
}
@Test
fun `count body is JSON with count key`() = runBlocking {
val t0 = Instant.parse("2026-09-22T10:00:00Z")
append("m1", "c1", "x", t0)
append("m2", "c1", "y", t0 + 1.seconds)
val client = client()
try {
val resp = client.get("http://127.0.0.1:$port/agentik/journal/conversations/c1/count")
assertEquals(HttpStatusCode.OK, resp.status)
val raw = resp.bodyAsText()
// Сырое JSON-тело — стабильный shape, не зависит от сериализатора.
assertEquals("""{"count":2}""", raw)
} finally {
client.close()
}
}
@Test
fun `count isolates by conversation`() = runBlocking {
val t = Instant.parse("2026-09-22T10:00:00Z")
for (i in 1..5) append("m$i", "c1", "x", t + i.seconds)
append("n1", "c2", "y", t)
val client = client()
try {
val r1 = client.get("http://127.0.0.1:$port/agentik/journal/conversations/c1/count").body<CountResponse>()
val r2 = client.get("http://127.0.0.1:$port/agentik/journal/conversations/c2/count").body<CountResponse>()
assertEquals(5L, r1.count)
assertEquals(1L, r2.count)
} finally {
client.close()
}
}
/**
* Минимальный stub-агент: только [JournalStore] нужен для теста
* `journalRoutes`. Outbox/ConversationStore — noop-объекты,
* чтобы старт [agentikAgent] не валился.
*/
private class FakeAgent(private val js: JournalStore) : Agent {
override val id: String = "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 = 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): 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() {}
}
}
-2
View File
@@ -4,7 +4,6 @@
Главный исполняемый модуль проекта — single-jar HTTP-сервер с: Главный исполняемый модуль проекта — single-jar HTTP-сервер с:
- **AG-UI** transport на `POST /agui` (SSE) + `GET /health`.
- **A2A** transport на `POST /` (JSON-RPC) + `GET /.well-known/agent-card.json`. - **A2A** transport на `POST /` (JSON-RPC) + `GET /.well-known/agent-card.json`.
- **`:proto`** transport на `POST /agentik/*` (HTTP+JSON+SSE) — наш stateful. - **`:proto`** transport на `POST /agentik/*` (HTTP+JSON+SSE) — наш stateful.
- **Embedded LLM backend**: `GOOGLE` (LiteRT) или `OPENAI`-совместимый - **Embedded LLM backend**: `GOOGLE` (LiteRT) или `OPENAI`-совместимый
@@ -98,7 +97,6 @@ AGENTIK_GOOGLE_MODEL_PATH=/root/gemma-4-E2B-it.litertlm \
| Метод | Путь | Transport | Описание | | Метод | Путь | Transport | Описание |
|---|---|---|---| |---|---|---|---|
| `GET` | `/health` | любой | health-check (`{"ok":true}`) | | `GET` | `/health` | любой | health-check (`{"ok":true}`) |
| `POST` | `/agui` | AG-UI | Стриминг run (SSE) |
| `POST` | `/` | A2A | JSON-RPC `message/send`, `tasks/get`, `tasks/cancel` | | `POST` | `/` | A2A | JSON-RPC `message/send`, `tasks/get`, `tasks/cancel` |
| `GET` | `/.well-known/agent-card.json` | A2A | Discovery | | `GET` | `/.well-known/agent-card.json` | A2A | Discovery |
| `POST` | `/agentik/conversations` | :proto | Создать диалог | | `POST` | `/agentik/conversations` | :proto | Создать диалог |
+1 -1
View File
@@ -70,7 +70,7 @@ kotlin {
// ksqlite объявлен как `implementation` (не `api`) в каждом // ksqlite объявлен как `implementation` (не `api`) в каждом
// ksqlite-модуле, поэтому SqliteStores здесь использует // ksqlite-модуле, поэтому SqliteStores здесь использует
// SQLiteConnection напрямую — фиксируем зависимость явно. // SQLiteConnection напрямую — фиксируем зависимость явно.
implementation("pw.binom.db:ksqlite:0.1.2") implementation("pw.binom.db:ksqlite:0.1.3")
// Bounded-tail live event stream + per-event TTL. // Bounded-tail live event stream + per-event TTL.
implementation(project(":outbox-inmemory")) implementation(project(":outbox-inmemory"))
@@ -84,7 +84,7 @@ fun main(args: Array<String>) {
private fun printHelp() { private fun printHelp() {
println(""" println("""
agentik standalone — usage: agentik standalone — usage:
java -jar agentik.jar Start HTTP server (AGUI + A2A + :proto) java -jar agentik.jar Start HTTP server (A2A + :proto)
java -jar agentik.jar pull-model Download the LiteRT-LM model from static.binom.pw java -jar agentik.jar pull-model Download the LiteRT-LM model from static.binom.pw
""".trimIndent()) """.trimIndent())
} }
+3 -3
View File
@@ -7,7 +7,7 @@ plugins {
// обычная SQLite БД. Подходит для small-to-medium масштабов (≤10K записей на // обычная SQLite БД. Подходит для small-to-medium масштабов (≤10K записей на
// embedding ~512d); для больших датасетов — sqlite-vec или JVector. // embedding ~512d); для больших датасетов — sqlite-vec или JVector.
// //
// Цели сборки — только те, для которых ksqlite 0.1.2 опубликован в Maven Central: // Цели сборки — только те, для которых ksqlite 0.1.3 опубликован в Maven Central:
// jvm + linuxX64/Arm64 + mingwX64. Apple/iOS/tvOS/watchOS — НЕ публикуются; // jvm + linuxX64/Arm64 + mingwX64. Apple/iOS/tvOS/watchOS — НЕ публикуются;
// для них использовать JVM-only `:vector-index-jvector`. // для них использовать JVM-only `:vector-index-jvector`.
@@ -20,8 +20,8 @@ kotlin {
sourceSets { sourceSets {
commonMain.dependencies { commonMain.dependencies {
// ksqlite 0.1.2 опубликован в Maven Central. // ksqlite 0.1.3 опубликован в Maven Central.
implementation("pw.binom.db:ksqlite:0.1.2") implementation("pw.binom.db:ksqlite:0.1.3")
api(project(":vector-index-api")) api(project(":vector-index-api"))
} }
commonTest.dependencies { commonTest.dependencies {