2 Commits
12 ... 15

11 changed files with 288 additions and 5 deletions
@@ -5,6 +5,7 @@ import io.ktor.client.call.body
import io.ktor.client.request.get
import io.ktor.client.request.parameter
import io.ktor.http.HttpStatusCode
import kotlinx.serialization.Serializable
import pw.binom.agentik.journal.JournalStore
import pw.binom.agentik.journal.MessageRecord
import kotlin.time.Instant
@@ -13,8 +14,11 @@ import kotlin.time.Instant
* HTTP-реализация [JournalStore] (append-only audit log сообщений диалога),
* ходящая в `:server`-фасад.
*
* **Endpoint**: `GET {baseUrl}/journal/conversations/{id}/messages?after=&offset=&limit=`
* (см. [pw.binom.agentik.server.journalRoutes]).
* **Endpoints** (см. [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 /
* ToolCall / ToolResult / Error). В отличие от `GET /conversations/{id}/messages`
@@ -54,7 +58,28 @@ internal class HttpJournalStore(
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() {
// HttpClient закрывает владелец (AgentClient / AgentikAgent).
}
}
@Serializable
private data class CountResponse(val count: Long)
@@ -18,6 +18,28 @@ interface JournalStore : AutoCloseable {
suspend fun list(conversationId: String, after: Instant, offset: Int, limit: Int): List<MessageRecord>
/**
* Сколько сообщений в диалоге [conversationId] всего.
*
* O(1) на SQL-бэкендах (`SELECT COUNT(*) ... WHERE conversation_id = ?`),
* O(N) на in-memory (size простого list'а с фильтром по conversationId).
* Не зависит от cursor'а [after] — для total-размера диалога.
*/
suspend fun count(conversationId: String): Long
/**
* Сколько сообщений в диалоге [conversationId] создано **позже** [after]
* (строго `createdAt > after`, как и в [list]).
*
* O(1) на SQL-бэкендах, O(N) на in-memory. Полезно для:
* - UI badge "N новых сообщений" — клиент знает последний `lastSeen`,
* сервер говорит `count(convId, after=lastSeen)`;
* - пагинации без получения самих записей: знаем лимит последней страницы,
* надо понять "есть ли ещё";
* - compaction-метрик: «сколько turn'ов осталось после cutoff».
*/
suspend fun count(conversationId: String, after: Instant): Long
/**
* Cold-flow paging через [list]. Default-реализация делает N+1 round-trip
* (по странице через `list()` пока не получит короткую страницу). Для
@@ -59,6 +59,14 @@ class InMemoryJournalStore : MutableJournalStore {
records.removeAll { it.conversationId == conversationId }
}
override suspend fun count(conversationId: String): Long = mutex.withLock {
records.count { it.conversationId == conversationId }.toLong()
}
override suspend fun count(conversationId: String, after: Instant): Long = mutex.withLock {
records.count { it.conversationId == conversationId && it.createdAt > after }.toLong()
}
/** Сколько записей сейчас в кэше. Для тестов/диагностики. */
suspend fun size(): Int = mutex.withLock { records.size }
@@ -77,4 +77,85 @@ class InMemoryJournalStoreTest {
assertTrue(store.list("c1", Instant.DISTANT_PAST, 0, 100).isEmpty())
assertEquals(listOf("m2"), store.list("c2", Instant.DISTANT_PAST, 0, 100).map { it.id })
}
@Test
fun `count returns total per conversation`() = runTest {
val store = InMemoryJournalStore()
val t0 = Instant.parse("2026-09-21T10:00:00Z")
assertEquals(0L, store.count("c1"))
store.append(userMsg("m1", "c1", "a", t0))
store.append(userMsg("m2", "c1", "b", t0 + 1.seconds))
store.append(userMsg("m3", "c2", "c", t0 + 2.seconds))
assertEquals(2L, store.count("c1"))
assertEquals(1L, store.count("c2"))
assertEquals(0L, store.count("never-existed"))
}
@Test
fun `count is unaffected by clear of another conversation`() = runTest {
val store = InMemoryJournalStore()
val t0 = Instant.parse("2026-09-21T10:00:00Z")
for (i in 1..20) {
store.append(userMsg("m$i", "c1", "x", t0 + i.seconds))
}
store.append(userMsg("n1", "c2", "y", t0))
assertEquals(20L, store.count("c1"))
assertEquals(1L, store.count("c2"))
store.clear("c1")
assertEquals(0L, store.count("c1"))
assertEquals(1L, store.count("c2"))
}
@Test
fun `count after cursor excludes earlier messages`() = runTest {
val store = InMemoryJournalStore()
val t0 = Instant.parse("2026-09-21T10:00:00Z")
store.append(userMsg("m1", "c1", "a", t0))
store.append(userMsg("m2", "c1", "b", t0 + 10.seconds))
store.append(userMsg("m3", "c1", "c", t0 + 20.seconds))
// strictly > t0
assertEquals(2L, store.count("c1", t0))
// strictly > t0+10s — only the last
assertEquals(1L, store.count("c1", t0 + 10.seconds))
// after last — empty
assertEquals(0L, store.count("c1", t0 + 20.seconds))
// distant past — all
assertEquals(3L, store.count("c1", Instant.DISTANT_PAST))
}
@Test
fun `count after cursor scopes to conversation`() = runTest {
val store = InMemoryJournalStore()
val t = Instant.parse("2026-09-21T10:00:00Z")
store.append(userMsg("m1", "c1", "a", t))
store.append(userMsg("m2", "c2", "b", t))
val future = Instant.parse("2099-01-01T00:00:00Z")
assertEquals(0L, store.count("c1", future))
assertEquals(0L, store.count("c2", future))
assertEquals(1L, store.count("c1", Instant.DISTANT_PAST))
assertEquals(1L, store.count("c2", Instant.DISTANT_PAST))
}
@Test
fun `count agrees with list size`() = runTest {
val store = InMemoryJournalStore()
val t0 = Instant.parse("2026-09-21T10:00:00Z")
for (i in 1..15) {
store.append(userMsg("m$i", "c1", "x", t0 + i.seconds))
}
assertEquals(
store.list("c1", Instant.DISTANT_PAST, 0, 1000).size.toLong(),
store.count("c1"),
)
val cutoff = t0 + 7.seconds
assertEquals(
store.list("c1", cutoff, 0, 1000).size.toLong(),
store.count("c1", cutoff),
)
}
}
@@ -100,6 +100,16 @@ class KsqliteJournalStore private constructor(
private val clearStmt: SQLitePreparedStatement = connection.prepare(
"DELETE FROM ${Schema.TABLE_MESSAGE} WHERE ${Schema.COL_CONVERSATION_ID} = ?"
)
private val countAllStmt: SQLitePreparedStatement = connection.prepare(
"SELECT COUNT(*) FROM ${Schema.TABLE_MESSAGE} WHERE ${Schema.COL_CONVERSATION_ID} = ?"
)
private val countAfterStmt: SQLitePreparedStatement = connection.prepare(
"""
SELECT COUNT(*) FROM ${Schema.TABLE_MESSAGE}
WHERE ${Schema.COL_CONVERSATION_ID} = ?
AND ${Schema.COL_CREATED_AT} > ?
""".trimIndent()
)
override suspend fun append(record: MessageRecord): Unit = withContext(Dispatchers.Default) {
val (kind, payload) = encodeRecord(record)
@@ -147,10 +157,37 @@ class KsqliteJournalStore private constructor(
}
}
override suspend fun count(conversationId: String): Long = withContext(Dispatchers.Default) {
mutex.withLock {
countAllStmt.reset()
countAllStmt.clearBindings()
countAllStmt.bindText(1, conversationId)
countAllStmt.executeQuery().use { rs ->
check(rs.next()) { "COUNT(*) must return at least one row" }
(rs.getLong(0) ?: 0L)
}
}
}
override suspend fun count(conversationId: String, after: Instant): Long = withContext(Dispatchers.Default) {
mutex.withLock {
countAfterStmt.reset()
countAfterStmt.clearBindings()
countAfterStmt.bindText(1, conversationId)
countAfterStmt.bindLong(2, after.toEpochMilliseconds())
countAfterStmt.executeQuery().use { rs ->
check(rs.next()) { "COUNT(*) must return at least one row" }
(rs.getLong(0) ?: 0L)
}
}
}
override fun close() {
insertStmt.close()
listStmt.close()
clearStmt.close()
countAllStmt.close()
countAfterStmt.close()
if (ownsConnection) {
connection.close()
}
@@ -96,4 +96,93 @@ class KsqliteJournalStoreTest {
assertEquals(emptyList(), store.listFlow("conv1", Instant.DISTANT_PAST).toList())
assertEquals(1, store.listFlow("conv2", Instant.DISTANT_PAST).toList().size)
}
@Test
fun testCountReturnsTotalForConversation() = runTest {
val t = Instant.parse("2026-09-15T10:00:00Z")
assertEquals(0L, store.count("conv1"))
store.append(MessageRecord.UserMessage("m1", "conv1", listOf(Content.Text("a")), t, null))
store.append(MessageRecord.UserMessage("m2", "conv1", listOf(Content.Text("b")), t, null))
store.append(MessageRecord.UserMessage("m3", "conv2", listOf(Content.Text("c")), t, null))
assertEquals(2L, store.count("conv1"))
assertEquals(1L, store.count("conv2"))
assertEquals(0L, store.count("missing"))
}
@Test
fun testCountIsolatedFromOtherConversations() = runTest {
val t = Instant.parse("2026-09-15T10:00:00Z")
// Bulk insert into conv1, single into conv2.
for (i in 1..50) {
store.append(MessageRecord.UserMessage("m$i", "conv1", listOf(Content.Text("x$i")), t, null))
}
store.append(MessageRecord.UserMessage("n1", "conv2", listOf(Content.Text("only")), t, null))
assertEquals(50L, store.count("conv1"))
assertEquals(1L, store.count("conv2"))
assertEquals(0L, store.count("never-existed"))
}
@Test
fun testCountAfterFiltersStrictly() = runTest {
val t1 = Instant.parse("2026-09-15T10:01:00Z")
val t2 = Instant.parse("2026-09-15T10:02:00Z")
val t3 = Instant.parse("2026-09-15T10:03:00Z")
store.append(MessageRecord.UserMessage("m1", "conv1", listOf(Content.Text("a")), t1, null))
store.append(MessageRecord.UserMessage("m2", "conv1", listOf(Content.Text("b")), t2, null))
store.append(MessageRecord.UserMessage("m3", "conv1", listOf(Content.Text("c")), t3, null))
// Strictly > — t1 is not counted when after=t1.
assertEquals(2L, store.count("conv1", t1))
// Strictly > — t2 IS counted (only t3 after).
assertEquals(1L, store.count("conv1", t2))
// After last — empty.
assertEquals(0L, store.count("conv1", t3))
// Distant past — all three.
assertEquals(3L, store.count("conv1", Instant.DISTANT_PAST))
}
@Test
fun testCountAfterIsolatesByConversation() = runTest {
val t = Instant.parse("2026-09-15T10:00:00Z")
store.append(MessageRecord.UserMessage("m1", "conv1", listOf(Content.Text("a")), t, null))
store.append(MessageRecord.UserMessage("m2", "conv2", listOf(Content.Text("b")), t, null))
// Cursor that excludes everything — both conversations read 0.
val future = Instant.parse("2099-01-01T00:00:00Z")
assertEquals(0L, store.count("conv1", future))
assertEquals(0L, store.count("conv2", future))
// Cursor that includes everything — only conv1's message matches its conversation.
assertEquals(1L, store.count("conv1", Instant.DISTANT_PAST))
assertEquals(1L, store.count("conv2", Instant.DISTANT_PAST))
}
@Test
fun testCountAgreesWithListSize() = runTest {
val t0 = Instant.parse("2026-09-15T10:00:00Z")
for (i in 1..10) {
store.append(
MessageRecord.UserMessage(
id = "m$i",
conversationId = "conv1",
content = listOf(Content.Text("x$i")),
createdAt = t0 + kotlin.time.Duration.parse("PT${i}S"),
context = null,
),
)
}
// count === list(..., 0, +∞).size
assertEquals(
store.list("conv1", Instant.DISTANT_PAST, 0, 1000).size.toLong(),
store.count("conv1"),
)
// count(after=t5) === list(..., after=t5, 0, +∞).size
val cutoff = t0 + kotlin.time.Duration.parse("PT5S")
assertEquals(
store.list("conv1", cutoff, 0, 1000).size.toLong(),
store.count("conv1", cutoff),
)
}
}
+3
View File
@@ -49,6 +49,9 @@ kotlin {
implementation(kotlin("test"))
implementation(libs.ktor.server.test.host)
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]), здесь клиент
* получает полный transcript с tool-call/tool-result/error payload-ами,
* turn-tokens и context-метаданными.
* - `GET /conversations/{id}/count?after=` — сколько сообщений в диалоге
* всего (без `after`) или строго позже `after` (с `after`). Лёгкий
* endpoint для UI-бейджей "N новых сообщений" и compaction-метрик;
* тело ответа — JSON `{"count": <Long>}`.
*
* **Read-only:** [JournalStore] не имеет `append` — запись только через
* writer-референс, который ChatAgent держит внутри (тип `MutableJournalStore`,
@@ -40,5 +44,19 @@ fun Route.journalRoutes(
val limit = call.request.queryParameters["limit"]?.toIntOrNull() ?: JournalStore.PAGE_SIZE
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 {
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() {}
}
-2
View File
@@ -4,7 +4,6 @@
Главный исполняемый модуль проекта — single-jar HTTP-сервер с:
- **AG-UI** transport на `POST /agui` (SSE) + `GET /health`.
- **A2A** transport на `POST /` (JSON-RPC) + `GET /.well-known/agent-card.json`.
- **`:proto`** transport на `POST /agentik/*` (HTTP+JSON+SSE) — наш stateful.
- **Embedded LLM backend**: `GOOGLE` (LiteRT) или `OPENAI`-совместимый
@@ -98,7 +97,6 @@ AGENTIK_GOOGLE_MODEL_PATH=/root/gemma-4-E2B-it.litertlm \
| Метод | Путь | Transport | Описание |
|---|---|---|---|
| `GET` | `/health` | любой | health-check (`{"ok":true}`) |
| `POST` | `/agui` | AG-UI | Стриминг run (SSE) |
| `POST` | `/` | A2A | JSON-RPC `message/send`, `tasks/get`, `tasks/cancel` |
| `GET` | `/.well-known/agent-card.json` | A2A | Discovery |
| `POST` | `/agentik/conversations` | :proto | Создать диалог |
@@ -84,7 +84,7 @@ fun main(args: Array<String>) {
private fun printHelp() {
println("""
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
""".trimIndent())
}