Добавляет модуль :sync — generic cursor-journal событий.

EventStore/MutableEventStore (get с EventSeqResult.expired, push, dropAllBefore, clear, min/max курсоров) и ksqlite-реализация KSqliteMutableEventStore: кэш граничных курсоров, пагинация чтения, требование возрастающего курсора. Тесты на in-memory SQLite.
This commit is contained in:
2026-10-03 06:04:27 +03:00
parent 5bdc517988
commit 22a167ed03
11 changed files with 559 additions and 0 deletions
+33
View File
@@ -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)
}
}
}
@@ -0,0 +1,3 @@
package pw.binom
class CursorExpiredException : RuntimeException()
@@ -0,0 +1,5 @@
package pw.binom
import pw.binom.agentik.cursor.Cursor
data class Event<EVENT>(val data: EVENT, val seq: Cursor)
@@ -0,0 +1,27 @@
package pw.binom
import kotlinx.coroutines.flow.Flow
import kotlin.jvm.JvmInline
@JvmInline
value class EventSeqResult<EVENT> private constructor(val raw: Any) {
companion object {
private object ExpiredMarker
fun <T> expired() = EventSeqResult<T>(ExpiredMarker)
fun <T> success(flow: Flow<Event<T>>) = EventSeqResult<T>(flow)
}
val isExpired
get() = raw === ExpiredMarker
@Suppress("UNCHECKED_CAST")
fun getOrException(): Flow<Event<EVENT>> {
if (isExpired) {
throw CursorExpiredException()
}
return raw as Flow<Event<EVENT>>
}
}
@@ -0,0 +1,14 @@
package pw.binom
import pw.binom.agentik.cursor.Cursor
interface EventStore<EVENT> : AutoCloseable {
val minCursor: Cursor?
val maxCursor: Cursor?
/**
* Получение всех событий ПОСЛЕ [cursor] включая сам [cursor]
*/
fun get(cursor: Cursor): EventSeqResult<EVENT>
}
@@ -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<EVENT>(
private val connection: SQLiteConnection,
private val ownsConnection: Boolean,
private val tableName: String,
private val serializer: KSerializer<EVENT>,
private val serializersModule: SerializersModule,
) : MutableEventStore<EVENT> {
constructor(
path: String, tableName: String, serializer: KSerializer<EVENT>,
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<EVENT>,
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<EVENT> {
val lower = minCursor
if (lower != null && cursor < lower) {
return EventSeqResult.expired()
}
return EventSeqResult.success(selectAfter(cursor))
}
private fun selectAfter(cursor: Cursor): Flow<Event<EVENT>> =
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()
}
}
}
@@ -0,0 +1,14 @@
package pw.binom
import pw.binom.agentik.cursor.Cursor
interface MutableEventStore<EVENT> : EventStore<EVENT> {
fun push(event: Event<EVENT>) = push(event.data, event.seq)
fun push(event: EVENT, cursor: Cursor)
fun clear()
/**
* Удаляет все события ДО [cursor] включая сам [cursor]
*/
fun dropAllBefore(cursor: Cursor)
}
@@ -0,0 +1,8 @@
package pw.binom
import pw.binom.agentik.cursor.Cursor
interface MutableStateStore<STATE, EVENT> : StateStore<STATE> {
fun apply(event: EVENT, cursor: Cursor)
fun apply(event: Event<EVENT>) = apply(event.data, event.seq)
}
@@ -0,0 +1,12 @@
package pw.binom
import pw.binom.agentik.cursor.Cursor
interface MutableStorage<STATE, EVENT>: AutoCloseable {
fun push(event: EVENT, cursor: Cursor)
/**
* Возвращает состояние на момент [cursor]
*/
fun getState(cursor: Cursor): STATE
}
@@ -0,0 +1,12 @@
package pw.binom
import pw.binom.agentik.cursor.Cursor
interface StateStore<STATE> {
/**
* Возвращает состояние на момент [cursor]
*/
fun getState(cursor: Cursor): STATE
}
@@ -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<TestEvent>) -> 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<TestEvent>.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<IllegalArgumentException> { store.push(TestEvent(2), cursor(5u)) }
assertFailsWith<IllegalArgumentException> { store.push(TestEvent(3), cursor(4u)) }
assertFailsWith<IllegalArgumentException> { store.push(TestEvent(4), cursor(6u, createdAt = EPOCH - 1)) }
}
}
}