diff --git a/journal-api/src/commonMain/kotlin/pw/binom/agentik/journal/JournalStore.kt b/journal-api/src/commonMain/kotlin/pw/binom/agentik/journal/JournalStore.kt index 74d8826..ef98e2b 100644 --- a/journal-api/src/commonMain/kotlin/pw/binom/agentik/journal/JournalStore.kt +++ b/journal-api/src/commonMain/kotlin/pw/binom/agentik/journal/JournalStore.kt @@ -18,6 +18,28 @@ interface JournalStore : AutoCloseable { suspend fun list(conversationId: String, after: Instant, offset: Int, limit: Int): List + /** + * Сколько сообщений в диалоге [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()` пока не получит короткую страницу). Для diff --git a/journal-inmemory/src/commonMain/kotlin/pw/binom/agentik/journal/inmemory/InMemoryJournalStore.kt b/journal-inmemory/src/commonMain/kotlin/pw/binom/agentik/journal/inmemory/InMemoryJournalStore.kt index 8494db5..4745bf4 100644 --- a/journal-inmemory/src/commonMain/kotlin/pw/binom/agentik/journal/inmemory/InMemoryJournalStore.kt +++ b/journal-inmemory/src/commonMain/kotlin/pw/binom/agentik/journal/inmemory/InMemoryJournalStore.kt @@ -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 } diff --git a/journal-inmemory/src/commonTest/kotlin/pw/binom/agentik/journal/inmemory/InMemoryJournalStoreTest.kt b/journal-inmemory/src/commonTest/kotlin/pw/binom/agentik/journal/inmemory/InMemoryJournalStoreTest.kt index 2a33912..b30e7aa 100644 --- a/journal-inmemory/src/commonTest/kotlin/pw/binom/agentik/journal/inmemory/InMemoryJournalStoreTest.kt +++ b/journal-inmemory/src/commonTest/kotlin/pw/binom/agentik/journal/inmemory/InMemoryJournalStoreTest.kt @@ -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), + ) + } } diff --git a/journal-ksqlite/src/commonMain/kotlin/pw/binom/agentik/journal/ksqlite/KsqliteJournalStore.kt b/journal-ksqlite/src/commonMain/kotlin/pw/binom/agentik/journal/ksqlite/KsqliteJournalStore.kt index 413f69d..d3dfad7 100644 --- a/journal-ksqlite/src/commonMain/kotlin/pw/binom/agentik/journal/ksqlite/KsqliteJournalStore.kt +++ b/journal-ksqlite/src/commonMain/kotlin/pw/binom/agentik/journal/ksqlite/KsqliteJournalStore.kt @@ -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() } diff --git a/journal-ksqlite/src/commonTest/kotlin/pw/binom/agentik/journal/ksqlite/KsqliteJournalStoreTest.kt b/journal-ksqlite/src/commonTest/kotlin/pw/binom/agentik/journal/ksqlite/KsqliteJournalStoreTest.kt index 7ab73eb..9681926 100644 --- a/journal-ksqlite/src/commonTest/kotlin/pw/binom/agentik/journal/ksqlite/KsqliteJournalStoreTest.kt +++ b/journal-ksqlite/src/commonTest/kotlin/pw/binom/agentik/journal/ksqlite/KsqliteJournalStoreTest.kt @@ -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), + ) + } }