Compare commits
3 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| f191387e99 | |||
| d1b4f897b7 | |||
| a81d92f489 |
@@ -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)
|
||||||
|
|||||||
@@ -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"))
|
||||||
|
|||||||
@@ -18,6 +18,28 @@ interface JournalStore : AutoCloseable {
|
|||||||
|
|
||||||
suspend fun list(conversationId: String, after: Instant, offset: Int, limit: Int): List<MessageRecord>
|
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
|
* Cold-flow paging через [list]. Default-реализация делает N+1 round-trip
|
||||||
* (по странице через `list()` пока не получит короткую страницу). Для
|
* (по странице через `list()` пока не получит короткую страницу). Для
|
||||||
|
|||||||
+8
@@ -59,6 +59,14 @@ class InMemoryJournalStore : MutableJournalStore {
|
|||||||
records.removeAll { it.conversationId == conversationId }
|
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 }
|
suspend fun size(): Int = mutex.withLock { records.size }
|
||||||
|
|
||||||
|
|||||||
+81
@@ -77,4 +77,85 @@ class InMemoryJournalStoreTest {
|
|||||||
assertTrue(store.list("c1", Instant.DISTANT_PAST, 0, 100).isEmpty())
|
assertTrue(store.list("c1", Instant.DISTANT_PAST, 0, 100).isEmpty())
|
||||||
assertEquals(listOf("m2"), store.list("c2", Instant.DISTANT_PAST, 0, 100).map { it.id })
|
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),
|
||||||
|
)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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"))
|
||||||
|
|||||||
+37
@@ -100,6 +100,16 @@ class KsqliteJournalStore private constructor(
|
|||||||
private val clearStmt: SQLitePreparedStatement = connection.prepare(
|
private val clearStmt: SQLitePreparedStatement = connection.prepare(
|
||||||
"DELETE FROM ${Schema.TABLE_MESSAGE} WHERE ${Schema.COL_CONVERSATION_ID} = ?"
|
"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) {
|
override suspend fun append(record: MessageRecord): Unit = withContext(Dispatchers.Default) {
|
||||||
val (kind, payload) = encodeRecord(record)
|
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() {
|
override fun close() {
|
||||||
insertStmt.close()
|
insertStmt.close()
|
||||||
listStmt.close()
|
listStmt.close()
|
||||||
clearStmt.close()
|
clearStmt.close()
|
||||||
|
countAllStmt.close()
|
||||||
|
countAfterStmt.close()
|
||||||
if (ownsConnection) {
|
if (ownsConnection) {
|
||||||
connection.close()
|
connection.close()
|
||||||
}
|
}
|
||||||
|
|||||||
+89
@@ -96,4 +96,93 @@ class KsqliteJournalStoreTest {
|
|||||||
assertEquals(emptyList(), store.listFlow("conv1", Instant.DISTANT_PAST).toList())
|
assertEquals(emptyList(), store.listFlow("conv1", Instant.DISTANT_PAST).toList())
|
||||||
assertEquals(1, store.listFlow("conv2", Instant.DISTANT_PAST).toList().size)
|
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),
|
||||||
|
)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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)
|
||||||
|
|
||||||
|
|||||||
@@ -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 типы в сигнатурах
|
||||||
|
|||||||
@@ -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() {}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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 | Создать диалог |
|
||||||
|
|||||||
@@ -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())
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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 {
|
||||||
|
|||||||
Reference in New Issue
Block a user