Add count methods to JournalStore API and implementations for message counting per conversation (count(conversationId) and count(conversationId, after)), with supporting tests.
release / Publish KMP libraries → caffeine Nexus (release) Failing after 32s
release / Publish KMP libraries → caffeine Nexus (release) Failing after 32s
This commit is contained in:
@@ -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),
|
||||||
|
)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
+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),
|
||||||
|
)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user