chore(deps): bump ksqlite 0.1.2 → 0.1.3 + journal count routes
release / Publish KMP libraries → caffeine Nexus (release) Successful in 1m18s
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.
This commit is contained in:
@@ -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() {}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user