From 22a167ed038eecb6a80bff2eaaccec2aa41c4237 Mon Sep 17 00:00:00 2001 From: subochev Date: Sat, 3 Oct 2026 06:04:27 +0300 Subject: [PATCH] =?UTF-8?q?=D0=94=D0=BE=D0=B1=D0=B0=D0=B2=D0=BB=D1=8F?= =?UTF-8?q?=D0=B5=D1=82=20=D0=BC=D0=BE=D0=B4=D1=83=D0=BB=D1=8C=20:sync=20?= =?UTF-8?q?=E2=80=94=20generic=20cursor-journal=20=D1=81=D0=BE=D0=B1=D1=8B?= =?UTF-8?q?=D1=82=D0=B8=D0=B9.?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit EventStore/MutableEventStore (get с EventSeqResult.expired, push, dropAllBefore, clear, min/max курсоров) и ksqlite-реализация KSqliteMutableEventStore: кэш граничных курсоров, пагинация чтения, требование возрастающего курсора. Тесты на in-memory SQLite. --- sync/build.gradle.kts | 33 +++ .../kotlin/pw/binom/CursorExpiredException.kt | 3 + sync/src/commonMain/kotlin/pw/binom/Event.kt | 5 + .../kotlin/pw/binom/EventSeqResult.kt | 27 ++ .../commonMain/kotlin/pw/binom/EventStore.kt | 14 + .../pw/binom/KSqliteMutableEventStore.kt | 257 ++++++++++++++++++ .../kotlin/pw/binom/MutableEventStore.kt | 14 + .../kotlin/pw/binom/MutableStateStore.kt | 8 + .../kotlin/pw/binom/MutableStorage.kt | 12 + .../commonMain/kotlin/pw/binom/StateStore.kt | 12 + .../pw/binom/KSqliteMutableEventStoreTest.kt | 174 ++++++++++++ 11 files changed, 559 insertions(+) create mode 100644 sync/build.gradle.kts create mode 100644 sync/src/commonMain/kotlin/pw/binom/CursorExpiredException.kt create mode 100644 sync/src/commonMain/kotlin/pw/binom/Event.kt create mode 100644 sync/src/commonMain/kotlin/pw/binom/EventSeqResult.kt create mode 100644 sync/src/commonMain/kotlin/pw/binom/EventStore.kt create mode 100644 sync/src/commonMain/kotlin/pw/binom/KSqliteMutableEventStore.kt create mode 100644 sync/src/commonMain/kotlin/pw/binom/MutableEventStore.kt create mode 100644 sync/src/commonMain/kotlin/pw/binom/MutableStateStore.kt create mode 100644 sync/src/commonMain/kotlin/pw/binom/MutableStorage.kt create mode 100644 sync/src/commonMain/kotlin/pw/binom/StateStore.kt create mode 100644 sync/src/commonTest/kotlin/pw/binom/KSqliteMutableEventStoreTest.kt diff --git a/sync/build.gradle.kts b/sync/build.gradle.kts new file mode 100644 index 0000000..f066c25 --- /dev/null +++ b/sync/build.gradle.kts @@ -0,0 +1,33 @@ +plugins { + alias(libs.plugins.kotlin.multiplatform) + alias(libs.plugins.kotlin.serialization) +} + +kotlin { + jvmToolchain(21) + + // Полный набор KMP-целей. Зеркалит :proto. + jvm() + macosX64() + macosArm64() + iosX64() + iosArm64() + iosSimulatorArm64() + linuxX64() + linuxArm64() + mingwX64() + + sourceSets { + commonMain.dependencies { + api(libs.kotlinx.coroutines.core) + api(libs.kotlinx.serialization.core) + api(libs.kotlinx.serialization.protobuf) + api(libs.ksqlite) + api(project(":cursor-api")) + } + commonTest.dependencies { + implementation(kotlin("test")) + implementation(libs.kotlinx.coroutines.test) + } + } +} diff --git a/sync/src/commonMain/kotlin/pw/binom/CursorExpiredException.kt b/sync/src/commonMain/kotlin/pw/binom/CursorExpiredException.kt new file mode 100644 index 0000000..17c2d18 --- /dev/null +++ b/sync/src/commonMain/kotlin/pw/binom/CursorExpiredException.kt @@ -0,0 +1,3 @@ +package pw.binom + +class CursorExpiredException : RuntimeException() \ No newline at end of file diff --git a/sync/src/commonMain/kotlin/pw/binom/Event.kt b/sync/src/commonMain/kotlin/pw/binom/Event.kt new file mode 100644 index 0000000..af0559a --- /dev/null +++ b/sync/src/commonMain/kotlin/pw/binom/Event.kt @@ -0,0 +1,5 @@ +package pw.binom + +import pw.binom.agentik.cursor.Cursor + +data class Event(val data: EVENT, val seq: Cursor) \ No newline at end of file diff --git a/sync/src/commonMain/kotlin/pw/binom/EventSeqResult.kt b/sync/src/commonMain/kotlin/pw/binom/EventSeqResult.kt new file mode 100644 index 0000000..e120b4d --- /dev/null +++ b/sync/src/commonMain/kotlin/pw/binom/EventSeqResult.kt @@ -0,0 +1,27 @@ +package pw.binom + +import kotlinx.coroutines.flow.Flow +import kotlin.jvm.JvmInline + +@JvmInline +value class EventSeqResult private constructor(val raw: Any) { + + companion object { + private object ExpiredMarker + + fun expired() = EventSeqResult(ExpiredMarker) + fun success(flow: Flow>) = EventSeqResult(flow) + } + + + val isExpired + get() = raw === ExpiredMarker + + @Suppress("UNCHECKED_CAST") + fun getOrException(): Flow> { + if (isExpired) { + throw CursorExpiredException() + } + return raw as Flow> + } +} \ No newline at end of file diff --git a/sync/src/commonMain/kotlin/pw/binom/EventStore.kt b/sync/src/commonMain/kotlin/pw/binom/EventStore.kt new file mode 100644 index 0000000..11e19c1 --- /dev/null +++ b/sync/src/commonMain/kotlin/pw/binom/EventStore.kt @@ -0,0 +1,14 @@ +package pw.binom + +import pw.binom.agentik.cursor.Cursor + +interface EventStore : AutoCloseable { + + val minCursor: Cursor? + val maxCursor: Cursor? + + /** + * Получение всех событий ПОСЛЕ [cursor] включая сам [cursor] + */ + fun get(cursor: Cursor): EventSeqResult +} \ No newline at end of file diff --git a/sync/src/commonMain/kotlin/pw/binom/KSqliteMutableEventStore.kt b/sync/src/commonMain/kotlin/pw/binom/KSqliteMutableEventStore.kt new file mode 100644 index 0000000..821ed13 --- /dev/null +++ b/sync/src/commonMain/kotlin/pw/binom/KSqliteMutableEventStore.kt @@ -0,0 +1,257 @@ +package pw.binom + +import kotlinx.coroutines.flow.Flow +import kotlinx.coroutines.flow.flow +import kotlinx.serialization.ExperimentalSerializationApi +import kotlinx.serialization.KSerializer +import kotlinx.serialization.modules.EmptySerializersModule +import kotlinx.serialization.modules.SerializersModule +import kotlinx.serialization.protobuf.ProtoBuf +import pw.binom.agentik.cursor.Cursor +import pw.binom.db.ksqlite.SQLiteConnection +import pw.binom.db.ksqlite.transaction + +/** + * SQLite-хранилище журнала событий. + * + * **Не потокобезопасно.** Работать с ним нужно строго в один поток: используются общие + * prepared statements и кэши курсоров, а ленивая коллекция [get] должна выполняться в том же + * потоке, что и запись. + * + * **Известные ограничения и открытые вопросы:** + * + * - Фронтир компакции тут не хранится: [get] выводит `expired` из минимальной выжившей строки. + * Это даёт ложный `expired` на самой границе компакции и пустой поток (уже без `expired`) + * после полной очистки таблицы. Открытый вопрос: должен ли класс вообще решать про `expired`, + * или курсором владеет внешний компонент и проверку надо вынести наружу (тогда `get` просто + * верит переданному курсору). + * - [push] требует строго возрастающий курсор (`cursor > upperCursor`); повторная вставка того же + * курсора бросает исключение, а не является идемпотентным no-op. После полной очистки таблицы + * кэш `upperCursor` обнуляется, и это ограничение ослабляется. + * - Кэши [minCursor]/[maxCursor] валидны, только пока таблицу пишет этот же экземпляр; + * другой писатель делает их устаревшими. + * - Пагинация [get] через `LIMIT/OFFSET` небезопасна, если между страницами прошёл [dropAllBefore] + * (события могут быть молча пропущены), и стоит O(n²). Keyset-пагинация была бы корректнее. + * - `expired` вычисляется в момент вызова [get], а SELECT'ы выполняются лениво: компакция между + * этими моментами может привести к неполному потоку без признака `expired`. + * - Одновременная коллекция двух [get] не поддерживается (общий prepared statement). + * - [close] во время активной коллекции ломает поток. + * - Кодек зафиксирован на ProtoBuf, версия схемы/миграции не предусмотрены; ошибка декодирования + * выбрасывается посреди потока. + */ +@OptIn(ExperimentalSerializationApi::class) +class KSqliteMutableEventStore( + private val connection: SQLiteConnection, + private val ownsConnection: Boolean, + private val tableName: String, + private val serializer: KSerializer, + private val serializersModule: SerializersModule, +) : MutableEventStore { + constructor( + path: String, tableName: String, serializer: KSerializer, + serializersModule: SerializersModule = EmptySerializersModule(), + ) : this( + connection = SQLiteConnection.open(path = path), + ownsConnection = true, + tableName = tableName, + serializer = serializer, + serializersModule = serializersModule, + ) + + /** + * Внешнее соединение — store НЕ закрывает его в [close]. + */ + constructor( + connection: SQLiteConnection, + tableName: String, + serializer: KSerializer, + serializersModule: SerializersModule = EmptySerializersModule(), + ) : this( + connection = connection, + ownsConnection = false, + tableName = tableName, + serializer = serializer, + serializersModule = serializersModule + ) + + companion object { + private const val DATA = "DATA" + private const val CURSOR_CREATED_AT = "CREATED_AT" + private const val CURSOR_OFFSET = "OFFSET" + private const val PAGE_SIZE = 100 + } + + private val proto = ProtoBuf { + serializersModule = this@KSqliteMutableEventStore.serializersModule + } + + init { + val sql = """ + CREATE TABLE IF NOT EXISTS $tableName ( + $CURSOR_CREATED_AT INTEGER NOT NULL, + $CURSOR_OFFSET INTEGER NOT NULL, + $DATA BLOB NOT NULL, + PRIMARY KEY ($CURSOR_CREATED_AT, $CURSOR_OFFSET) + ); + """.trimIndent() + connection.transaction { + connection.exec(sql) + } + } + + private val pushStatement = connection.prepare( + """ + insert into $tableName ($CURSOR_CREATED_AT,$CURSOR_OFFSET,$DATA) values (?,?,?) + """.trimIndent() + ) + + private val dropAllBeforeStatement = connection.prepare( + """ + delete from $tableName + where $CURSOR_CREATED_AT < ? or ($CURSOR_CREATED_AT = ? and $CURSOR_OFFSET <= ?) + """.trimIndent() + ) + + private val clearStatement = connection.prepare( + """ + delete from $tableName + """.trimIndent() + ) + + private val getStatement = connection.prepare( + """ + select $CURSOR_CREATED_AT,$CURSOR_OFFSET,$DATA from $tableName + where $CURSOR_CREATED_AT > ? or ($CURSOR_CREATED_AT = ? and $CURSOR_OFFSET >= ?) + order by $CURSOR_CREATED_AT, $CURSOR_OFFSET + limit ? offset ? + """.trimIndent() + ) + + private val getLowerCursorStatement = connection.prepare( + """ + select $CURSOR_CREATED_AT,$CURSOR_OFFSET from $tableName + order by $CURSOR_CREATED_AT, $CURSOR_OFFSET + limit 1 + """.trimIndent() + ) + + private val getUpperCursorStatement = connection.prepare( + """ + select $CURSOR_CREATED_AT,$CURSOR_OFFSET from $tableName + order by $CURSOR_CREATED_AT desc, $CURSOR_OFFSET desc + limit 1 + """.trimIndent() + ) + + private var lowerCursor: Cursor? = null + private var upperCursor: Cursor? = null + + init { + lowerCursor = loadLowerCursor() + upperCursor = loadUpperCursor() + } + + private fun loadLowerCursor(): Cursor? = + getLowerCursorStatement.executeQuery().use { rs -> + if (!rs.next()) null + else Cursor(createdAt = rs.getLong(0)!!, offset = rs.getLong(1)!!.toULong()) + } + + private fun loadUpperCursor(): Cursor? = + getUpperCursorStatement.executeQuery().use { rs -> + if (!rs.next()) null + else Cursor(createdAt = rs.getLong(0)!!, offset = rs.getLong(1)!!.toULong()) + } + + override fun push(event: EVENT, cursor: Cursor) { + val previous = upperCursor + require(previous == null || cursor > previous) { + "cursor $cursor must be greater than the last cursor $previous" + } + val data = proto.encodeToByteArray(serializer, event) + pushStatement.bindLong(1, cursor.createdAt) + pushStatement.bindLong(2, cursor.offset.toLong()) + pushStatement.bindBlob(3, data) + pushStatement.executeUpdate() + upperCursor = cursor + if (lowerCursor == null) { + lowerCursor = cursor + } + } + + override fun dropAllBefore(cursor: Cursor) { + dropAllBeforeStatement.bindLong(1, cursor.createdAt) + dropAllBeforeStatement.bindLong(2, cursor.createdAt) + dropAllBeforeStatement.bindLong(3, cursor.offset.toLong()) + dropAllBeforeStatement.executeUpdate() + lowerCursor = loadLowerCursor() + upperCursor = loadUpperCursor() + } + + override fun clear() { + clearStatement.executeUpdate() + lowerCursor = null + upperCursor = null + } + + /** + * Возвращает минимальный курсор. Либо null если записей нет + */ + override val minCursor: Cursor? + get() = lowerCursor + + /** + * Возвращает максимальный курсор. Либо null если записей нет + */ + override val maxCursor: Cursor? + get() = upperCursor + + override fun get(cursor: Cursor): EventSeqResult { + val lower = minCursor + if (lower != null && cursor < lower) { + return EventSeqResult.expired() + } + return EventSeqResult.success(selectAfter(cursor)) + } + + private fun selectAfter(cursor: Cursor): Flow> = + flow { + var pageOffset = 0L + while (true) { + var count = 0 + getStatement.bindLong(1, cursor.createdAt) + getStatement.bindLong(2, cursor.createdAt) + getStatement.bindLong(3, cursor.offset.toLong()) + getStatement.bindLong(4, PAGE_SIZE.toLong()) + getStatement.bindLong(5, pageOffset) + getStatement.executeQuery().use { rs -> + while (rs.next()) { + count++ + emit( + Event( + data = proto.decodeFromByteArray(serializer, rs.getBlob(2)!!), + seq = Cursor( + createdAt = rs.getLong(0)!!, + offset = rs.getLong(1)!!.toULong(), + ), + ) + ) + } + } + if (count < PAGE_SIZE) break + pageOffset += count + } + } + + override fun close() { + pushStatement.close() + dropAllBeforeStatement.close() + clearStatement.close() + getStatement.close() + getLowerCursorStatement.close() + getUpperCursorStatement.close() + if (ownsConnection) { + connection.close() + } + } +} diff --git a/sync/src/commonMain/kotlin/pw/binom/MutableEventStore.kt b/sync/src/commonMain/kotlin/pw/binom/MutableEventStore.kt new file mode 100644 index 0000000..06af8fe --- /dev/null +++ b/sync/src/commonMain/kotlin/pw/binom/MutableEventStore.kt @@ -0,0 +1,14 @@ +package pw.binom + +import pw.binom.agentik.cursor.Cursor + +interface MutableEventStore : EventStore { + fun push(event: Event) = push(event.data, event.seq) + fun push(event: EVENT, cursor: Cursor) + fun clear() + + /** + * Удаляет все события ДО [cursor] включая сам [cursor] + */ + fun dropAllBefore(cursor: Cursor) +} \ No newline at end of file diff --git a/sync/src/commonMain/kotlin/pw/binom/MutableStateStore.kt b/sync/src/commonMain/kotlin/pw/binom/MutableStateStore.kt new file mode 100644 index 0000000..c38b9c9 --- /dev/null +++ b/sync/src/commonMain/kotlin/pw/binom/MutableStateStore.kt @@ -0,0 +1,8 @@ +package pw.binom + +import pw.binom.agentik.cursor.Cursor + +interface MutableStateStore : StateStore { + fun apply(event: EVENT, cursor: Cursor) + fun apply(event: Event) = apply(event.data, event.seq) +} \ No newline at end of file diff --git a/sync/src/commonMain/kotlin/pw/binom/MutableStorage.kt b/sync/src/commonMain/kotlin/pw/binom/MutableStorage.kt new file mode 100644 index 0000000..8c49b7b --- /dev/null +++ b/sync/src/commonMain/kotlin/pw/binom/MutableStorage.kt @@ -0,0 +1,12 @@ +package pw.binom + +import pw.binom.agentik.cursor.Cursor + +interface MutableStorage: AutoCloseable { + fun push(event: EVENT, cursor: Cursor) + + /** + * Возвращает состояние на момент [cursor] + */ + fun getState(cursor: Cursor): STATE +} \ No newline at end of file diff --git a/sync/src/commonMain/kotlin/pw/binom/StateStore.kt b/sync/src/commonMain/kotlin/pw/binom/StateStore.kt new file mode 100644 index 0000000..5f28b62 --- /dev/null +++ b/sync/src/commonMain/kotlin/pw/binom/StateStore.kt @@ -0,0 +1,12 @@ +package pw.binom + +import pw.binom.agentik.cursor.Cursor + +interface StateStore { + + + /** + * Возвращает состояние на момент [cursor] + */ + fun getState(cursor: Cursor): STATE +} \ No newline at end of file diff --git a/sync/src/commonTest/kotlin/pw/binom/KSqliteMutableEventStoreTest.kt b/sync/src/commonTest/kotlin/pw/binom/KSqliteMutableEventStoreTest.kt new file mode 100644 index 0000000..26c974a --- /dev/null +++ b/sync/src/commonTest/kotlin/pw/binom/KSqliteMutableEventStoreTest.kt @@ -0,0 +1,174 @@ +package pw.binom + +import kotlinx.coroutines.flow.toList +import kotlinx.coroutines.test.runTest +import kotlinx.serialization.Serializable +import pw.binom.agentik.cursor.Cursor +import pw.binom.db.ksqlite.SQLiteConnection +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertFailsWith +import kotlin.test.assertFalse +import kotlin.test.assertNull +import kotlin.test.assertTrue + +private const val EPOCH = 1_700_000_000_000L + +@Serializable +private data class TestEvent(val value: Int) + +class KSqliteMutableEventStoreTest { + + private fun cursor(offset: ULong, createdAt: Long = EPOCH) = Cursor(createdAt = createdAt, offset = offset) + + private suspend fun withStore(block: suspend (KSqliteMutableEventStore) -> Unit) { + val conn = SQLiteConnection.memory("sync-${kotlin.random.Random.nextLong()}") + val store = KSqliteMutableEventStore( + connection = conn, + tableName = "events", + serializer = TestEvent.serializer(), + ) + try { + block(store) + } finally { + store.close() + conn.close() + } + } + + private fun KSqliteMutableEventStore.pushRange(from: ULong, to: ULong) { + for (offset in from..to) { + push(TestEvent(offset.toInt()), cursor(offset)) + } + } + + @Test + fun emptyStoreHasNoCursorsAndEmptyStream() = runTest { + withStore { store -> + assertNull(store.minCursor) + assertNull(store.maxCursor) + + val result = store.get(cursor(1u)) + assertFalse(result.isExpired) + assertTrue(result.getOrException().toList().isEmpty()) + } + } + + @Test + fun returnsAllFromLowerCursorInOrder() = runTest { + withStore { store -> + store.pushRange(1u, 3u) + + val events = store.get(cursor(1u)).getOrException().toList() + + assertEquals(listOf(1uL, 2uL, 3uL), events.map { it.seq.offset }) + assertEquals(listOf(1, 2, 3), events.map { it.data.value }) + } + } + + @Test + fun returnsTailIncludingCursor() = runTest { + withStore { store -> + store.pushRange(1u, 5u) + + val events = store.get(cursor(3u)).getOrException().toList() + + assertEquals(listOf(3uL, 4uL, 5uL), events.map { it.seq.offset }) + } + } + + @Test + fun expiredWhenCursorIsBeforeLower() = runTest { + withStore { store -> + store.pushRange(2u, 4u) + + assertTrue(store.get(cursor(1u)).isExpired) + assertFalse(store.get(cursor(2u)).isExpired) + } + } + + @Test + fun pagesThroughMoreThanOnePage() = runTest { + withStore { store -> + store.pushRange(1u, 250u) + + val offsets = store.get(cursor(1u)).getOrException().toList().map { it.seq.offset } + + assertEquals((1uL..250uL).toList(), offsets) + } + } + + @Test + fun boundsTrackPushedEvents() = runTest { + withStore { store -> + store.pushRange(3u, 7u) + + assertEquals(cursor(3u), store.minCursor) + assertEquals(cursor(7u), store.maxCursor) + } + } + + @Test + fun pushAcceptsEventWrapper() = runTest { + withStore { store -> + store.push(Event(TestEvent(42), cursor(1u))) + + val events = store.get(cursor(1u)).getOrException().toList() + + assertEquals(listOf(42), events.map { it.data.value }) + } + } + + @Test + fun dropAllBeforeRemovesUpToCursorInclusive() = runTest { + withStore { store -> + store.pushRange(1u, 5u) + + store.dropAllBefore(cursor(3u)) + + assertEquals(cursor(4u), store.minCursor) + assertEquals(cursor(5u), store.maxCursor) + val events = store.get(cursor(4u)).getOrException().toList() + assertEquals(listOf(4uL, 5uL), events.map { it.seq.offset }) + } + } + + @Test + fun dropAllBeforeEverythingEmptiesJournal() = runTest { + withStore { store -> + store.pushRange(1u, 3u) + + store.dropAllBefore(cursor(3u)) + + assertNull(store.minCursor) + assertNull(store.maxCursor) + val result = store.get(cursor(3u)) + assertFalse(result.isExpired) + assertTrue(result.getOrException().toList().isEmpty()) + } + } + + @Test + fun clearRemovesEverything() = runTest { + withStore { store -> + store.pushRange(1u, 3u) + + store.clear() + + assertNull(store.minCursor) + assertNull(store.maxCursor) + assertTrue(store.get(cursor(1u)).getOrException().toList().isEmpty()) + } + } + + @Test + fun pushRequiresStrictlyIncreasingCursor() = runTest { + withStore { store -> + store.push(TestEvent(1), cursor(5u)) + + assertFailsWith { store.push(TestEvent(2), cursor(5u)) } + assertFailsWith { store.push(TestEvent(3), cursor(4u)) } + assertFailsWith { store.push(TestEvent(4), cursor(6u, createdAt = EPOCH - 1)) } + } + } +}