remove :storage-ksqlite (conversation/message) and related tests; decouple schema from journal
ci / JVM build + tests (push) Successful in 6m6s

This commit is contained in:
2026-09-22 16:01:05 +03:00
parent 639c7d1748
commit 84f5fd84f3
39 changed files with 1569 additions and 616 deletions
@@ -26,11 +26,38 @@ import kotlinx.coroutines.withContext
* *
* ВНИМАНИЕ: `:storage-ksqlite/KsqliteWorkingMemoryStore.kt` остаётся на диске — * ВНИМАНИЕ: `:storage-ksqlite/KsqliteWorkingMemoryStore.kt` остаётся на диске —
* это копия, не замена. Не удалять старый файл; миграция consumers'ов — отдельно. * это копия, не замена. Не удалять старый файл; миграция consumers'ов — отдельно.
*
* ## Lifecycle соединения
*
* Семантика владения connection'ом идентична
* `pw.binom.agentik.journal.ksqlite.KsqliteJournalStore`:
* - `KsqliteContextStore(connection)` — внешнее соединение, store НЕ
* закрывает его в [close].
* - `KsqliteContextStore(path)` — открывает файловое соединение,
* закрывает его в [close].
* - `KsqliteContextStore.memory(name)` — in-memory, закрывает в [close].
*
* [Schema.migrate] прогоняется ВСЕГДА при конструировании (idempotent).
*/ */
class KsqliteContextStore( class KsqliteContextStore private constructor(
private val connection: SQLiteConnection, private val connection: SQLiteConnection,
private val ownsConnection: Boolean,
) : ContextStore { ) : ContextStore {
constructor(path: String) : this(
connection = SQLiteConnection.open(path = path),
ownsConnection = true,
)
constructor(connection: SQLiteConnection) : this(
connection = connection,
ownsConnection = false,
)
init {
Schema.migrate(connection)
}
private val mutex = Mutex() private val mutex = Mutex()
private val json = Json { ignoreUnknownKeys = true } private val json = Json { ignoreUnknownKeys = true }
@@ -179,6 +206,15 @@ class KsqliteContextStore(
maxOrderIdxStmt.close() maxOrderIdxStmt.close()
dropFromIdxStmt.close() dropFromIdxStmt.close()
insertSummaryStmt.close() insertSummaryStmt.close()
if (ownsConnection) connection.close()
}
companion object {
fun memory(name: String? = null): KsqliteContextStore =
KsqliteContextStore(
connection = SQLiteConnection.memory(name),
ownsConnection = true,
)
} }
private fun maxOrderIdx(conversationId: String): Long { private fun maxOrderIdx(conversationId: String): Long {
@@ -12,16 +12,9 @@ import kotlin.test.assertTrue
import kotlin.time.Instant import kotlin.time.Instant
/** /**
* Тесты для [KsqliteContextStore] — точная копия * Тесты для [KsqliteContextStore]. Автономная фикстура: in-memory
* `KsqliteWorkingMemoryStoreTest` из `:storage-ksqlite`, с переименованием * SQLiteConnection + конструктор `KsqliteContextStore(connection)` — store сам
* типов (`WorkingMemoryStore` → `ContextStore`) и обновлённым пакетом для * прогоняет `Schema.migrate` в init, явный вызов не нужен.
* `Content` (`pw.binom.agentik.journal` — новый canonical, но структура
* та же).
*
* Тестовая фикстура: in-memory SQLiteConnection, [Schema.migrate] в @BeforeTest,
* `KsqliteContextStore(conn)` + ручной close в @AfterTest. Никакой внешней
* зависимости от `KsqliteStores` из `:storage-ksqlite` — этот модуль
* автономный.
*/ */
class KsqliteContextStoreTest { class KsqliteContextStoreTest {
@@ -31,7 +24,6 @@ class KsqliteContextStoreTest {
@BeforeTest @BeforeTest
fun setup() { fun setup() {
conn = SQLiteConnection.memory("ctx-${kotlin.random.Random.nextLong()}") conn = SQLiteConnection.memory("ctx-${kotlin.random.Random.nextLong()}")
Schema.migrate(conn)
store = KsqliteContextStore(conn) store = KsqliteContextStore(conn)
} }
@@ -14,4 +14,12 @@ package pw.binom.agentik.journal
*/ */
interface MutableJournalStore : JournalStore { interface MutableJournalStore : JournalStore {
suspend fun append(record: MessageRecord) suspend fun append(record: MessageRecord)
/**
* Удалить все сообщения диалога [conversationId]. Используется
* владельцем lifecycle диалога при его удалении (каскад из
* ChatAgent.deleteConversation). Append-only природа audit log'а
* не нарушается — это bulk-clear, а не редактирование.
*/
suspend fun clear(conversationId: String)
} }
@@ -54,9 +54,9 @@ class InMemoryJournalStore : MutableJournalStore {
.toList() .toList()
} }
/** Сбросить кэш (например, когда диалог удалён). */ /** Удалить все записи диалога (каскад из ChatAgent.deleteConversation). */
suspend fun clear(): Unit = mutex.withLock { override suspend fun clear(conversationId: String): Unit = mutex.withLock {
records.clear() records.removeAll { it.conversationId == conversationId }
} }
/** Сколько записей сейчас в кэше. Для тестов/диагностики. */ /** Сколько записей сейчас в кэше. Для тестов/диагностики. */
@@ -29,7 +29,7 @@ class InMemoryMutableConversationStore(
private val clock: Clock = Clock.System, private val clock: Clock = Clock.System,
) : MutableConversationStore { ) : MutableConversationStore {
private val byId: MutableMap<String, ConversationRecord> = mutableMapOf() private val byId = mutableMapOf<String, ConversationRecord>()
private val mutex = Mutex() private val mutex = Mutex()
override suspend fun upsert(record: ConversationRecord) { override suspend fun upsert(record: ConversationRecord) {
@@ -66,13 +66,15 @@ class InMemoryJournalStoreTest {
} }
@Test @Test
fun `clear empties the cache`() = runTest { fun `clear empties a conversation only`() = runTest {
val store = InMemoryJournalStore() val store = InMemoryJournalStore()
val t0 = Instant.parse("2026-09-21T10:00:00Z") val t0 = Instant.parse("2026-09-21T10:00:00Z")
store.append(userMsg("m1", "c1", "x", t0)) store.append(userMsg("m1", "c1", "x", t0))
store.append(userMsg("m2", "c2", "y", t0))
assertEquals(2, store.size())
store.clear("c1")
assertEquals(1, store.size()) assertEquals(1, store.size())
store.clear()
assertEquals(0, store.size())
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 })
} }
} }
@@ -14,31 +14,67 @@ import kotlinx.coroutines.withContext
/** /**
* ksqlite-реализация [MutableJournalStore] (append-only audit log). * ksqlite-реализация [MutableJournalStore] (append-only audit log).
* *
* Структура — копия [pw.binom.agentik.storage.ksqlite.KsqliteMessageStore] * Миграция завершена: это ЕДИНСТВЕННЫЙ класс для message-таблицы. Параллельная
* из `:storage-ksqlite`, но: * копия `:storage-ksqlite/KsqliteMessageStore.kt` удалена вместе со своим
* - лежит в собственном модуле `:journal-ksqlite`; * тестом (consumer `KsqliteStores.assemble()` уже мигрировал на этот класс).
* - реализует переименованный [MutableJournalStore] (раньше был *
* `MutableMessageStore`, теперь главный класс — `JournalStore` / * ## Lifecycle соединения
* `MutableJournalStore`); сам тип записи [MessageRecord] не *
* переименовывался. * Три формы конструктора с разной семантикой владения:
* - `KsqliteJournalStore(connection)` — внешнее соединение, store НЕ закрывает
* его в [close]. Для shared-connection bundles (`KsqliteStores.assemble`),
* где один connection используется многими store'ами и закрывается bundle'ом.
* - `KsqliteJournalStore(path)` — открывает файловое соединение, закрывает
* его в [close].
* - `KsqliteJournalStore.memory(name)` — открывает in-memory соединение,
* закрывает его в [close].
*
* ## Миграция
*
* [Schema.migrate] прогоняется ВСЕГДА при конструировании — это idempotent
* (CREATE TABLE / INDEX IF NOT EXISTS), так что лишних эффектов нет ни в
* standalone-форме, ни в shared-connection bundle'е, где несколько store'ов
* прогоняют миграцию одной и той же схемы по очереди.
* *
* Prepared statements (insert / list / clear) препарируются один раз в * Prepared statements (insert / list / clear) препарируются один раз в
* конструкторе и закрываются в [close]. Без этого GC финалайзеры каждого * конструкторе и закрываются в [close] ДО закрытия owned connection. Без этого
* StmtHolder'а пытаются `sqlite3_finalize` stmt, чей parent connection уже * GC финалайзеры каждого StmtHolder'а пытаются `sqlite3_finalize` stmt, чей
* закрыт → SIGSEGV в `pthread_mutex_lock` (см. [pw.binom.db.ksqlite.StmtHolder]). * parent connection уже закрыт → SIGSEGV в `pthread_mutex_lock`
* (см. [pw.binom.db.ksqlite.StmtHolder]).
* *
* `payloadJson` хранит JSON-сериализованные kind-specific поля. encoding * `payloadJson` хранит JSON-сериализованные kind-specific поля. encoding
* helpers (`encodeRecord` / `toMessageRecord` / `CallPayload` / ...) лежат * helpers (`encodeRecord` / `toMessageRecord` / `CallPayload` / ...) лежат
* в [MessageCodecs.kt] рядом. * в [MessageCodecs.kt] рядом.
*
* ВНИМАНИЕ: `:storage-ksqlite/KsqliteMessageStore.kt` остаётся на диске —
* это копия, не замена. Не удалять старый файл; миграция consumers'ов —
* отдельно.
*/ */
class KsqliteJournalStore internal constructor( class KsqliteJournalStore private constructor(
private val connection: SQLiteConnection, private val connection: SQLiteConnection,
private val ownsConnection: Boolean,
) : MutableJournalStore { ) : MutableJournalStore {
/**
* Открывает файловое соединение через [SQLiteConnection.open] и берёт на
* себя его закрытие в [close]. Для standalone использования, когда у
* store'а нет bundle'а-владельца connection'а.
*/
constructor(path: String) : this(
connection = SQLiteConnection.open(path = path),
ownsConnection = true,
)
/**
* Внешнее соединение — store НЕ закрывает его в [close]. Для
* shared-connection bundles (`KsqliteStores.assemble`), где один
* connection используется многими store'ами и закрывается bundle'ом.
*/
constructor(connection: SQLiteConnection) : this(
connection = connection,
ownsConnection = false,
)
init {
Schema.migrate(connection)
}
private val mutex = Mutex() private val mutex = Mutex()
private val json = Json { ignoreUnknownKeys = true } private val json = Json { ignoreUnknownKeys = true }
@@ -94,13 +130,15 @@ class KsqliteJournalStore internal constructor(
listStmt.bindLong(4, offset.toLong()) listStmt.bindLong(4, offset.toLong())
val out = mutableListOf<MessageRecord>() val out = mutableListOf<MessageRecord>()
listStmt.executeQuery().use { rs -> listStmt.executeQuery().use { rs ->
while (rs.next()) out.add(rs.toMessageRecord(json)) while (rs.next()) {
out.add(rs.toMessageRecord(json))
}
} }
out out
} }
} }
internal suspend fun clear(conversationId: String): Unit = withContext(Dispatchers.Default) { override suspend fun clear(conversationId: String): Unit = withContext(Dispatchers.Default) {
mutex.withLock { mutex.withLock {
clearStmt.reset() clearStmt.reset()
clearStmt.clearBindings() clearStmt.clearBindings()
@@ -113,5 +151,21 @@ class KsqliteJournalStore internal constructor(
insertStmt.close() insertStmt.close()
listStmt.close() listStmt.close()
clearStmt.close() clearStmt.close()
if (ownsConnection) {
connection.close()
}
}
companion object {
/**
* Открывает in-memory соединение через [SQLiteConnection.memory] и
* берёт на себя его закрытие в [close]. Удобно для тестов и ephemeral
* runtime.
*/
fun memory(name: String? = null) =
KsqliteJournalStore(
connection = SQLiteConnection.memory(name),
ownsConnection = true,
)
} }
} }
@@ -1,4 +1,4 @@
package pw.binom.agentik.storage.ksqlite package pw.binom.agentik.journal.ksqlite
import kotlin.time.Clock import kotlin.time.Clock
import kotlin.time.Instant import kotlin.time.Instant
@@ -14,13 +14,50 @@ import kotlinx.coroutines.withContext
* ksqlite-реализация [MutableConversationStore]. Схема таблицы `conversation` живёт * ksqlite-реализация [MutableConversationStore]. Схема таблицы `conversation` живёт
* в [Schema] (миграция через PRAGMA user_version) — этот класс только * в [Schema] (миграция через PRAGMA user_version) — этот класс только
* готовит и выполняет SQL, ссылаясь на `Schema.COL_*` / `Schema.TABLE_*`. * готовит и выполняет SQL, ссылаясь на `Schema.COL_*` / `Schema.TABLE_*`.
*
* Каскадное удаление связанных данных (message + working_memory) делает
* владелец lifecycle диалога (см. ChatAgent.deleteConversation) — этот
* store знает только про свою таблицу.
*
* ## Lifecycle соединения
*
* Семантика владения connection'ом идентична [KsqliteJournalStore]:
* - `KsqliteMutableConversationStore(connection)` — внешнее соединение,
* store НЕ закрывает его в [close] (используется shared-connection
* bundle'ом `KsqliteStores.assemble`).
* - `KsqliteMutableConversationStore(path)` — открывает файловое соединение,
* закрывает его в [close].
* - `KsqliteMutableConversationStore.memory(name)` — in-memory, закрывает
* в [close].
*/ */
class KsqliteMutableConversationStore( class KsqliteMutableConversationStore private constructor(
private val connection: SQLiteConnection, private val connection: SQLiteConnection,
private val messageStore: KsqliteMessageStore? = null, private val ownsConnection: Boolean,
private val workingMemoryStore: KsqliteWorkingMemoryStore? = null,
) : MutableConversationStore { ) : MutableConversationStore {
/**
* Открывает файловое соединение через [SQLiteConnection.open] и берёт на
* себя его закрытие в [close]. Для standalone использования, когда у
* store'а нет bundle'а-владельца connection'а.
*/
constructor(path: String) : this(
connection = SQLiteConnection.open(path = path),
ownsConnection = true,
)
/**
* Внешнее соединение — store НЕ закрывает его в [close]. Для
* shared-connection bundles (`KsqliteStores.assemble`).
*/
constructor(connection: SQLiteConnection) : this(
connection = connection,
ownsConnection = false,
)
init {
Schema.migrate(connection)
}
private val mutex = Mutex() private val mutex = Mutex()
// pre-prepare всех statement'ов — аналогично KsqliteMessageStore (см. // pre-prepare всех statement'ов — аналогично KsqliteMessageStore (см.
@@ -125,8 +162,6 @@ class KsqliteMutableConversationStore(
// Проверяем существование через raw query, НЕ через get() — get() тоже // Проверяем существование через raw query, НЕ через get() — get() тоже
// берёт mutex (не реентрант), что привело бы к deadlock. // берёт mutex (не реентрант), что привело бы к deadlock.
if (!execExists(id)) return@withContext false if (!execExists(id)) return@withContext false
messageStore?.clear(id)
workingMemoryStore?.clear(id)
deleteStmt.reset() deleteStmt.reset()
deleteStmt.clearBindings() deleteStmt.clearBindings()
deleteStmt.bindText(1, id) deleteStmt.bindText(1, id)
@@ -188,6 +223,15 @@ class KsqliteMutableConversationStore(
renameStmt.close() renameStmt.close()
renameUpdatedAtStmt.close() renameUpdatedAtStmt.close()
touchStmt.close() touchStmt.close()
if (ownsConnection) connection.close()
}
companion object {
fun memory(name: String? = null): KsqliteMutableConversationStore =
KsqliteMutableConversationStore(
connection = SQLiteConnection.memory(name),
ownsConnection = true,
)
} }
private fun execExists(id: String): Boolean { private fun execExists(id: String): Boolean {
@@ -5,32 +5,51 @@ import pw.binom.db.ksqlite.SQLiteConnection
/** /**
* Имена таблиц/колонок/индексов для ksqlite-бэкенда `:journal-api`. * Имена таблиц/колонок/индексов для ksqlite-бэкенда `:journal-api`.
* *
* Минимум — только то, что относится к `message` (append-only audit log). * Владеет двумя таблицами:
* Остальные таблицы агента (`conversation`, `working_memory`, `reflection`) * - `conversation` — реестр диалогов агента (см. ConversationRecord);
* живут в других ksqlite-модулях. * - `message` — append-only audit log сообщений диалогов.
*
* `working_memory` и `reflection` живут в других ksqlite-модулях.
* *
* Все DDL/DML в этом модуле должны ссылаться на эти константы — никаких * Все DDL/DML в этом модуле должны ссылаться на эти константы — никаких
* хардкоженных литералов в `prepare("SELECT ... FROM foo ...")` в store'е. * хардкоженных литералов в `prepare("SELECT ... FROM foo ...")` в store'е.
*/ */
internal object Schema { object Schema {
/** Версия схемы модуля. Увеличивать при ЛЮБОМ изменении DDL. */ /** Версия схемы модуля. Увеличивать при ЛЮБОМ изменении DDL. */
const val CURRENT_VERSION: Int = 1 const val CURRENT_VERSION: Int = 1
// ───── Таблица ───── // ───── Таблицы ─────
const val TABLE_CONVERSATION = "conversation"
const val TABLE_MESSAGE = "message" const val TABLE_MESSAGE = "message"
// ───── Колонки ───── // ───── Колонки conversation ─────
const val COL_ID = "id" const val COL_ID = "id"
const val COL_TITLE = "title"
const val COL_IS_TEMPORAL = "is_temporal"
const val COL_CREATED_AT = "created_at"
const val COL_UPDATED_AT = "updated_at"
// ───── Колонки message ─────
const val COL_CONVERSATION_ID = "conversation_id" const val COL_CONVERSATION_ID = "conversation_id"
const val COL_KIND = "kind" const val COL_KIND = "kind"
const val COL_PAYLOAD_JSON = "payload_json" const val COL_PAYLOAD_JSON = "payload_json"
const val COL_CREATED_AT = "created_at"
// ───── Индексы ───── // ───── Индексы ─────
const val IDX_CONV_UPDATED = "idx_conv_updated"
const val IDX_MSG_CONV = "idx_msg_conv" const val IDX_MSG_CONV = "idx_msg_conv"
private val v1Ddl = """ private val v1ConversationDdl = """
CREATE TABLE IF NOT EXISTS $TABLE_CONVERSATION (
$COL_ID TEXT NOT NULL PRIMARY KEY,
$COL_TITLE TEXT,
$COL_IS_TEMPORAL INTEGER NOT NULL DEFAULT 0,
$COL_CREATED_AT INTEGER NOT NULL,
$COL_UPDATED_AT INTEGER NOT NULL
);
"""
private val v1MessageDdl = """
CREATE TABLE IF NOT EXISTS $TABLE_MESSAGE ( CREATE TABLE IF NOT EXISTS $TABLE_MESSAGE (
$COL_ID TEXT NOT NULL PRIMARY KEY, $COL_ID TEXT NOT NULL PRIMARY KEY,
$COL_CONVERSATION_ID TEXT NOT NULL, $COL_CONVERSATION_ID TEXT NOT NULL,
@@ -38,54 +57,43 @@ internal object Schema {
$COL_PAYLOAD_JSON TEXT NOT NULL, $COL_PAYLOAD_JSON TEXT NOT NULL,
$COL_CREATED_AT INTEGER NOT NULL $COL_CREATED_AT INTEGER NOT NULL
); );
""".trimIndent() """
private val v1IndexesDdl = """ private val v1IndexesDdl = """
CREATE INDEX IF NOT EXISTS $IDX_CONV_UPDATED
ON $TABLE_CONVERSATION($COL_UPDATED_AT DESC);
-- Главный hot-path индекс для list/сообщений: фильтр по conv +
-- сортировка по created_at (используется list(), cascade-clear, etc.)
CREATE INDEX IF NOT EXISTS $IDX_MSG_CONV CREATE INDEX IF NOT EXISTS $IDX_MSG_CONV
ON $TABLE_MESSAGE($COL_CONVERSATION_ID, $COL_CREATED_AT); ON $TABLE_MESSAGE($COL_CONVERSATION_ID, $COL_CREATED_AT);
""".trimIndent() """
/** /**
* Прогоняет миграцию схемы до [CURRENT_VERSION] на пустой или существующей БД. * Прогоняет миграцию схемы до [CURRENT_VERSION] на пустой или существующей БД.
* *
* Версия хранится в `PRAGMA user_version` (стандартный SQLite-механизм, * Гарантии:
* 32-bit int в заголовке БД — без своей таблицы). Каждая миграция — * - идемпотентность: `CREATE TABLE/INDEX IF NOT EXISTS` — безопасно на
* блок DDL под номером `fromV+1`, выполняется в транзакции. Если миграция * уже-мигрированной БД;
* упадёт посередине — `ROLLBACK` оставит БД на предыдущей версии. * - атомарность: каждая миграция в BEGIN/COMMIT — упал посреди →
* ROLLBACK оставит БД консистентной.
* *
* Идемпотентен: повторный вызов на уже мигрированной БД — no-op. **NOTE**: в сплит-мире (4 ksqlite-модуля, каждый владеет своей таблицей)
* user_version как gate перестал работать — два модуля ставят его в 1,
* второй вызов short-circuit'ит. Поэтому migrate() просто прогоняет DDL
* idempotently; координация multi-module миграций — ответственность
* вызывающего (см. `KsqliteStores.open()` в `:storage-ksqlite`).
*/ */
fun migrate(conn: SQLiteConnection) { fun migrate(conn: SQLiteConnection) {
val current = readUserVersion(conn)
if (current >= CURRENT_VERSION) return
conn.exec("BEGIN") conn.exec("BEGIN")
try { try {
if (current < 1) { conn.exec(v1ConversationDdl)
conn.exec(v1Ddl) conn.exec(v1MessageDdl)
conn.exec(v1IndexesDdl) conn.exec(v1IndexesDdl)
}
// future: if (current < 2) { conn.exec(v2Ddl) }
writeUserVersion(conn, CURRENT_VERSION)
conn.exec("COMMIT") conn.exec("COMMIT")
} catch (t: Throwable) { } catch (t: Throwable) {
runCatching { conn.exec("ROLLBACK") } runCatching { conn.exec("ROLLBACK") }
throw t throw t
} }
} }
private fun readUserVersion(conn: SQLiteConnection): Int {
conn.prepare("PRAGMA user_version").use { stmt ->
stmt.executeQuery().use { rs ->
if (rs.next()) return rs.getLong(0)?.toInt() ?: 0
}
}
return 0
}
private fun writeUserVersion(conn: SQLiteConnection, version: Int) {
// SQLite PRAGMA с literal-аргументом нельзя параметризовать через `?`,
// поэтому собираем SQL строкой (значение контролируемое, не user input).
conn.exec("PRAGMA user_version = $version")
}
} }
@@ -13,15 +13,9 @@ import kotlin.test.assertEquals
import kotlin.time.Instant import kotlin.time.Instant
/** /**
* Тесты для [KsqliteJournalStore] — точная копия * Тесты для [KsqliteJournalStore]. Автономная фикстура: in-memory
* `KsqliteMessageStoreTest` из `:storage-ksqlite`, с переименованием типов * SQLiteConnection + конструктор `KsqliteJournalStore(connection)` — store сам
* (`MessageStore` → `JournalStore`) и автономной фикстурой (in-memory * прогоняет `Schema.migrate` в init, явный вызов не нужен.
* SQLiteConnection + Schema.migrate).
*
* Тест `testClearRemovesByConversation` из оригинала использовал
* `stores.conversations.delete(...)` (cascade через `KsqliteStores`) — здесь
* он заменён на прямой вызов `store.clear(...)`, потому что `:journal-ksqlite`
* автономен и не знает про ConversationStore.
*/ */
class KsqliteJournalStoreTest { class KsqliteJournalStoreTest {
@@ -31,7 +25,6 @@ class KsqliteJournalStoreTest {
@BeforeTest @BeforeTest
fun setup() { fun setup() {
conn = SQLiteConnection.memory("journal-${kotlin.random.Random.nextLong()}") conn = SQLiteConnection.memory("journal-${kotlin.random.Random.nextLong()}")
Schema.migrate(conn)
store = KsqliteJournalStore(conn) store = KsqliteJournalStore(conn)
} }
@@ -33,8 +33,10 @@ data class Reflection(
interface ReflectionStore : AutoCloseable { interface ReflectionStore : AutoCloseable {
suspend fun insert(reflection: Reflection) suspend fun insert(reflection: Reflection)
suspend fun get(id: String): Reflection? suspend fun get(id: String): Reflection?
/** Самые свежие рефлексии (по всему агенту). */ /** Самые свежие рефлексии (по всему агенту). */
suspend fun listRecent(limit: Int = 10): List<Reflection> suspend fun listRecent(limit: Int = 10): List<Reflection>
/** Рефлексии для конкретного диалога. */ /** Рефлексии для конкретного диалога. */
suspend fun listForConversation(conversationId: String, limit: Int = 10): List<Reflection> suspend fun listForConversation(conversationId: String, limit: Int = 10): List<Reflection>
suspend fun deleteOlderThan(cutoff: Instant) suspend fun deleteOlderThan(cutoff: Instant)
@@ -55,5 +57,5 @@ sealed interface ReflectionEvent {
* БД-схемах было видно сразу. * БД-схемах было видно сразу.
*/ */
object Ids { object Ids {
fun new(): String = "refl-${kotlin.uuid.Uuid.random()}" fun new() = "refl-${kotlin.uuid.Uuid.random()}"
} }
+11
View File
@@ -127,3 +127,14 @@ include(":journal-ksqlite")
// тулы enable_toolset/disable_toolset. KMP, не зависит от :standalone, может быть // тулы enable_toolset/disable_toolset. KMP, не зависит от :standalone, может быть
// переиспользован в Android-сборке. Интеграция с ChatAgent — commit 5+. // переиспользован в Android-сборке. Интеграция с ChatAgent — commit 5+.
include(":agent-toolsets") include(":agent-toolsets")
// API-модуль vector-index: интерфейсы ANN-индекса, общие для любых бэкендов
// (JVector, sqlite-vec, inmemory, ...). Pure KMP, без реализаций.
include(":vector-index-api")
// JVector-реализация `MutableVectorIndexStore`. JVM-only — JVector публикуется
// только под JVM (нет KMP-таргетов). In-RAM (без persistence в v1).
include(":vector-index-jvector")
// ksqlite-реализация `MutableVectorIndexStore`. KMP (все 9 целей) — brute-force
// ANN: SELECT всех записей + cosine similarity в Kotlin. Persistence — обычная
// SQLite БД. Подходит для small-to-medium масштабов; для больших — sqlite-vec
// или JVector (выше).
include(":vector-index-ksqlite")
@@ -317,6 +317,11 @@ class ChatAgent(
conv?.close() conv?.close()
val ok = mutableConversationStore.delete(id) val ok = mutableConversationStore.delete(id)
if (ok) { if (ok) {
// Каскад — здесь, потому что владелец lifecycle диалога
// (ChatAgent) имеет доступ ко всем трём store'ам, а каждый
// отдельный store должен знать только про свою таблицу.
messageStore.clear(id)
workingMemoryStore.clear(id)
val event = AgentEvent.Deleted(date = now(), id = id) val event = AgentEvent.Deleted(date = now(), id = id)
eventStore.append(CommonEvent.Agent(date = now(), event = event)) eventStore.append(CommonEvent.Agent(date = now(), event = event))
} }
@@ -116,6 +116,11 @@ class PersistenceTest {
val removed = stores.conversations.delete("c1") val removed = stores.conversations.delete("c1")
assertTrue(removed) assertTrue(removed)
assertNull(stores.conversations.get("c1")) assertNull(stores.conversations.get("c1"))
// conversations.delete больше не каскадит — это делает владелец
// lifecycle (ChatAgent). Здесь проверяем, что каждый store делает
// только свою таблицу чистой при ручном вызове.
stores.messages.clear("c1")
stores.workingMemory.clear("c1")
assertEquals(emptyList(), stores.messages.listFlow("c1", Instant.DISTANT_PAST).toList()) assertEquals(emptyList(), stores.messages.listFlow("c1", Instant.DISTANT_PAST).toList())
assertEquals(emptyList(), stores.workingMemory.list("c1")) assertEquals(emptyList(), stores.workingMemory.list("c1"))
} }
@@ -37,6 +37,12 @@ class InMemoryMessageStore : MutableJournalStore {
} }
} }
override suspend fun clear(conversationId: String) {
mutex.withLock {
byConv.remove(conversationId)
}
}
override fun close() { override fun close() {
// no-op // no-op
} }
+1
View File
@@ -42,6 +42,7 @@ kotlin {
implementation(libs.kotlinx.serialization.json) implementation(libs.kotlinx.serialization.json)
api(project(":journal-api")) api(project(":journal-api"))
api(project(":journal-ksqlite"))
api(project(":reflection-api")) api(project(":reflection-api"))
api(project(":context-api")) api(project(":context-api"))
} }
@@ -1,105 +0,0 @@
package pw.binom.agentik.storage.ksqlite
import kotlinx.serialization.json.Json
import pw.binom.agentik.journal.MessageRecord
import pw.binom.agentik.journal.MutableJournalStore
import pw.binom.db.ksqlite.SQLiteConnection
import pw.binom.db.ksqlite.SQLitePreparedStatement
import kotlin.time.Instant
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock
import kotlinx.coroutines.withContext
/**
* ksqlite-реализация [MutableMessageStore] (append-only audit log).
*
* Prepared statements (insert / list / clear) препарируются один раз в
* конструкторе и закрываются в [close]. Без этого GC финалайзеры каждого
* StmtHolder'а пытаются `sqlite3_finalize` stmt, чей parent connection уже
* закрыт → SIGSEGV в `pthread_mutex_lock` (см. [pw.binom.db.ksqlite.StmtHolder]).
*
* `payloadJson` хранит JSON-сериализованные kind-specific поля. encoding
* helpers (`encodeRecord` / `toMessageRecord` / `CallPayload` / ...) лежат
* в [MessageCodecs.kt] рядом.
*/
class KsqliteMessageStore internal constructor(
private val connection: SQLiteConnection,
) : MutableJournalStore {
private val mutex = Mutex()
private val json = Json { ignoreUnknownKeys = true }
private val insertStmt: SQLitePreparedStatement = connection.prepare(
"""
INSERT INTO ${Schema.TABLE_MESSAGE}
(${Schema.COL_ID}, ${Schema.COL_CONVERSATION_ID}, ${Schema.COL_KIND},
${Schema.COL_PAYLOAD_JSON}, ${Schema.COL_CREATED_AT})
VALUES (?, ?, ?, ?, ?)
""".trimIndent()
)
private val listStmt: SQLitePreparedStatement = connection.prepare(
"""
SELECT ${Schema.COL_ID}, ${Schema.COL_CONVERSATION_ID}, ${Schema.COL_KIND},
${Schema.COL_PAYLOAD_JSON}, ${Schema.COL_CREATED_AT}
FROM ${Schema.TABLE_MESSAGE}
WHERE ${Schema.COL_CONVERSATION_ID} = ?
AND ${Schema.COL_CREATED_AT} > ?
ORDER BY ${Schema.COL_CREATED_AT} ASC, ${Schema.COL_ID} ASC
LIMIT ? OFFSET ?
""".trimIndent()
)
private val clearStmt: SQLitePreparedStatement = connection.prepare(
"DELETE FROM ${Schema.TABLE_MESSAGE} WHERE ${Schema.COL_CONVERSATION_ID} = ?"
)
override suspend fun append(record: MessageRecord): Unit = withContext(Dispatchers.Default) {
val (kind, payload) = encodeRecord(record)
mutex.withLock {
insertStmt.reset()
insertStmt.clearBindings()
insertStmt.bindText(1, record.id)
insertStmt.bindText(2, record.conversationId)
insertStmt.bindText(3, kind)
insertStmt.bindText(4, payload)
insertStmt.bindLong(5, record.createdAt.toEpochMilliseconds())
insertStmt.executeUpdate()
}
}
override suspend fun list(
conversationId: String,
after: Instant,
offset: Int,
limit: Int,
): List<MessageRecord> = withContext(Dispatchers.Default) {
mutex.withLock {
listStmt.reset()
listStmt.clearBindings()
listStmt.bindText(1, conversationId)
listStmt.bindLong(2, after.toEpochMilliseconds())
listStmt.bindLong(3, limit.toLong())
listStmt.bindLong(4, offset.toLong())
val out = mutableListOf<MessageRecord>()
listStmt.executeQuery().use { rs ->
while (rs.next()) out.add(rs.toMessageRecord(json))
}
out
}
}
internal suspend fun clear(conversationId: String): Unit = withContext(Dispatchers.Default) {
mutex.withLock {
clearStmt.reset()
clearStmt.clearBindings()
clearStmt.bindText(1, conversationId)
clearStmt.executeUpdate()
}
}
override fun close() {
insertStmt.close()
listStmt.close()
clearStmt.close()
}
}
@@ -16,11 +16,42 @@ import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock import kotlinx.coroutines.sync.withLock
import kotlinx.coroutines.withContext import kotlinx.coroutines.withContext
class KsqliteReflectionStore( /**
* ksqlite-реализация [ReflectionStore].
*
* ## Lifecycle соединения
*
* Семантика владения connection'ом идентична
* `pw.binom.agentik.journal.ksqlite.KsqliteJournalStore`:
* - `KsqliteReflectionStore(connection)` — внешнее соединение, store НЕ
* закрывает его в [close] и НЕ прогоняет миграцию (shared-connection bundle).
* - `KsqliteReflectionStore(path)` — открывает файловое соединение,
* прогоняет [Schema.migrate], закрывает соединение в [close].
* - `KsqliteReflectionStore.memory(name)` — in-memory, мигрирует,
* закрывает в [close].
*/
class KsqliteReflectionStore private constructor(
private val connection: SQLiteConnection, private val connection: SQLiteConnection,
private val ownsConnection: Boolean,
private val clock: Clock = Clock.System, private val clock: Clock = Clock.System,
) : ReflectionStore { ) : ReflectionStore {
constructor(path: String, clock: Clock = Clock.System) : this(
connection = SQLiteConnection.open(path = path),
ownsConnection = true,
clock = clock,
)
constructor(connection: SQLiteConnection, clock: Clock = Clock.System) : this(
connection = connection,
ownsConnection = false,
clock = clock,
)
init {
Schema.migrate(connection)
}
private val mutex = Mutex() private val mutex = Mutex()
private val ev = MutableSharedFlow<ReflectionEvent>(extraBufferCapacity = 16) private val ev = MutableSharedFlow<ReflectionEvent>(extraBufferCapacity = 16)
@@ -155,6 +186,16 @@ class KsqliteReflectionStore(
listForConvStmt.close() listForConvStmt.close()
deleteOlderThanStmt.close() deleteOlderThanStmt.close()
countStmt.close() countStmt.close()
if (ownsConnection) connection.close()
}
companion object {
fun memory(name: String? = null, clock: Clock = Clock.System): KsqliteReflectionStore =
KsqliteReflectionStore(
connection = SQLiteConnection.memory(name),
ownsConnection = true,
clock = clock,
)
} }
private fun SQLiteResultSet.toDomain(): Reflection = Reflection( private fun SQLiteResultSet.toDomain(): Reflection = Reflection(
@@ -3,15 +3,19 @@ package pw.binom.agentik.storage.ksqlite
import pw.binom.agentik.context.ContextStore import pw.binom.agentik.context.ContextStore
import pw.binom.agentik.journal.MutableJournalStore import pw.binom.agentik.journal.MutableJournalStore
import pw.binom.agentik.journal.MutableConversationStore import pw.binom.agentik.journal.MutableConversationStore
import pw.binom.agentik.journal.ksqlite.KsqliteMutableConversationStore as JournalKsqliteConversationStore
import pw.binom.agentik.journal.ksqlite.KsqliteJournalStore
import pw.binom.agentik.reflection.ReflectionStore import pw.binom.agentik.reflection.ReflectionStore
import pw.binom.db.ksqlite.SQLiteConnection import pw.binom.db.ksqlite.SQLiteConnection
/** /**
* Фабрика 4 store'ов поверх ksqlite. * Фабрика 4 store'ов поверх ksqlite.
* *
* Lifecycle: открывает [SQLiteConnection], прогоняет [Schema.migrate] (создаёт * Lifecycle: открывает [SQLiteConnection] и возвращает bundle из 4 store'ов.
* таблицы/индексы если их нет, догоняет версию схемы до [Schema.CURRENT_VERSION]), * Каждый store сам прогоняет свою схему в конструкторе
* возвращает bundle из 4 store'ов. Caller ДОЛЖЕН вызвать [close] при завершении. * (`Schema.migrate(connection)` — idempotent `CREATE TABLE IF NOT EXISTS`), так
* что явных вызовов миграции в bundle'е нет. Caller ДОЛЖЕН вызвать [close]
* при завершении.
* *
* @param path путь к .db файлу, либо URI для in-memory/shared-cache. * @param path путь к .db файлу, либо URI для in-memory/shared-cache.
*/ */
@@ -34,28 +38,18 @@ class KsqliteStores internal constructor(
companion object { companion object {
fun open(path: String): KsqliteStores { fun open(path: String): KsqliteStores =
val conn = SQLiteConnection.open(path) assemble(SQLiteConnection.open(path))
Schema.migrate(conn)
return assemble(conn)
}
fun inMemory(name: String = "agentik-test"): KsqliteStores { fun inMemory(name: String = "agentik-test"): KsqliteStores =
val conn = SQLiteConnection.memory(name) assemble(SQLiteConnection.memory(name))
Schema.migrate(conn)
return assemble(conn)
}
private fun assemble(conn: SQLiteConnection): KsqliteStores { private fun assemble(conn: SQLiteConnection): KsqliteStores = KsqliteStores(
val messages = KsqliteMessageStore(conn)
val working = KsqliteWorkingMemoryStore(conn)
return KsqliteStores(
connection = conn, connection = conn,
conversations = KsqliteMutableConversationStore(conn, messages, working), conversations = JournalKsqliteConversationStore(conn),
messages = messages, messages = KsqliteJournalStore(conn),
workingMemory = working, workingMemory = KsqliteWorkingMemoryStore(conn),
reflections = KsqliteReflectionStore(conn), reflections = KsqliteReflectionStore(conn),
) )
} }
} }
}
@@ -14,10 +14,39 @@ import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock import kotlinx.coroutines.sync.withLock
import kotlinx.coroutines.withContext import kotlinx.coroutines.withContext
class KsqliteWorkingMemoryStore( /**
* ksqlite-реализация [ContextStore] для `working_memory`.
*
* ## Lifecycle соединения
*
* Семантика владения connection'ом идентична
* `pw.binom.agentik.journal.ksqlite.KsqliteJournalStore`:
* - `KsqliteWorkingMemoryStore(connection)` — внешнее соединение, store НЕ
* закрывает его в [close] и НЕ прогоняет миграцию (shared-connection bundle).
* - `KsqliteWorkingMemoryStore(path)` — открывает файловое соединение,
* прогоняет [Schema.migrate], закрывает соединение в [close].
* - `KsqliteWorkingMemoryStore.memory(name)` — in-memory, мигрирует,
* закрывает в [close].
*/
class KsqliteWorkingMemoryStore private constructor(
private val connection: SQLiteConnection, private val connection: SQLiteConnection,
private val ownsConnection: Boolean,
) : ContextStore { ) : ContextStore {
constructor(path: String) : this(
connection = SQLiteConnection.open(path = path),
ownsConnection = true,
)
constructor(connection: SQLiteConnection) : this(
connection = connection,
ownsConnection = false,
)
init {
Schema.migrate(connection)
}
private val mutex = Mutex() private val mutex = Mutex()
private val json = Json { ignoreUnknownKeys = true } private val json = Json { ignoreUnknownKeys = true }
@@ -166,6 +195,15 @@ class KsqliteWorkingMemoryStore(
maxOrderIdxStmt.close() maxOrderIdxStmt.close()
dropFromIdxStmt.close() dropFromIdxStmt.close()
insertSummaryStmt.close() insertSummaryStmt.close()
if (ownsConnection) connection.close()
}
companion object {
fun memory(name: String? = null): KsqliteWorkingMemoryStore =
KsqliteWorkingMemoryStore(
connection = SQLiteConnection.memory(name),
ownsConnection = true,
)
} }
private fun maxOrderIdx(conversationId: String): Long { private fun maxOrderIdx(conversationId: String): Long {
@@ -3,12 +3,16 @@ package pw.binom.agentik.storage.ksqlite
import pw.binom.db.ksqlite.SQLiteConnection import pw.binom.db.ksqlite.SQLiteConnection
/** /**
* Имена таблиц/колонок/индексов для ksqlite-бэкенда agentik'а. * Имена таблиц/колонок/индексов для ksqlite-бэкенда agentik'а,
* которыми владеет `:storage-ksqlite`: `working_memory` + `reflection`.
* *
* Все DDL/DML через ksqlite должны ссылаться на эти константы — никаких * `conversation` + `message` уехали в `:journal-ksqlite` —
* хардкоженных литералов в `prepare("SELECT ... FROM foo ...")` в каждом * см. `pw.binom.agentik.journal.ksqlite.Schema`. Управляющий
* store'е. Это (1) даёт единую точку правды при будущих миграциях и * [KsqliteStores.open] / [KsqliteStores.inMemory] прогоняет ОБА
* (2) делает rename'ы безопасными (компилятор поймает все использования). * `Schema.migrate(conn)` подряд.
*
* Все DDL/DML в этом модуле ссылаются на эти константы — никаких
* хардкоженных литералов в `prepare("SELECT ... FROM foo ...")` в store'е.
*/ */
internal object Schema { internal object Schema {
@@ -16,27 +20,17 @@ internal object Schema {
const val CURRENT_VERSION: Int = 1 const val CURRENT_VERSION: Int = 1
// ───── Таблицы ───── // ───── Таблицы ─────
const val TABLE_CONVERSATION = "conversation"
const val TABLE_MESSAGE = "message"
const val TABLE_WORKING_MEMORY = "working_memory" const val TABLE_WORKING_MEMORY = "working_memory"
const val TABLE_REFLECTION = "reflection" const val TABLE_REFLECTION = "reflection"
const val TABLE_MIGRATION = "migration" // (резерв на будущее, сейчас версия в user_version)
// ───── Колонки conversation ─────
const val COL_ID = "id"
const val COL_TITLE = "title"
const val COL_IS_TEMPORAL = "is_temporal"
const val COL_CREATED_AT = "created_at"
const val COL_UPDATED_AT = "updated_at"
// ───── Колонки message ─────
const val COL_CONVERSATION_ID = "conversation_id"
const val COL_KIND = "kind"
const val COL_PAYLOAD_JSON = "payload_json"
// ───── Колонки working_memory ───── // ───── Колонки working_memory ─────
const val COL_ID = "id"
const val COL_CONVERSATION_ID = "conversation_id"
const val COL_ORDER_IDX = "order_idx" const val COL_ORDER_IDX = "order_idx"
const val COL_SOURCE_MESSAGE_ID = "source_message_id" const val COL_SOURCE_MESSAGE_ID = "source_message_id"
const val COL_KIND = "kind"
const val COL_PAYLOAD_JSON = "payload_json"
const val COL_CREATED_AT = "created_at"
// ───── Колонки reflection ───── // ───── Колонки reflection ─────
const val COL_TURNS_ANALYZED = "turns_analyzed" const val COL_TURNS_ANALYZED = "turns_analyzed"
@@ -45,8 +39,6 @@ internal object Schema {
const val COL_WEAK_SPOTS_JSON = "weak_spots_json" const val COL_WEAK_SPOTS_JSON = "weak_spots_json"
// ───── Индексы ───── // ───── Индексы ─────
const val IDX_CONV_UPDATED = "idx_conv_updated"
const val IDX_MSG_CONV = "idx_msg_conv"
const val IDX_WM_UNIQUE = "idx_wm_unique" const val IDX_WM_UNIQUE = "idx_wm_unique"
const val IDX_WM_CONV = "idx_wm_conv" const val IDX_WM_CONV = "idx_wm_conv"
const val IDX_REFLECTION_CREATED = "idx_reflection_created" const val IDX_REFLECTION_CREATED = "idx_reflection_created"
@@ -55,90 +47,34 @@ internal object Schema {
/** /**
* Прогоняет миграцию схемы до [CURRENT_VERSION] на пустой или существующей БД. * Прогоняет миграцию схемы до [CURRENT_VERSION] на пустой или существующей БД.
* *
* Версия хранится в `PRAGMA user_version` (стандартный SQLite-механизм,
* 32-bit int в заголовке БД — без своей таблицы). Каждая миграция —
* блок DDL+данных под номером `fromV+1`, выполняется в транзакции.
*
* Гарантии: * Гарантии:
* - идемпотентность: повторный вызов после достижения текущей версии — * - идемпотентность: `CREATE TABLE/INDEX IF NOT EXISTS` — безопасно на
* no-op (`PRAGMA user_version` совпадает с CURRENT_VERSION); * уже-мигрированной БД;
* - атомарность: каждая миграция в BEGIN/COMMIT — упал посреди → * - атомарность: каждая миграция в BEGIN/COMMIT — упал посреди →
* PRAGMA остаётся на предыдущей версии, БД консистентна; * ROLLBACK оставит БД консистентной.
* - CREATE TABLE/INDEX через `IF NOT EXISTS` — безопасно на partially- *
* migrated БД (если ручной ROLLBACK оставил схему в полупосаженном виде). **NOTE**: в сплит-мире user_version как gate перестал работать
* (соседний `:journal-ksqlite` тоже ставит user_version=1, и второй
* вызов short-circuit'ит). Поэтому migrate() просто прогоняет DDL
* idempotently; координация multi-module миграций — ответственность
* вызывающего (см. [KsqliteStores.open] / [KsqliteStores.inMemory]).
*/ */
fun migrate(conn: SQLiteConnection) { fun migrate(conn: SQLiteConnection) {
val current = readUserVersion(conn)
if (current >= CURRENT_VERSION) return
// Миграции строго последовательны — каждая стартует с (current) и
// выставляет user_version = current+1 в конце (внутри транзакции).
if (current < 1) {
conn.exec("BEGIN") conn.exec("BEGIN")
try { try {
conn.exec(v1ConversationDdl)
conn.exec(v1MessageDdl)
conn.exec(v1WorkingMemoryDdl) conn.exec(v1WorkingMemoryDdl)
conn.exec(v1ReflectionDdl) conn.exec(v1ReflectionDdl)
conn.exec(v1IndexesDdl) conn.exec(v1IndexesDdl)
writeUserVersion(conn, 1)
conn.exec("COMMIT") conn.exec("COMMIT")
} catch (t: Throwable) { } catch (t: Throwable) {
runCatching { conn.exec("ROLLBACK") } runCatching { conn.exec("ROLLBACK") }
throw t throw t
} }
} }
// sanity check — после всех миграций обязаны достичь CURRENT_VERSION
check(readUserVersion(conn) == CURRENT_VERSION) {
"Schema migration failed to reach version $CURRENT_VERSION"
}
}
private fun readUserVersion(conn: SQLiteConnection): Int {
// user_version живёт в sqlite_master-метаданных; PRAGMA возвращает его
// как row column "user_version". Используем прямой SELECT к internal
// pragma function через ksqlite (prepared + executeQuery).
var version = 0
conn.prepare("PRAGMA user_version").use { stmt ->
stmt.executeQuery().use { rs ->
if (rs.next()) {
version = (rs.getLong(0) ?: 0L).toInt()
}
}
}
return version
}
private fun writeUserVersion(conn: SQLiteConnection, version: Int) {
// SQLite PRAGMA с literal-аргументом нельзя параметризовать через `?`,
// поэтому собираем SQL строкой (значение контролируемое, не user input).
conn.exec("PRAGMA user_version = $version")
}
// ───── DDL миграций ───── // ───── DDL миграций ─────
// Ниже идут блоки по одной миграции. Конкатенация в [v1*] — потому что // v1 — начальная схема модуля. Пара `CREATE IF NOT EXISTS` →
// v1 — начальная схема (нет pre-existing DB с user_version=0 в проде, // миграция idempotent без всяких version-checks.
// но мы поддерживаем эту ветку на случай dev-БД под `agentik-dev.db`).
private val v1ConversationDdl = """
CREATE TABLE IF NOT EXISTS $TABLE_CONVERSATION (
$COL_ID TEXT NOT NULL PRIMARY KEY,
$COL_TITLE TEXT,
$COL_IS_TEMPORAL INTEGER NOT NULL DEFAULT 0,
$COL_CREATED_AT INTEGER NOT NULL,
$COL_UPDATED_AT INTEGER NOT NULL
);
"""
private val v1MessageDdl = """
CREATE TABLE IF NOT EXISTS $TABLE_MESSAGE (
$COL_ID TEXT NOT NULL PRIMARY KEY,
$COL_CONVERSATION_ID TEXT NOT NULL,
$COL_KIND TEXT NOT NULL,
$COL_PAYLOAD_JSON TEXT NOT NULL,
$COL_CREATED_AT INTEGER NOT NULL
);
"""
private val v1WorkingMemoryDdl = """ private val v1WorkingMemoryDdl = """
CREATE TABLE IF NOT EXISTS $TABLE_WORKING_MEMORY ( CREATE TABLE IF NOT EXISTS $TABLE_WORKING_MEMORY (
@@ -165,14 +101,6 @@ internal object Schema {
""" """
private val v1IndexesDdl = """ private val v1IndexesDdl = """
CREATE INDEX IF NOT EXISTS $IDX_CONV_UPDATED
ON $TABLE_CONVERSATION($COL_UPDATED_AT DESC);
-- Главный hot-path индекс для list/conversation: фильтр по conv +
-- сортировка по created_at (используется list(), cascade-delete, etc.)
CREATE INDEX IF NOT EXISTS $IDX_MSG_CONV
ON $TABLE_MESSAGE($COL_CONVERSATION_ID, $COL_CREATED_AT);
-- Working memory: гарантия уникального order_idx внутри conv'а -- Working memory: гарантия уникального order_idx внутри conv'а
-- (порядок имеет значение — compaction полагается на монотонность). -- (порядок имеет значение — compaction полагается на монотонность).
CREATE UNIQUE INDEX IF NOT EXISTS $IDX_WM_UNIQUE CREATE UNIQUE INDEX IF NOT EXISTS $IDX_WM_UNIQUE
@@ -1,96 +0,0 @@
package pw.binom.agentik.storage.ksqlite
import kotlinx.coroutines.flow.toList
import kotlinx.coroutines.test.runTest
import pw.binom.agentik.journal.Content
import pw.binom.agentik.journal.MessageRecord
import pw.binom.agentik.journal.TurnTokens
import kotlin.test.AfterTest
import kotlin.test.BeforeTest
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.time.Instant
class KsqliteMessageStoreTest {
private lateinit var stores: KsqliteStores
@BeforeTest
fun setup() {
stores = KsqliteStores.inMemory("msg-${kotlin.random.Random.nextLong()}")
kotlinx.coroutines.runBlocking {
stores.conversations.upsert(
pw.binom.agentik.journal.ConversationRecord(
"conv1", null, false,
Instant.parse("2026-09-15T10:00:00Z"),
Instant.parse("2026-09-15T10:00:00Z"),
)
)
}
}
@AfterTest
fun tearDown() = stores.close()
@Test
fun testAppendUserAndRetrieve() = runTest {
stores.messages.append(MessageRecord.UserMessage(
id = "m1",
conversationId = "conv1",
content = listOf(Content.Text("hello")),
createdAt = Instant.parse("2026-09-15T10:01:00Z"),
context = null,
))
val list = stores.messages.listFlow("conv1", Instant.DISTANT_PAST).toList()
assertEquals(1, list.size)
val msg = list[0]
assertEquals("m1", msg.id)
assertEquals(MessageRecord.UserMessage::class, msg::class)
}
@Test
fun testAppendAssistantWithTokens() = runTest {
stores.messages.append(MessageRecord.AssistantMessage(
id = "m1",
conversationId = "conv1",
content = listOf(Content.Text("hi")),
createdAt = Instant.parse("2026-09-15T10:01:00Z"),
tokens = TurnTokens(input = 50, output = 30),
))
val list = stores.messages.listFlow("conv1", Instant.DISTANT_PAST).toList()
assertEquals(1, list.size)
val msg = list[0] as MessageRecord.AssistantMessage
assertEquals(TurnTokens(input = 50, output = 30), msg.tokens)
}
@Test
fun testListAfterFiltersByTimestamp() = 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")
stores.messages.append(MessageRecord.UserMessage("m1", "conv1", listOf(Content.Text("a")), t1, null))
stores.messages.append(MessageRecord.UserMessage("m2", "conv1", listOf(Content.Text("b")), t2, null))
stores.messages.append(MessageRecord.UserMessage("m3", "conv1", listOf(Content.Text("c")), t3, null))
val after = stores.messages.list("conv1", after = t1, offset = 0, limit = 10)
assertEquals(2, after.size)
assertEquals(listOf("m2", "m3"), after.map { it.id })
}
@Test
fun testListFlowReturnsAllInOrder() = runTest {
val t = Instant.parse("2026-09-15T10:00:00Z")
for (i in 1..3) stores.messages.append(
MessageRecord.UserMessage("m$i", "conv1", listOf(Content.Text("x$i")), t + kotlin.time.Duration.parse("PT${i}S"), null)
)
assertEquals(listOf("m1", "m2", "m3"), stores.messages.listFlow("conv1", Instant.DISTANT_PAST).toList().map { it.id })
}
@Test
fun testClearRemovesByConversation() = runTest {
stores.messages.append(MessageRecord.UserMessage("m1", "conv1", listOf(Content.Text("a")), Instant.parse("2026-09-15T10:00:00Z"), null))
stores.messages.append(MessageRecord.UserMessage("m2", "conv2", listOf(Content.Text("b")), Instant.parse("2026-09-15T10:00:00Z"), null))
stores.conversations.delete("conv1")
assertEquals(emptyList(), stores.messages.listFlow("conv1", Instant.DISTANT_PAST).toList())
assertEquals(1, stores.messages.listFlow("conv2", Instant.DISTANT_PAST).toList().size)
}
}
@@ -1,105 +0,0 @@
package pw.binom.agentik.storage.ksqlite
import kotlinx.coroutines.test.runTest
import pw.binom.agentik.journal.ConversationRecord
import kotlin.test.AfterTest
import kotlin.test.BeforeTest
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertNotNull
import kotlin.test.assertNull
import kotlin.test.assertTrue
import kotlin.time.Duration
import kotlin.time.Instant
class KsqliteMutableConversationStoreTest {
private lateinit var stores: KsqliteStores
@BeforeTest
fun setup() {
stores = KsqliteStores.inMemory("conv-${kotlin.random.Random.nextLong()}")
}
@AfterTest
fun tearDown() = stores.close()
private fun rec(id: String, title: String? = null, ts: Instant = Instant.parse("2026-09-15T10:00:00Z")) =
ConversationRecord(id, title, false, ts, ts)
@Test
fun testUpsertAndGetRoundtrip() = runTest {
stores.conversations.upsert(rec("c1", "test"))
assertEquals(rec("c1", "test"), stores.conversations.get("c1"))
}
@Test
fun testGetReturnsNullForMissing() = runTest {
assertNull(stores.conversations.get("nope"))
}
@Test
fun testDeleteRemovesAndReturnsTrue() = runTest {
stores.conversations.upsert(rec("c1"))
assertTrue(stores.conversations.delete("c1"))
assertNull(stores.conversations.get("c1"))
assertEquals(false, stores.conversations.delete("c1"))
}
@Test
fun testListSortsByUpdatedAtDesc() = runTest {
val t0 = Instant.parse("2026-09-15T10:00:00Z")
stores.conversations.upsert(rec("c1", ts = t0))
stores.conversations.upsert(rec("c2", ts = t0))
stores.conversations.upsert(rec("c3", ts = t0))
stores.conversations.upsert(rec("c4", ts = t0))
stores.conversations.touch("c2", t0 + Duration.parse("PT60S"))
stores.conversations.touch("c3", t0 + Duration.parse("PT120S"))
stores.conversations.touch("c4", t0 + Duration.parse("PT180S"))
val page = stores.conversations.list(offset = 0, limit = 4)
assertEquals(listOf("c4", "c3", "c2", "c1"), page.map { it.id })
}
@Test
fun testListRespectsOffsetAndLimit() = runTest {
val t0 = Instant.parse("2026-09-15T10:00:00Z")
for (i in 1..5) stores.conversations.upsert(rec("c$i", ts = t0 + Duration.parse("PT${i}S")))
val p0 = stores.conversations.list(offset = 0, limit = 2)
assertEquals(2, p0.size)
val p2 = stores.conversations.list(offset = 4, limit = 2)
assertEquals(1, p2.size)
}
@Test
fun testRenameUpdatesTitleAndUpdatedAt() = runTest {
stores.conversations.upsert(rec("c1"))
val newTs = stores.conversations.rename("c1", "new title")
assertNotNull(newTs)
assertEquals("new title", stores.conversations.get("c1")?.title)
}
@Test
fun testRenameWithNullClearsTitle() = runTest {
stores.conversations.upsert(rec("c1", "old"))
stores.conversations.rename("c1", null)
assertNull(stores.conversations.get("c1")?.title)
}
@Test
fun testRenameReturnsNullForMissing() = runTest {
assertNull(stores.conversations.rename("nope", "x"))
}
@Test
fun testTouchUpdatesUpdatedAtOnly() = runTest {
val t0 = Instant.parse("2026-09-15T10:00:00Z")
val t1 = Instant.parse("2026-09-15T10:01:00Z")
stores.conversations.upsert(ConversationRecord("c1", "title", false, t0, t0))
stores.conversations.touch("c1", t1)
val got = stores.conversations.get("c1")
assertEquals("title", got?.title)
assertEquals(t1, got?.updatedAt)
assertEquals(t0, got?.createdAt)
}
}
@@ -1,33 +1,35 @@
package pw.binom.agentik.storage.ksqlite package pw.binom.agentik.storage.ksqlite
import kotlinx.coroutines.test.runTest import kotlinx.coroutines.test.runTest
import pw.binom.agentik.context.WorkingMemoryEntry
import pw.binom.agentik.journal.Content
import pw.binom.agentik.journal.ksqlite.Schema as JournalSchema
import pw.binom.db.ksqlite.SQLiteConnection import pw.binom.db.ksqlite.SQLiteConnection
import kotlin.test.Test import kotlin.test.Test
import kotlin.test.assertEquals import kotlin.test.assertEquals
import kotlin.test.assertTrue import kotlin.test.assertTrue
/** /**
* Тесты на Schema.migrate(): * Тесты на Schema.migrate() в `:storage-ksqlite`:
* - fresh DB → создаются все 4 таблицы + индексы + user_version = CURRENT_VERSION; * - fresh DB → создаются `working_memory` + `reflection` + индексы;
* - уже мигрированная БД → migrate() идемпотентен (no-op, не падает на * - уже мигрированная БД → migrate() идемпотентен (no-op);
* повторных CREATE);
* - DB, открытая напрямую через SQLiteConnection (минуя KsqliteStores), * - DB, открытая напрямую через SQLiteConnection (минуя KsqliteStores),
* migrate() приводит её в боевое состояние. * migrate() приводит её в боевое состояние.
* *
* Также проверяем что наличие индекса idx_msg_conv (conversation_id + * `conversation` + `message` тестируются в
* created_at) — обязательный hot-path для list()/cascade-delete. * `pw.binom.agentik.journal.ksqlite.SchemaMigrationTest`.
*/ */
class SchemaMigrationTest { class SchemaMigrationTest {
@Test @Test
fun `fresh DB gets all tables indexes and CURRENT_VERSION`() = runTest { fun `fresh DB gets working_memory and reflection tables and indexes`() = runTest {
val conn = SQLiteConnection.memory("mig-fresh-${kotlin.random.Random.nextLong()}") val conn = SQLiteConnection.memory("mig-fresh-${kotlin.random.Random.nextLong()}")
try { try {
// Поднимаем ОБА schema — то же делает KsqliteStores.open().
JournalSchema.migrate(conn)
Schema.migrate(conn) Schema.migrate(conn)
for (table in listOf( for (table in listOf(
Schema.TABLE_CONVERSATION,
Schema.TABLE_MESSAGE,
Schema.TABLE_WORKING_MEMORY, Schema.TABLE_WORKING_MEMORY,
Schema.TABLE_REFLECTION, Schema.TABLE_REFLECTION,
)) { )) {
@@ -35,8 +37,6 @@ class SchemaMigrationTest {
} }
for (index in listOf( for (index in listOf(
Schema.IDX_CONV_UPDATED,
Schema.IDX_MSG_CONV,
Schema.IDX_WM_UNIQUE, Schema.IDX_WM_UNIQUE,
Schema.IDX_WM_CONV, Schema.IDX_WM_CONV,
Schema.IDX_REFLECTION_CREATED, Schema.IDX_REFLECTION_CREATED,
@@ -44,8 +44,6 @@ class SchemaMigrationTest {
)) { )) {
assertTrue(indexExists(conn, index), "index '$index' should exist after migrate()") assertTrue(indexExists(conn, index), "index '$index' should exist after migrate()")
} }
assertEquals(Schema.CURRENT_VERSION, readUserVersion(conn))
} finally { } finally {
conn.close() conn.close()
} }
@@ -55,13 +53,15 @@ class SchemaMigrationTest {
fun `migrate is idempotent on already-migrated DB`() = runTest { fun `migrate is idempotent on already-migrated DB`() = runTest {
val conn = SQLiteConnection.memory("mig-idem-${kotlin.random.Random.nextLong()}") val conn = SQLiteConnection.memory("mig-idem-${kotlin.random.Random.nextLong()}")
try { try {
JournalSchema.migrate(conn)
Schema.migrate(conn) Schema.migrate(conn)
val versionAfterFirst = readUserVersion(conn)
// повторный вызов не должен ни упасть, ни изменить версию, ни // повторный вызов не должен ни упасть, ни пересоздать таблицы
// пересоздать таблицы/индексы (CREATE IF NOT EXISTS — no-op) // (CREATE IF NOT EXISTS — no-op)
Schema.migrate(conn) Schema.migrate(conn)
assertEquals(versionAfterFirst, readUserVersion(conn)) Schema.migrate(conn)
assertTrue(tableExists(conn, Schema.TABLE_WORKING_MEMORY))
assertTrue(tableExists(conn, Schema.TABLE_REFLECTION))
} finally { } finally {
conn.close() conn.close()
} }
@@ -69,63 +69,31 @@ class SchemaMigrationTest {
@Test @Test
fun `raw SQLiteConnection plus migrate gives working bundle`() = runTest { fun `raw SQLiteConnection plus migrate gives working bundle`() = runTest {
// Имитируем сценарий: существующая БД без schema, открываем через
// ksqlite и прогоняем migrate руками (тот же путь, что в
// KsqliteStores.open, но без зависимости от фабрики).
val conn = SQLiteConnection.memory("mig-bundle-${kotlin.random.Random.nextLong()}") val conn = SQLiteConnection.memory("mig-bundle-${kotlin.random.Random.nextLong()}")
JournalSchema.migrate(conn)
Schema.migrate(conn) Schema.migrate(conn)
// Сборка bundle через internal-конструктор — KsqliteStores primary
// constructor internal, тест в том же модуле и может его звать.
val stores = KsqliteStores( val stores = KsqliteStores(
connection = conn, connection = conn,
conversations = KsqliteMutableConversationStore(conn), conversations = pw.binom.agentik.journal.ksqlite.KsqliteMutableConversationStore(conn),
messages = KsqliteMessageStore(conn), messages = pw.binom.agentik.journal.ksqlite.KsqliteJournalStore(conn),
workingMemory = KsqliteWorkingMemoryStore(conn), workingMemory = KsqliteWorkingMemoryStore(conn),
reflections = KsqliteReflectionStore(conn), reflections = KsqliteReflectionStore(conn),
) )
try { try {
// bundle работает end-to-end — conversation upsert + message append + // bundle работает end-to-end — working memory append. Никаких
// list. Никаких "no such table" или подобного. // "no such table" или подобного.
stores.conversations.upsert( stores.workingMemory.append(
pw.binom.agentik.journal.ConversationRecord( conversationId = "c1",
id = "c1", title = "t", isTemporal = false, entry = WorkingMemoryEntry.Summary(text = "init"),
createdAt = kotlin.time.Instant.parse("2026-09-15T10:00:00Z"), now = kotlin.time.Instant.parse("2026-09-15T10:00:00Z"),
updatedAt = kotlin.time.Instant.parse("2026-09-15T10:00:00Z"),
) )
) assertEquals(1, stores.workingMemory.list("c1").size)
stores.messages.append(
pw.binom.agentik.journal.MessageRecord.UserMessage(
id = "m1", conversationId = "c1",
content = listOf(pw.binom.agentik.journal.Content.Text("hi")),
createdAt = kotlin.time.Instant.parse("2026-09-15T10:00:01Z"),
)
)
val got = stores.messages.list("c1", kotlin.time.Instant.DISTANT_PAST, offset = 0, limit = 10)
assertEquals(1, got.size)
assertEquals("m1", got[0].id)
} finally { } finally {
// Закрываем store'ы → они закроют свои pre-prepared statements
// (StmtHolder.finalize увидит isOpen == false и не полезет в
// нативный sqlite3_finalize с уже-разрушенным db mutex).
stores.close() stores.close()
} }
} }
@Test
fun `idx_msg_conv covers conversation_id and created_at columns`() = runTest {
// Проверяем что индекс действительно покрывает обе колонки — без
// этого list()/cascade-delete будут делать full-scan по message.
val conn = SQLiteConnection.memory("mig-idx-${kotlin.random.Random.nextLong()}")
try {
Schema.migrate(conn)
val cols = indexColumns(conn, Schema.IDX_MSG_CONV)
assertEquals(listOf(Schema.COL_CONVERSATION_ID, Schema.COL_CREATED_AT), cols)
} finally {
conn.close()
}
}
private fun tableExists(conn: SQLiteConnection, name: String): Boolean { private fun tableExists(conn: SQLiteConnection, name: String): Boolean {
conn.prepare( conn.prepare(
"SELECT 1 FROM sqlite_master WHERE type = 'table' AND name = ?" "SELECT 1 FROM sqlite_master WHERE type = 'table' AND name = ?"
@@ -143,31 +111,4 @@ class SchemaMigrationTest {
stmt.executeQuery().use { rs -> return rs.next() } stmt.executeQuery().use { rs -> return rs.next() }
} }
} }
private fun indexColumns(conn: SQLiteConnection, indexName: String): List<String> {
// PRAGMA index_info возвращает одну строку на колонку индекса
// (seqno, cid, name). Параметризовать через `?` нельзя — собираем
// строку (name — контролируемая константа, не user input).
val cols = mutableListOf<String>()
conn.prepare("PRAGMA index_info($indexName)").use { stmt ->
stmt.executeQuery().use { rs ->
while (rs.next()) {
rs.getText(2)?.let(cols::add)
}
}
}
return cols
}
private fun readUserVersion(conn: SQLiteConnection): Int {
var version = 0
conn.prepare("PRAGMA user_version").use { stmt ->
stmt.executeQuery().use { rs ->
if (rs.next()) {
version = (rs.getLong(0) ?: 0L).toInt()
}
}
}
return version
}
} }
+26
View File
@@ -0,0 +1,26 @@
plugins {
alias(libs.plugins.kotlin.multiplatform)
}
kotlin {
jvmToolchain(21)
jvm()
macosX64()
macosArm64()
iosX64()
iosArm64()
iosSimulatorArm64()
linuxX64()
linuxArm64()
mingwX64()
sourceSets {
commonMain.dependencies {
api(libs.kotlinx.coroutines.core)
}
commonTest.dependencies {
implementation(kotlin("test"))
implementation(libs.kotlinx.coroutines.test)
}
}
}
@@ -0,0 +1,29 @@
package pw.binom.agentik.vectorindex
interface MutableVectorIndexStore : VectorIndexStore {
/**
* Добавляет запись в индекс.
*
* @param id Идентификатор записи. Если `null` — реализация генерирует
* его сама (формат остаётся на усмотрение реализации).
* @param embedding Вектор. Размер должен совпадать с [VectorIndexStore.dimension].
* @param payload Пользовательские данные (любая сериализуемая структура —
* реализация хранит opaque string и не интерпретирует).
* @return Добавленная запись (с сгенерированным `id`, если был передан `null`).
* @throws VectorIndexAlreadyExistsException запись с указанным `id` уже есть в индексе.
*/
suspend fun add(id: String?, embedding: FloatArray, payload: String?): VectorIndex
/**
* Удаляет запись из индекса.
* @return `true`, если запись была удалена, `false`, если записи с таким `id` не было.
*/
suspend fun delete(id: String): Boolean
/**
* Полностью очищает индекс.
* @return количество удалённых записей.
*/
suspend fun clear(): Long
}
@@ -0,0 +1,37 @@
package pw.binom.agentik.vectorindex
/**
* Запись в vector-индексе: идентификатор + embedding + опциональный
* пользовательский payload.
*
* НЕ `data class` потому что [embedding] это `FloatArray`, а дефолтные
* `equals`/`hashCode` data-class'а сравнивают массивы по reference (а не
* по содержимому) — два `VectorIndex` с идентичными числами оказались бы
* `!=`. Вместо этого [equals]/[hashCode]/[toString] переопределены руками.
*/
class VectorIndex(
val id: String,
val embedding: FloatArray,
val payload: String?,
) {
override fun equals(other: Any?): Boolean {
if (this === other) return true
if (other == null || this::class != other::class) return false
other as VectorIndex
if (id != other.id) return false
if (!embedding.contentEquals(other.embedding)) return false
if (payload != other.payload) return false
return true
}
override fun hashCode(): Int {
var result = id.hashCode()
result = 31 * result + embedding.contentHashCode()
result = 31 * result + (payload?.hashCode() ?: 0)
return result
}
override fun toString(): String =
"VectorIndex(id='$id', embedding.size=${embedding.size}, payload=$payload)"
}
@@ -0,0 +1,3 @@
package pw.binom.agentik.vectorindex
class VectorIndexAlreadyExistsException : Exception()
@@ -0,0 +1,22 @@
package pw.binom.agentik.vectorindex
/**
* Интерфейс read-only доступа к ANN-индексу. Реализации могут держать
* native resources (mmap'нутые файлы, sqlite handles) — поэтому [close]
* обязателен и наследуется от [AutoCloseable].
*/
interface VectorIndexStore : AutoCloseable {
/** Размерность embeddings, которые принимает этот индекс. */
val dimension: Int
/** Текущее количество записей в индексе. */
suspend fun getSize(): Long
/**
* Возвращает top-[limit] ближайших к [embedding] записей, отсортированных
* по убыванию [VectorSearchResult.score]. Если в индексе меньше [limit]
* записей — возвращается меньший список.
*/
suspend fun search(embedding: FloatArray, limit: Int): List<VectorSearchResult>
}
@@ -0,0 +1,6 @@
package pw.binom.agentik.vectorindex
data class VectorSearchResult(
val index: VectorIndex,
val score: Float,
)
+25
View File
@@ -0,0 +1,25 @@
plugins {
alias(libs.plugins.kotlin.multiplatform)
}
kotlin {
jvmToolchain(21)
// JVector (Datadog) публикуется только под JVM — нет KMP-таргетов.
// Если в будущем понадобится натив — отдельная реализация в vector-index-*.
jvm()
sourceSets {
commonMain.dependencies {
api(project(":vector-index-api"))
implementation(libs.kotlinx.coroutines.core)
}
commonTest.dependencies {
implementation(kotlin("test"))
implementation(libs.kotlinx.coroutines.test)
}
jvmMain.dependencies {
implementation(libs.jvector)
}
}
}
@@ -0,0 +1,199 @@
package pw.binom.agentik.vectorindex.jvector
import io.github.jbellis.jvector.graph.GraphIndexBuilder
import io.github.jbellis.jvector.graph.GraphSearcher
import io.github.jbellis.jvector.graph.ListRandomAccessVectorValues
import io.github.jbellis.jvector.graph.OnHeapGraphIndex
import io.github.jbellis.jvector.graph.SearchResult
import io.github.jbellis.jvector.graph.similarity.BuildScoreProvider
import io.github.jbellis.jvector.util.Bits
import io.github.jbellis.jvector.vector.VectorizationProvider
import io.github.jbellis.jvector.vector.VectorSimilarityFunction
import pw.binom.agentik.vectorindex.MutableVectorIndexStore
import pw.binom.agentik.vectorindex.VectorIndex
import pw.binom.agentik.vectorindex.VectorIndexAlreadyExistsException
import pw.binom.agentik.vectorindex.VectorSearchResult
import java.util.UUID
import java.util.concurrent.locks.ReentrantReadWriteLock
import kotlin.concurrent.read
import kotlin.concurrent.write
/**
* In-RAM [MutableVectorIndexStore] поверх JVector (Datadog ANN library, jvm-only).
*
* Семантика хранения: всё держится в heap'е — [OnHeapGraphIndex] + mapping
* id ↔ ordinal + payload. Граф перестраивается с нуля на каждом
* [add]/[delete]/[clear] (для 10K vectors <100ms — паттерн скопирован из
* `:memory-vector/JVectorMemoryIndex`).
*
* **Persist НЕ поддерживается** в этой версии. Если понадобится — пара
* `(JVectorVectorIndexStore, KVectorStore)` с SQLite BLOB для embeddings
* и graph-rebuild на старте. Подход с `OnDiskGraphIndex` отброшен: требует
* Feature `INLINE_VECTORS`, который в JVector 3.x конфигурируется отдельно
* и нестабилен.
*
* Потокобезопасность: [ReentrantReadWriteLock] — параллельные [search] ок,
* [add]/[delete]/[clear] — эксклюзивно (как в `:memory-vector/JVectorMemoryIndex`).
*/
class JVectorVectorIndexStore(
override val dimension: Int,
seedEntries: List<SeedEntry> = emptyList(),
) : MutableVectorIndexStore {
/**
* Запись для seed'а индекса при конструировании.
* Полезно, когда граф надо построить сразу из уже-имеющихся данных
* (например, при старте агента из snapshot'а).
*/
data class SeedEntry(
val id: String,
val embedding: FloatArray,
val payload: String?,
)
init {
require(seedEntries.all { it.embedding.size == dimension }) {
"all seed embeddings must have dimension=$dimension"
}
require(seedEntries.map { it.id }.toSet().size == seedEntries.size) {
"duplicate ids in seedEntries"
}
}
private val rwLock = ReentrantReadWriteLock()
private val vts = VectorizationProvider.getInstance().getVectorTypeSupport()
private val similarity = VectorSimilarityFunction.COSINE
// In-RAM state. Защищён rwLock.
private val idToOrdinal = LinkedHashMap<String, Int>()
private val ordinalToId = ArrayList<String>(seedEntries.size + 16)
private val ordinalToVector = ArrayList<FloatArray>(seedEntries.size + 16)
private val ordinalToPayload = ArrayList<String?>(seedEntries.size + 16)
private val deleted = java.util.BitSet()
private var graph: OnHeapGraphIndex? = null
init {
seedEntries.forEach { entry ->
val ord = ordinalToId.size
idToOrdinal[entry.id] = ord
ordinalToId.add(entry.id)
ordinalToVector.add(entry.embedding)
ordinalToPayload.add(entry.payload)
}
if (ordinalToId.isNotEmpty()) {
graph = rebuildFromScratch()
}
}
override suspend fun getSize(): Long = rwLock.read {
(ordinalToId.size - deleted.cardinality()).toLong()
}
override suspend fun add(
id: String?,
embedding: FloatArray,
payload: String?,
): VectorIndex = rwLock.write {
require(embedding.size == dimension) {
"embedding size ${embedding.size} != dimension $dimension"
}
val resolvedId = id ?: UUID.randomUUID().toString()
if (idToOrdinal.containsKey(resolvedId)) {
throw VectorIndexAlreadyExistsException()
}
val ord = ordinalToId.size
idToOrdinal[resolvedId] = ord
ordinalToId.add(resolvedId)
ordinalToVector.add(embedding)
ordinalToPayload.add(payload)
rebuildAndSwapGraph()
VectorIndex(id = resolvedId, embedding = embedding, payload = payload)
}
override suspend fun delete(id: String): Boolean = rwLock.write {
val ord = idToOrdinal.remove(id) ?: return false
// ordinal в ordinalToId/Vector/Payload НЕ удаляем — JVector rebuild
// опирается на стабильные ordinal'ы; помечаем в BitSet и rebuild.
deleted.set(ord)
rebuildAndSwapGraph()
true
}
override suspend fun clear(): Long = rwLock.write {
val count = (ordinalToId.size - deleted.cardinality()).toLong()
idToOrdinal.clear()
ordinalToId.clear()
ordinalToVector.clear()
ordinalToPayload.clear()
deleted.clear()
graph?.close()
graph = null
count
}
override suspend fun search(
embedding: FloatArray,
limit: Int,
): List<VectorSearchResult> = rwLock.read {
require(embedding.size == dimension) {
"query size ${embedding.size} != dimension $dimension"
}
if (limit <= 0 || graph == null) return@read emptyList()
val activeOrdinals = (0 until ordinalToId.size).filter { !deleted.get(it) }
if (activeOrdinals.isEmpty()) return@read emptyList()
val vectors = activeOrdinals.map { vts.createFloatVector(ordinalToVector[it]) }
val ravv = ListRandomAccessVectorValues(vectors, dimension)
val queryVec = vts.createFloatVector(embedding)
val result: SearchResult = GraphSearcher.search(
queryVec,
limit.coerceAtMost(activeOrdinals.size),
ravv,
similarity,
graph!!,
Bits.ALL,
)
val nodes: Array<SearchResult.NodeScore> = result.getNodes()
nodes.map { ns ->
val realOrd = activeOrdinals[ns.node]
VectorSearchResult(
index = VectorIndex(
id = ordinalToId[realOrd],
embedding = ordinalToVector[realOrd],
payload = ordinalToPayload[realOrd],
),
score = ns.score,
)
}
}
override fun close() {
rwLock.write {
graph?.close()
graph = null
}
}
private fun rebuildAndSwapGraph() {
val newGraph = rebuildFromScratch()
val old = graph
graph = newGraph
old?.close()
}
private fun rebuildFromScratch(): OnHeapGraphIndex {
val activeOrdinals = (0 until ordinalToId.size).filter { !deleted.get(it) }
val vectors = activeOrdinals.map { vts.createFloatVector(ordinalToVector[it]) }
val ravv = ListRandomAccessVectorValues(vectors, dimension)
val bsp = BuildScoreProvider.randomAccessScoreProvider(ravv, similarity)
// Параметры графа по умолчанию (JVector README):
// M (max degree) = 16
// efConstruction = 100
// neighborOverflow = 1.2f
// alpha = 1.2f
// Для масштабов до ~10K vectors дают хороший баланс точность/скорость.
return GraphIndexBuilder(bsp, dimension, 16, 100, 1.2f, 1.2f).use { builder ->
builder.build(ravv)
}
}
}
@@ -0,0 +1,141 @@
package pw.binom.agentik.vectorindex.jvector
import kotlinx.coroutines.test.runTest
import pw.binom.agentik.vectorindex.VectorIndexAlreadyExistsException
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertFailsWith
import kotlin.test.assertNotNull
import kotlin.test.assertTrue
class JVectorVectorIndexStoreTest {
private fun vec(vararg values: Float) = values
@Test
fun `add with explicit id and search returns it as top-1`() = runTest {
val store = JVectorVectorIndexStore(dimension = 4)
try {
val a = store.add(id = "a", embedding = vec(1f, 0f, 0f, 0f), payload = "first")
assertEquals("a", a.id)
assertEquals("first", a.payload)
val results = store.search(embedding = vec(1f, 0f, 0f, 0f), limit = 1)
assertEquals(1, results.size)
assertEquals("a", results[0].index.id)
// cosine similarity of identical unit vectors = 1.0
assertTrue(results[0].score > 0.99f, "score=${results[0].score} should be ~1.0")
} finally {
store.close()
}
}
@Test
fun `add with null id auto-generates`() = runTest {
val store = JVectorVectorIndexStore(dimension = 2)
try {
val r1 = store.add(id = null, embedding = vec(1f, 0f), payload = null)
val r2 = store.add(id = null, embedding = vec(0f, 1f), payload = null)
assertNotNull(r1.id)
assertNotNull(r2.id)
assertTrue(r1.id != r2.id, "auto-generated ids must differ")
assertEquals(2L, store.getSize())
} finally {
store.close()
}
}
@Test
fun `add with duplicate explicit id throws`() = runTest {
val store = JVectorVectorIndexStore(dimension = 2)
try {
store.add(id = "dup", embedding = vec(1f, 0f), payload = null)
assertFailsWith<VectorIndexAlreadyExistsException> {
store.add(id = "dup", embedding = vec(0f, 1f), payload = null)
}
} finally {
store.close()
}
}
@Test
fun `delete removes and clears count`() = runTest {
val store = JVectorVectorIndexStore(dimension = 2)
try {
store.add(id = "x", embedding = vec(1f, 0f), payload = null)
store.add(id = "y", embedding = vec(0f, 1f), payload = null)
assertEquals(2L, store.getSize())
assertTrue(store.delete("x"))
assertEquals(1L, store.getSize())
assertEquals(false, store.delete("x"), "second delete returns false")
} finally {
store.close()
}
}
@Test
fun `clear empties and returns count`() = runTest {
val store = JVectorVectorIndexStore(dimension = 2)
try {
store.add(id = "a", embedding = vec(1f, 0f), payload = null)
store.add(id = "b", embedding = vec(0f, 1f), payload = null)
store.add(id = "c", embedding = vec(-1f, 0f), payload = null)
val cleared = store.clear()
assertEquals(3L, cleared)
assertEquals(0L, store.getSize())
assertEquals(emptyList(), store.search(vec(1f, 0f), limit = 5))
} finally {
store.close()
}
}
@Test
fun `search returns top-K sorted by descending score`() = runTest {
val store = JVectorVectorIndexStore(dimension = 3)
try {
// Один точный матч + два шумовых, отдалённых от query.
store.add(id = "exact", embedding = vec(1f, 0f, 0f), payload = null)
store.add(id = "noise1", embedding = vec(-1f, 0f, 0f), payload = null)
store.add(id = "noise2", embedding = vec(0f, 1f, 0f), payload = null)
val results = store.search(embedding = vec(1f, 0f, 0f), limit = 3)
assertEquals(3, results.size)
// Score-ы монотонно убывают.
assertTrue(results[0].score >= results[1].score)
assertTrue(results[1].score >= results[2].score)
assertEquals("exact", results[0].index.id)
} finally {
store.close()
}
}
@Test
fun `search with limit greater than index returns all entries`() = runTest {
val store = JVectorVectorIndexStore(dimension = 2)
try {
store.add(id = "a", embedding = vec(1f, 0f), payload = null)
store.add(id = "b", embedding = vec(0f, 1f), payload = null)
val results = store.search(vec(1f, 0f), limit = 100)
assertEquals(2, results.size)
} finally {
store.close()
}
}
@Test
fun `payload round-trips through search`() = runTest {
val store = JVectorVectorIndexStore(dimension = 2)
try {
store.add(id = "p", embedding = vec(1f, 0f), payload = """{"k":"v"}""")
val results = store.search(vec(1f, 0f), limit = 1)
assertEquals(1, results.size)
assertEquals("p", results[0].index.id)
assertEquals("""{"k":"v"}""", results[0].index.payload)
} finally {
store.close()
}
}
}
+32
View File
@@ -0,0 +1,32 @@
plugins {
alias(libs.plugins.kotlin.multiplatform)
}
// KMP-реализация :vector-index-api (MutableVectorIndexStore) поверх ksqlite.
// Brute-force ANN: SELECT всех записей + cosine similarity в Kotlin. Persistence —
// обычная SQLite БД. Подходит для small-to-medium масштабов (≤10K записей на
// embedding ~512d); для больших датасетов — sqlite-vec или JVector.
//
// Цели сборки — только те, для которых ksqlite 0.1.2 опубликован в Maven Central:
// jvm + linuxX64/Arm64 + mingwX64. Apple/iOS/tvOS/watchOS — НЕ публикуются;
// для них использовать JVM-only `:vector-index-jvector`.
kotlin {
jvmToolchain(21)
jvm()
linuxX64()
linuxArm64()
mingwX64()
sourceSets {
commonMain.dependencies {
// ksqlite 0.1.2 опубликован в Maven Central.
implementation("pw.binom.db:ksqlite:0.1.2")
api(project(":vector-index-api"))
}
commonTest.dependencies {
implementation(kotlin("test"))
implementation(libs.kotlinx.coroutines.test)
}
}
}
@@ -0,0 +1,245 @@
package pw.binom.agentik.vectorindex.ksqlite
import pw.binom.agentik.vectorindex.MutableVectorIndexStore
import pw.binom.agentik.vectorindex.VectorIndex
import pw.binom.agentik.vectorindex.VectorIndexAlreadyExistsException
import pw.binom.agentik.vectorindex.VectorSearchResult
import pw.binom.db.ksqlite.SQLiteConnection
import pw.binom.db.ksqlite.SQLitePreparedStatement
import pw.binom.db.ksqlite.SQLiteResultSet
import kotlin.uuid.Uuid
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock
import kotlinx.coroutines.withContext
import pw.binom.agentik.vectorindex.ksqlite.utils.toFloatArray
import pw.binom.agentik.vectorindex.ksqlite.utils.toLittleEndianBytes
import kotlin.math.sqrt
/**
* ksqlite-реализация [MutableVectorIndexStore].
*
* **Стратегия поиска: brute-force cosine similarity.**
* Один SELECT всех записей → декодирование BLOB → cosine sim в Kotlin →
* sort desc → top-[limit]. Просто и KMP-совместимо. Для масштабов >10K
* записей становится узким местом — переезжать на `:vector-index-jvector`
* или sqlite-vec (см. KDoc [Schema]).
*
* ## Lifecycle соединения
*
* Семантика владения connection'ом идентична другим ksqlite-store'ам
* (`pw.binom.agentik.journal.ksqlite.KsqliteJournalStore` и т.п.):
* - `KsqliteVectorIndexStore(dimension, path)` — открывает файловое
* соединение, прогоняет [Schema.migrate], закрывает в [close].
* - `KsqliteVectorIndexStore(dimension, connection)` — внешнее соединение,
* store НЕ закрывает его в [close].
* - `KsqliteVectorIndexStore.memory(dimension, name)` — in-memory, мигрирует,
* закрывает в [close].
*/
class KsqliteVectorIndexStore private constructor(
override val dimension: Int,
private val connection: SQLiteConnection,
private val ownsConnection: Boolean,
) : MutableVectorIndexStore {
init {
require(dimension > 0) { "dimension must be positive, got $dimension" }
Schema.migrate(connection)
}
/**
* Открывает файловое соединение через [SQLiteConnection.open] и берёт
* на себя его закрытие в [close].
*/
constructor(dimension: Int, path: String) : this(
dimension = dimension,
connection = SQLiteConnection.open(path = path),
ownsConnection = true,
)
/**
* Внешнее соединение — store НЕ закрывает его в [close].
*/
constructor(dimension: Int, connection: SQLiteConnection) : this(
dimension = dimension,
connection = connection,
ownsConnection = false,
)
private val mutex = Mutex()
// pre-prepare (см. KDoc KsqliteJournalStore — почему это критично против
// SIGSEGV в StmtHolder.finalize на закрытой connection).
private val existsStmt: SQLitePreparedStatement = connection.prepare(
"SELECT 1 FROM ${Schema.TABLE_VECTOR_INDEX} WHERE ${Schema.COL_ID} = ?"
)
private val insertStmt: SQLitePreparedStatement = connection.prepare(
"""
INSERT INTO ${Schema.TABLE_VECTOR_INDEX}
(${Schema.COL_ID}, ${Schema.COL_DIMENSION},
${Schema.COL_EMBEDDING}, ${Schema.COL_PAYLOAD})
VALUES (?, ?, ?, ?)
""".trimIndent()
)
private val deleteStmt: SQLitePreparedStatement = connection.prepare(
"DELETE FROM ${Schema.TABLE_VECTOR_INDEX} WHERE ${Schema.COL_ID} = ?"
)
private val clearStmt: SQLitePreparedStatement = connection.prepare(
"DELETE FROM ${Schema.TABLE_VECTOR_INDEX}"
)
private val countStmt: SQLitePreparedStatement = connection.prepare(
"SELECT COUNT(*) FROM ${Schema.TABLE_VECTOR_INDEX}"
)
private val allStmt: SQLitePreparedStatement = connection.prepare(
"""
SELECT ${Schema.COL_ID}, ${Schema.COL_EMBEDDING}, ${Schema.COL_PAYLOAD}
FROM ${Schema.TABLE_VECTOR_INDEX}
""".trimIndent()
)
override suspend fun getSize(): Long = withContext(Dispatchers.Default) {
mutex.withLock {
countStmt.reset()
countStmt.clearBindings()
countStmt.executeQuery().use { rs ->
if (rs.next()) rs.getLong(0) ?: 0L else 0L
}
}
}
override suspend fun add(
id: String?,
embedding: FloatArray,
payload: String?,
): VectorIndex = withContext(Dispatchers.Default) {
require(embedding.size == dimension) {
"embedding size ${embedding.size} != dimension $dimension"
}
val resolvedId = id ?: Uuid.random().toString()
mutex.withLock {
if (execExists(resolvedId)) {
throw VectorIndexAlreadyExistsException()
}
insertStmt.reset()
insertStmt.clearBindings()
insertStmt.bindText(1, resolvedId)
insertStmt.bindInt(2, dimension)
insertStmt.bindBlob(3, embedding.toLittleEndianBytes())
if (payload != null) insertStmt.bindText(4, payload) else insertStmt.bindNull(4)
insertStmt.executeUpdate()
VectorIndex(id = resolvedId, embedding = embedding, payload = payload)
}
}
override suspend fun delete(id: String): Boolean = withContext(Dispatchers.Default) {
mutex.withLock {
deleteStmt.reset()
deleteStmt.clearBindings()
deleteStmt.bindText(1, id)
deleteStmt.executeUpdate() > 0
}
}
override suspend fun clear(): Long = withContext(Dispatchers.Default) {
mutex.withLock {
val before = getSizeNoLock()
clearStmt.reset()
clearStmt.clearBindings()
clearStmt.executeUpdate()
before
}
}
override suspend fun search(
embedding: FloatArray,
limit: Int,
): List<VectorSearchResult> = withContext(Dispatchers.Default) {
require(embedding.size == dimension) {
"query size ${embedding.size} != dimension $dimension"
}
if (limit <= 0) return@withContext emptyList()
mutex.withLock {
allStmt.reset()
allStmt.clearBindings()
val scored = mutableListOf<ScoredRow>()
allStmt.executeQuery().use { rs ->
while (rs.next()) {
val row = rs.toRow() ?: continue
val score = cosineSimilarity(embedding, row.embedding)
scored.add(ScoredRow(row, score))
}
}
scored.sortByDescending { it.score }
scored.take(limit).map { it.toResult() }
}
}
override fun close() {
existsStmt.close()
insertStmt.close()
deleteStmt.close()
clearStmt.close()
countStmt.close()
allStmt.close()
if (ownsConnection) connection.close()
}
companion object {
fun memory(dimension: Int, name: String? = null): KsqliteVectorIndexStore =
KsqliteVectorIndexStore(
dimension = dimension,
connection = SQLiteConnection.memory(name),
ownsConnection = true,
)
}
private fun execExists(id: String): Boolean {
existsStmt.reset()
existsStmt.clearBindings()
existsStmt.bindText(1, id)
existsStmt.executeQuery().use { rs -> return rs.next() }
}
private fun getSizeNoLock(): Long {
countStmt.reset()
countStmt.clearBindings()
countStmt.executeQuery().use { rs ->
if (rs.next()) return rs.getLong(0) ?: 0L
}
return 0L
}
private fun SQLiteResultSet.toRow(): Row? {
val id = getText(0) ?: return null
val blob = getBlob(1) ?: return null
val payload = getText(2)
return Row(id = id, embedding = blob.toFloatArray(), payload = payload)
}
private data class Row(val id: String, val embedding: FloatArray, val payload: String?)
private data class ScoredRow(val row: Row, val score: Float) {
fun toResult(): VectorSearchResult = VectorSearchResult(
index = VectorIndex(id = row.id, embedding = row.embedding, payload = row.payload),
score = score,
)
}
}
/**
* Cosine similarity двух одинаковой длины векторов.
* Возвращает 0, если один из векторов — нулевой (безопасно для пустых запросов).
*/
private fun cosineSimilarity(a: FloatArray, b: FloatArray): Float {
require(a.size == b.size) { "dimension mismatch: ${a.size} vs ${b.size}" }
var dot = 0f
var normA = 0f
var normB = 0f
for (i in a.indices) {
dot += a[i] * b[i]
normA += a[i] * a[i]
normB += b[i] * b[i]
}
val denom = sqrt(normA) * sqrt(normB)
return if (denom == 0f) 0f else dot / denom
}
@@ -0,0 +1,78 @@
package pw.binom.agentik.vectorindex.ksqlite
import pw.binom.db.ksqlite.SQLiteConnection
/**
* Имена таблиц/колонок для ksqlite-бэкенда `:vector-index-api`.
*
* Хранит embeddings как BLOB (raw float32 LE, см. [pw.binom.agentik.vectorindex.ksqlite.utils.toLittleEndianBytes])
* + dimension для sanity-check на insert + опциональный payload как TEXT.
*
* Brute-force search не использует индексов на стороне SQLite — все embeddings
* загружаются одним SELECT и cosine similarity считается в Kotlin. Для
* масштабов >10K записей это становится узким местом; тогда мигрировать на
* `:vector-index-jvector` или sqlite-vec.
*/
internal object Schema {
/** Версия схемы. Увеличивать при ЛЮБОМ изменении DDL. */
const val CURRENT_VERSION: Int = 1
const val TABLE_VECTOR_INDEX = "vector_index"
const val COL_ID = "id"
const val COL_DIMENSION = "dimension"
const val COL_EMBEDDING = "embedding"
const val COL_PAYLOAD = "payload"
private val v1Ddl = """
CREATE TABLE IF NOT EXISTS $TABLE_VECTOR_INDEX (
$COL_ID TEXT NOT NULL PRIMARY KEY,
$COL_DIMENSION INTEGER NOT NULL,
$COL_EMBEDDING BLOB NOT NULL,
$COL_PAYLOAD TEXT
);
""".trimIndent()
/**
* Прогоняет миграцию схемы до [CURRENT_VERSION] на пустой или существующей БД.
*
* Версия хранится в `PRAGMA user_version` (стандартный SQLite-механизм,
* 32-bit int в заголовке БД — без своей таблицы). Каждая миграция —
* блок DDL под номером `fromV+1`, выполняется в транзакции. Если миграция
* упадёт посередине — `ROLLBACK` оставит БД на предыдущей версии.
*
* Идемпотентен: повторный вызов на уже мигрированной БД — no-op.
*/
fun migrate(conn: SQLiteConnection) {
val current = readUserVersion(conn)
if (current >= CURRENT_VERSION) return
conn.exec("BEGIN")
try {
if (current < 1) {
conn.exec(v1Ddl)
}
writeUserVersion(conn, CURRENT_VERSION)
conn.exec("COMMIT")
} catch (t: Throwable) {
runCatching { conn.exec("ROLLBACK") }
throw t
}
}
private fun readUserVersion(conn: SQLiteConnection): Int {
conn.prepare("PRAGMA user_version").use { stmt ->
stmt.executeQuery().use { rs ->
if (rs.next()) return rs.getLong(0)?.toInt() ?: 0
}
}
return 0
}
private fun writeUserVersion(conn: SQLiteConnection, version: Int) {
// SQLite PRAGMA с literal-аргументом нельзя параметризовать через `?`,
// поэтому собираем SQL строкой (значение контролируемое, не user input).
conn.exec("PRAGMA user_version = $version")
}
}
@@ -0,0 +1,36 @@
package pw.binom.agentik.vectorindex.ksqlite.utils
/**
* Multiplatform-safe сериализация [FloatArray] ↔ [ByteArray] в little-endian.
*
* Используется для хранения embeddings в SQLite BLOB-колонке. Float занимает
* ровно 4 байта (IEEE-754 binary32). На JVM little-endian совпадает с native
* byte order; на linux/macos/ios/mingw — то же самое. Выбран little-endian
* как наиболее распространённый вариант и для совместимости с sqlite-vec
* (см. `pw.binom.db.ksqlite.SQLitePreparedStatement.bindVector`).
*/
internal fun FloatArray.toLittleEndianBytes(): ByteArray {
val out = ByteArray(size * 4)
for (i in indices) {
val bits = this[i].toRawBits()
out[i * 4 + 0] = bits.toByte()
out[i * 4 + 1] = (bits ushr 8).toByte()
out[i * 4 + 2] = (bits ushr 16).toByte()
out[i * 4 + 3] = (bits ushr 24).toByte()
}
return out
}
internal fun ByteArray.toFloatArray(): FloatArray {
require(size % 4 == 0) { "byte array size $size not a multiple of 4" }
val out = FloatArray(size / 4)
for (i in out.indices) {
val off = i * 4
val bits = (this[off].toInt() and 0xFF).toLong() or
((this[off + 1].toInt() and 0xFF).toLong() shl 8) or
((this[off + 2].toInt() and 0xFF).toLong() shl 16) or
((this[off + 3].toInt() and 0xFF).toLong() shl 24)
out[i] = Float.fromBits(bits.toInt())
}
return out
}
@@ -0,0 +1,179 @@
package pw.binom.agentik.vectorindex.ksqlite
import kotlinx.coroutines.test.runTest
import pw.binom.agentik.vectorindex.VectorIndexAlreadyExistsException
import pw.binom.agentik.vectorindex.ksqlite.utils.toFloatArray
import pw.binom.db.ksqlite.SQLiteConnection
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertFailsWith
import kotlin.test.assertNotNull
import kotlin.test.assertTrue
class KsqliteVectorIndexStoreTest {
private fun vec(vararg values: Float) = values
private fun newStore(dim: Int = 4): KsqliteVectorIndexStore =
KsqliteVectorIndexStore.memory(dimension = dim, name = "test-${kotlin.random.Random.nextLong()}")
@Test
fun `add with explicit id and search returns it as top-1`() = runTest {
val store = newStore(4)
try {
val a = store.add(id = "a", embedding = vec(1f, 0f, 0f, 0f), payload = "first")
assertEquals("a", a.id)
assertEquals("first", a.payload)
val results = store.search(embedding = vec(1f, 0f, 0f, 0f), limit = 1)
assertEquals(1, results.size)
assertEquals("a", results[0].index.id)
assertTrue(results[0].score > 0.99f, "score=${results[0].score} should be ~1.0")
} finally {
store.close()
}
}
@Test
fun `add with null id auto-generates`() = runTest {
val store = newStore(2)
try {
val r1 = store.add(id = null, embedding = vec(1f, 0f), payload = null)
val r2 = store.add(id = null, embedding = vec(0f, 1f), payload = null)
assertNotNull(r1.id)
assertNotNull(r2.id)
assertTrue(r1.id != r2.id, "auto-generated ids must differ")
assertEquals(2L, store.getSize())
} finally {
store.close()
}
}
@Test
fun `add with duplicate explicit id throws`() = runTest {
val store = newStore(2)
try {
store.add(id = "dup", embedding = vec(1f, 0f), payload = null)
assertFailsWith<VectorIndexAlreadyExistsException> {
store.add(id = "dup", embedding = vec(0f, 1f), payload = null)
}
} finally {
store.close()
}
}
@Test
fun `delete removes and clears count`() = runTest {
val store = newStore(2)
try {
store.add(id = "x", embedding = vec(1f, 0f), payload = null)
store.add(id = "y", embedding = vec(0f, 1f), payload = null)
assertEquals(2L, store.getSize())
assertTrue(store.delete("x"))
assertEquals(1L, store.getSize())
assertEquals(false, store.delete("x"), "second delete returns false")
// После delete можно заново add с тем же id.
store.add(id = "x", embedding = vec(-1f, 0f), payload = null)
assertEquals(2L, store.getSize())
} finally {
store.close()
}
}
@Test
fun `clear empties and returns count`() = runTest {
val store = newStore(2)
try {
store.add(id = "a", embedding = vec(1f, 0f), payload = null)
store.add(id = "b", embedding = vec(0f, 1f), payload = null)
store.add(id = "c", embedding = vec(-1f, 0f), payload = null)
val cleared = store.clear()
assertEquals(3L, cleared)
assertEquals(0L, store.getSize())
assertEquals(emptyList(), store.search(vec(1f, 0f), limit = 5))
} finally {
store.close()
}
}
@Test
fun `search returns top-K sorted by descending score`() = runTest {
val store = newStore(3)
try {
store.add(id = "exact", embedding = vec(1f, 0f, 0f), payload = null)
store.add(id = "noise1", embedding = vec(-1f, 0f, 0f), payload = null)
store.add(id = "noise2", embedding = vec(0f, 1f, 0f), payload = null)
val results = store.search(embedding = vec(1f, 0f, 0f), limit = 3)
assertEquals(3, results.size)
assertTrue(results[0].score >= results[1].score)
assertTrue(results[1].score >= results[2].score)
assertEquals("exact", results[0].index.id)
} finally {
store.close()
}
}
@Test
fun `search with limit greater than index returns all entries`() = runTest {
val store = newStore(2)
try {
store.add(id = "a", embedding = vec(1f, 0f), payload = null)
store.add(id = "b", embedding = vec(0f, 1f), payload = null)
val results = store.search(vec(1f, 0f), limit = 100)
assertEquals(2, results.size)
} finally {
store.close()
}
}
@Test
fun `payload round-trips through search`() = runTest {
val store = newStore(2)
try {
store.add(id = "p", embedding = vec(1f, 0f), payload = """{"k":"v"}""")
val results = store.search(vec(1f, 0f), limit = 1)
assertEquals(1, results.size)
assertEquals("p", results[0].index.id)
assertEquals("""{"k":"v"}""", results[0].index.payload)
} finally {
store.close()
}
}
@Test
fun `embedding BLOB round-trips bit-exact through raw SQL`() = runTest {
// Открываем свой connection, создаём store, пишем, потом читаем через
// raw SQL на том же connection — проверяем что FloatArray<->BLOB encoding
// round-trip'ит без потерь.
val conn = SQLiteConnection.memory("roundtrip-${kotlin.random.Random.nextLong()}")
try {
val store = KsqliteVectorIndexStore(dimension = 3, connection = conn)
store.add(id = "p", embedding = vec(1.5f, -2.25f, 3.875f), payload = null)
store.close()
val stmt = conn.prepare(
"SELECT id, embedding, payload FROM vector_index WHERE id = ?"
)
stmt.bindText(1, "p")
stmt.executeQuery().use { rs ->
assertTrue(rs.next())
assertEquals("p", rs.getText(0))
val blob = rs.getBlob(1)!!
val floats = blob.toFloatArray()
assertEquals(3, floats.size)
// Float сравниваем по toRawBits из-за float-округления.
assertEquals(1.5f.toRawBits(), floats[0].toRawBits())
assertEquals((-2.25f).toRawBits(), floats[1].toRawBits())
assertEquals(3.875f.toRawBits(), floats[2].toRawBits())
assertEquals(null, rs.getText(2))
}
stmt.close()
} finally {
conn.close()
}
}
}
@@ -0,0 +1,92 @@
package pw.binom.agentik.vectorindex.ksqlite
import kotlinx.coroutines.test.runTest
import pw.binom.db.ksqlite.SQLiteConnection
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertTrue
/**
* Тесты Schema.migrate(): idempotency, current_user_version, table existence.
*
* Принцип: каждый тест использует свежий in-memory connection, чтобы не
* зависеть от порядка выполнения.
*/
class SchemaMigrationTest {
private fun freshConn(name: String = "mig-${kotlin.random.Random.nextLong()}"): SQLiteConnection =
SQLiteConnection.memory(name)
@Test
fun `fresh DB gets the table and CURRENT_VERSION`() = runTest {
val conn = freshConn()
try {
Schema.migrate(conn)
assertTrue(tableExists(conn, Schema.TABLE_VECTOR_INDEX), "vector_index table should exist")
assertEquals(Schema.CURRENT_VERSION, readUserVersion(conn))
} finally {
conn.close()
}
}
@Test
fun `migrate is idempotent on already-migrated DB`() = runTest {
val conn = freshConn()
try {
Schema.migrate(conn)
// Добавим запись, чтобы убедиться что вторая migrate ничего не сломала.
val store = KsqliteVectorIndexStore(dimension = 2, connection = conn)
store.add(id = "x", embedding = floatArrayOf(1f, 0f), payload = null)
store.close()
// Повторный migrate должен быть no-op.
Schema.migrate(conn)
assertEquals(Schema.CURRENT_VERSION, readUserVersion(conn))
// Запись должна быть на месте.
val checkStmt = conn.prepare("SELECT COUNT(*) FROM ${Schema.TABLE_VECTOR_INDEX}")
checkStmt.executeQuery().use { rs ->
assertTrue(rs.next())
assertEquals(1L, rs.getLong(0))
}
checkStmt.close()
} finally {
conn.close()
}
}
@Test
fun `migrate leaves user_version alone if already at CURRENT_VERSION`() = runTest {
val conn = freshConn()
try {
// Имитируем БД, которая уже мигрирована (выставляем user_version вручную).
conn.exec("PRAGMA user_version = ${Schema.CURRENT_VERSION}")
assertEquals(Schema.CURRENT_VERSION, readUserVersion(conn))
// migrate должен быть no-op — таблица НЕ создаётся.
Schema.migrate(conn)
assertEquals(false, tableExists(conn, Schema.TABLE_VECTOR_INDEX),
"migrate on already-current DB must not create tables")
} finally {
conn.close()
}
}
private fun tableExists(conn: SQLiteConnection, name: String): Boolean {
val stmt = conn.prepare(
"SELECT 1 FROM sqlite_master WHERE type='table' AND name=?"
)
stmt.bindText(1, name)
stmt.executeQuery().use { rs -> return rs.next() }
}
private fun readUserVersion(conn: SQLiteConnection): Int {
conn.prepare("PRAGMA user_version").use { stmt ->
stmt.executeQuery().use { rs ->
if (rs.next()) return rs.getLong(0)?.toInt() ?: 0
}
}
return 0
}
}