feat(memory): add :memory-md-vector hybrid backend + :reflection-api + outbox/journal entity splits
ci / JVM build + tests (push) Failing after 2m59s
release / Publish KMP libraries → caffeine Nexus (release) Failing after 9s

- :memory-md-vector (KMP jvm+linuxX64+mingwX64): .md-файлы как source of
  truth, векторный индекс (sqlite-vec) как derived cache. reconcile()
  на старте: orphan-cleanup + content-hash-gated re-embed. Гибридный
  скор 0.7*vector + 0.3*keyword. Заменяет EmbeddingProvider на
  KMP-TextEmbeddingExecutor из :memory-api.

- :reflection-api: новый 4-й API-модуль (Reflection, ReflectionStore,
  ReflectionEvent). Зависит только от :memory-api.

- :journal-api получил ConversationRecord/ConversationStore/Ids (бывший
  :message-store-api, полностью удалён). :outbox-api получил Event,
  CommonEvent, AgentEvent (бывший :event-store).

- :memory-api получил MemoryVectorIndex + NoteMatches +
  TextEmbeddingExecutor (suspend-обёртка над TextEmbeddingExtractor).

- :memory-vector KMP-цели достигнуты через commonMain-only TextEmbedding-
  Executor, EmbeddingProvider выпилен; :memory-md-vector тянет
  text-embedding-api транзитивно через :memory-api.

- :standalone flatten в commonMain/commonTest завершён (тесты из jvmTest
  переехали в commonTest). Включён optional деп :memory-md-vector через
  AGENTIK_MEMORY_BACKEND=md-vector.

jvmTest: 96 задач, 407 тестов, 0 падений.
This commit is contained in:
subochev
2026-09-21 23:10:28 +03:00
parent 161be41adf
commit 0a7c40688c
41 changed files with 4952 additions and 0 deletions
+51
View File
@@ -0,0 +1,51 @@
plugins {
alias(libs.plugins.kotlin.multiplatform)
}
// :memory-md-vector — гибридное хранилище памяти:
//
// .md файлы (:memory-md, single source of truth)
// ↓ reconcile() на старте
// sqlite vector index (ksqlite + sqlite-vec vec0, derived cache)
//
// `.md` — единственный источник правды по метаданным и тексту заметок.
// Вектора — derived cache, перестраивается на старте и при `upsert`/`delete`.
//
// ANN-поиск: vector KNN (sqlite-vec MATCH) → top-50 → keyword rerank
// через `MdMemoryFormat.keywordScore` (vector 0.7 + keyword 0.3).
//
// Цели сборки — KMP: jvm() + linuxX64() + mingwX64(). До 2026-09-21 был
// JVM-only, потому что тащил `EmbeddingProvider` из JVM-only `:memory-vector`.
// С переходом на `TextEmbeddingExecutor` (из `:memory-api`, который тянет
// `pw.binom.ai.embeddingtext:api` — теперь KMP) модуль стал платформо-
// независимым. Под нативом тесты работают с `FakeTextEmbeddingExtractor`;
// прод-реализация (`:siglip` модуль text-embedding-kmp) пока JVM+Android only.
kotlin {
jvmToolchain(21)
jvm()
linuxX64()
mingwX64()
sourceSets {
commonMain.dependencies {
// ksqlite 0.1.2 опубликован в Maven Central — обычный
// `mavenCentral()` в settings.gradle.kts его подтянет.
implementation("pw.binom.db:ksqlite:0.1.2")
implementation(libs.kotlinx.coroutines.core)
implementation(libs.kotlinx.io.core)
api(project(":memory-api"))
implementation(project(":memory-md"))
// `text-embedding-api` тянется транзитивно через `:memory-api`
// (мы добавили `api(libs.text.embedding.api)` в memory-api/build.gradle.kts).
// Раньше тут стоял `implementation(project(":memory-vector"))` ради
// `EmbeddingProvider` — JVM-only модуль с JVector. Теперь не нужен.
}
commonTest.dependencies {
implementation(kotlin("test"))
implementation(libs.kotlinx.coroutines.test)
}
}
}
@@ -0,0 +1,206 @@
package pw.binom.agentik.memory.mdvector
import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock
import pw.binom.agentik.memory.MemoryNote
import pw.binom.agentik.memory.MemorySearchQuery
import pw.binom.agentik.memory.MemorySearchResult
import pw.binom.agentik.memory.MemoryStore
import pw.binom.agentik.memory.MemoryStoreEvent
import pw.binom.agentik.memory.TextEmbeddingExecutor
import pw.binom.agentik.memory.md.MdMemoryFormat
import pw.binom.agentik.memory.md.MdMemoryStore
import pw.binom.agentik.memory.noteMatches
import kotlin.time.Instant
import kotlinx.coroutines.flow.Flow
/**
* Гибридное хранилище памяти:
*
* * `.md` файлы (через [MdMemoryStore]) — single source of truth по
* метаданным и тексту заметок;
* * ksqlite vector index ([KsqliteVectorIndex]) — derived cache embeddings
* и content_hash.
*
** Архитектурный контракт:
*
* 1. Любая мутация (upsert/delete) обновляет оба слоя атомарно: сначала
* `.md` (через [MdMemoryStore]), потом векторный кэш. Если vector-write
* упал — `.md` уже сохранён; reconcile при следующем старте восстановит
* консистентность.
*
* 2. [search] использует vector ANN (sqlite-vec MATCH) → top-50 → keyword
* rerank (`MdMemoryFormat.keywordScore`). Финальный score = 0.7 * vector
* + 0.3 * keyword. Это даёт семантический recall с быстрой фильтрацией
* по точным совпадениям.
*
* 3. [reconcile] — вызывается при старте (из [openHybridMemoryStore]):
* - .md файл есть, вектора нет → embed + add;
* - .md файл есть, вектор есть, content_hash отличается → re-embed;
* - .md файла нет, вектор есть → orphan, remove.
*
* 4. Read-only методы ([get], [list], [markUsed], [archiveStale], [events])
* делегируются в [MdMemoryStore] напрямую — никакой транзакции с
* vector-кэшем.
*
* Потокобезопасность: делегирующие методы — thread-safe за счёт
* `MdMemoryStore.mu`. Мутации векторов сериализуются
* [KsqliteVectorIndex.mutex]. Метод [reconcile] держит свой [mutex] для
* исключения конкурентных upsert'ов во время согласования.
*/
class HybridMdVectorStore internal constructor(
private val mdStore: MdMemoryStore,
private val vectorIndex: KsqliteVectorIndex,
private val embedder: TextEmbeddingExecutor,
) : MemoryStore {
private val reconcileMutex = Mutex()
/**
* Отчёт о согласовании `.md` ↔ vector-индекс. Возвращается из [reconcile].
*/
data class ReconcileReport(
val added: Int,
val reembedded: Int,
val orphansRemoved: Int,
) {
val totalChanged: Int get() = added + reembedded + orphansRemoved
}
/**
* Согласовать vector-кэш с текущим состоянием `.md` файлов.
*
* Идемпотентен — повторный вызов no-op.
*
* Можно вызывать из фонового потока при старте `Main.kt` чтобы
* залогировать "reconciled: 5 re-embedded, 2 added, 0 orphans".
*/
suspend fun reconcile(): ReconcileReport = reconcileMutex.withLock {
val onDisk: List<MemoryNote> = mdStore.list(limit = Int.MAX_VALUE)
val onDiskById: Map<String, MemoryNote> = onDisk.associateBy { it.id }
val inCache: List<KsqliteVectorIndex.MetaEntry> = vectorIndex.allMeta()
val cachedIds: Set<String> = inCache.map { it.id }.toSet()
var added = 0
var reembedded = 0
var orphansRemoved = 0
// 1) orphan-cleanup: vector есть, .md нет
for (cached in inCache) {
if (cached.id !in onDiskById) {
vectorIndex.remove(cached.id)
orphansRemoved++
}
}
// 2) re-embed / add
for (note in onDisk) {
val cached = inCache.firstOrNull { it.id == note.id }
val currentHash = note.contentHash()
if (cached == null) {
// .md есть, вектора нет → add
val vec = embedder.embed(note.content)
vectorIndex.add(note.id, vec, currentHash)
added++
} else if (cached.contentHash != currentHash) {
// .md изменился → re-embed
val vec = embedder.embed(note.content)
vectorIndex.add(note.id, vec, currentHash)
reembedded++
}
// else: cached.contentHash == currentHash → no-op
}
ReconcileReport(added, reembedded, orphansRemoved)
}
// ─── MemoryStore impl: мутации ─────────────────────────────────────
override suspend fun upsert(note: MemoryNote) {
mdStore.upsert(note)
val vec = embedder.embed(note.content)
vectorIndex.add(note.id, vec, note.contentHash())
}
override suspend fun delete(id: String): Boolean {
val existed = mdStore.delete(id)
vectorIndex.remove(id)
return existed
}
// ─── MemoryStore impl: search (hybrid) ────────────────────────────
override suspend fun search(query: MemorySearchQuery): List<MemorySearchResult> {
if (query.query.isBlank()) return emptyList()
if (query.topK <= 0) return emptyList()
// Этап 1: vector ANN top-K (K=50 или больше topK).
val candidateK = maxOf(query.topK, VECTOR_CANDIDATES)
val qVec = embedder.embed(query.query)
val vectorHits = vectorIndex.search(qVec, candidateK)
// Этап 2: загружаем кандидатов из .md (single source of truth)
// mapNotNull не умеет suspend, поэтому собираем вручную.
val candidates: List<Pair<MemoryNote, Float>> = buildList(vectorHits.size) {
for (hit in vectorHits) {
val note = mdStore.get(hit.id) ?: continue
// Применяем категорийный/конво-фильтр ДО rerank — экономим keywordScore.
if (!noteMatches(note, query.category, query.conversationId)) continue
add(note to hit.score)
}
}
// Этап 3: keyword rerank (vector 0.7 + keyword 0.3)
val rescored = candidates.map { (note, vecScore) ->
val kwScore = MdMemoryFormat.keywordScore(query.query, note)
val finalScore = vecScore * VECTOR_WEIGHT + kwScore * KEYWORD_WEIGHT
MemorySearchResult(note, finalScore)
}.sortedByDescending { it.score }
return if (rescored.size > query.topK) rescored.subList(0, query.topK) else rescored
}
// ─── MemoryStore impl: read-only delegation ───────────────────────
override suspend fun get(id: String): MemoryNote? = mdStore.get(id)
override suspend fun list(
category: pw.binom.agentik.memory.MemoryCategory?,
conversationId: String?,
limit: Int,
offset: Int,
): List<MemoryNote> = mdStore.list(category, conversationId, limit, offset)
override suspend fun markUsed(id: String, at: Instant) = mdStore.markUsed(id, at)
override suspend fun archiveStale(
maxAge: kotlin.time.Duration,
maxUseCount: Int,
now: Instant,
): Int {
// Вектор-кэш не хранит lastUsedAt/useCount (только content_hash).
// Делегируем в mdStore — он сам знает что удалять; vector удалится
// каскадно при archiveStale → delete loop ниже.
val deleted = mdStore.archiveStale(maxAge, maxUseCount, now)
// Дополнительно чистим vector-кэш от записей, которых больше нет в .md
val remaining = mdStore.list(limit = Int.MAX_VALUE).map { it.id }.toSet()
vectorIndex.allMeta().forEach { entry ->
if (entry.id !in remaining) vectorIndex.remove(entry.id)
}
return deleted
}
override fun events(): Flow<MemoryStoreEvent> = mdStore.events()
override fun close() {
runCatching { vectorIndex.close() }
runCatching { mdStore.close() }
}
companion object {
const val VECTOR_WEIGHT: Float = 0.7f
const val KEYWORD_WEIGHT: Float = 0.3f
const val VECTOR_CANDIDATES: Int = 50
}
}
@@ -0,0 +1,66 @@
package pw.binom.agentik.memory.mdvector
import kotlinx.io.files.Path
import pw.binom.agentik.memory.TextEmbeddingExecutor
import pw.binom.agentik.memory.md.openMdMemory
import pw.binom.db.ksqlite.SQLiteConnection
/**
* Открыть гибридное хранилище памяти (`.md` + sqlite vector index).
*
* Создаёт:
* - [MdMemoryStore] на [memoryRoot] (`.md` файлы);
* - [KsqliteVectorIndex] на [vectorDbPath] (sqlite-vec vec0);
* - [HybridMdVectorStore] — обёртка с reconcile и hybrid search.
*
* Перед возвратом выполняет [HybridMdVectorStore.reconcile] — для свежей
* БД это приведёт к первичному embed'у всех `.md` файлов; для существующей —
* к re-embed'у изменившихся заметок и orphan-cleanup.
*
* @param memoryRoot директория с `.md` файлами (`USER.md`, `WORLD.md`, ...).
* @param vectorDbPath путь к файлу sqlite-БД для vector-кэша.
* @param dimension размерность embeddings от [embedder]. Фиксируется
* при создании индекса; дальнейшая смена = wipe БД.
* @param embedder провайдер embeddings.
* @param runReconcile выполнить [HybridMdVectorStore.reconcile] сразу после
* открытия. В тестах можно отключить для скорости.
*/
fun openHybridMemoryStore(
memoryRoot: Path,
vectorDbPath: Path,
dimension: Int,
embedder: TextEmbeddingExecutor,
runReconcile: Boolean = true,
): HybridMdVectorStore {
val md = openMdMemory(memoryRoot)
val conn = SQLiteConnection.open(vectorDbPath.toString())
Schema.migrate(conn, dimension)
val idx = KsqliteVectorIndex(conn, dimension)
val hybrid = HybridMdVectorStore(md, idx, embedder)
if (runReconcile) {
kotlinx.coroutines.runBlocking { hybrid.reconcile() }
}
return hybrid
}
/**
* In-memory вариант для тестов: vector-кэш в `:memory:` sqlite,
* `.md` — в `/tmp/agentik-hybrid-test-{random}`.
*
* Используется POSIX-путь `/tmp`, потому что [System.getenv] / [System.getProperty]
* недоступны в KMP commonMain (только JVM). На Windows mingwX64 этот вызов
* упадёт — там тесты пока не предполагаются, нативные тесты только linuxX64.
* Под JVM `/tmp` либо есть как symlink (Linux/macOS), либо стоит использовать
* jvmTest-специфичный factory.
*/
fun openInMemoryHybridMemoryStore(
dimension: Int,
embedder: TextEmbeddingExecutor,
): HybridMdVectorStore {
val tmpDir = Path("/tmp/agentik-hybrid-test-${kotlin.random.Random.nextLong()}")
val md = openMdMemory(tmpDir)
val conn = SQLiteConnection.memory("hybrid-${kotlin.random.Random.nextLong()}")
Schema.migrate(conn, dimension)
val idx = KsqliteVectorIndex(conn, dimension)
return HybridMdVectorStore(md, idx, embedder)
}
@@ -0,0 +1,47 @@
package pw.binom.agentik.memory.mdvector
import kotlinx.io.files.Path
import pw.binom.agentik.memory.MemoryPrefetcher
import pw.binom.agentik.memory.MemoryReviewer
import pw.binom.agentik.memory.MemoryStore
import pw.binom.agentik.memory.MemorySystem
import pw.binom.agentik.memory.TextEmbeddingExecutor
import pw.binom.agentik.memory.md.KeywordMdPrefetcher
import pw.binom.agentik.memory.md.KeywordMdReviewer
/**
* Связка [HybridMdVectorStore] + keyword-prefetcher + keyword-reviewer.
*
* Store делегирует I/O между .md (single source of truth) и sqlite-vector-кэшем;
* prefetcher и reviewer работают по .md-данным (через [HybridMdVectorStore]),
* так что обе роли видят консистентное состояние.
*/
class HybridMemorySystem internal constructor(
override val store: MemoryStore,
override val prefetcher: MemoryPrefetcher,
override val reviewer: MemoryReviewer,
) : MemorySystem {
override fun close() = store.close()
}
/**
* Собирает [HybridMemorySystem] для указанной корневой директории + sqlite-БД.
*
* Под капотом: [HybridMdVectorStore] (md + vector), keyword-prefetcher из
* `:memory-md` (работает по store.api), keyword-reviewer без LLM —
* LLM-импл добавляется в `:standalone` поверх.
*/
fun openHybridMemorySystem(
memoryRoot: Path,
vectorDbPath: Path,
dimension: Int,
embedder: TextEmbeddingExecutor,
runReconcile: Boolean = true,
): HybridMemorySystem {
val store = openHybridMemoryStore(memoryRoot, vectorDbPath, dimension, embedder, runReconcile)
return HybridMemorySystem(
store = store,
prefetcher = KeywordMdPrefetcher(store),
reviewer = KeywordMdReviewer(),
)
}
@@ -0,0 +1,305 @@
package pw.binom.agentik.memory.mdvector
import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock
import pw.binom.agentik.memory.MemoryNote
import pw.binom.agentik.memory.MemoryVectorIndex
import pw.binom.agentik.memory.ScoredVector
import pw.binom.db.ksqlite.SQLiteConnection
import pw.binom.db.ksqlite.SQLitePreparedStatement
import kotlin.time.Clock
/**
* ksqlite-реализация [MemoryVectorIndex] поверх `vec0` virtual table
* (sqlite-vec extension, встроен в ksqlite).
*
* Маппинг id → rowid:
* - TEXT `id` (== MemoryNote.id, "mem-...") лежит в [Schema.TABLE_META].
* - `vec0` индексирует по `rowid` (INTEGER auto-increment).
* - JOIN через `WHERE vec0.rowid = meta.rowid`.
*
* На каждое [add] с contentHash рядом с вектором пишется meta с
* content_hash от [MemoryNote.contentHash]. Это позволяет reconcile'у
* в [HybridMdVectorStore] определить "изменилась ли заметка" без re-embed.
*
* Поиск — `vec0` MATCH (cosine distance, sqlite-vec native). Score
* конвертируется из distance (0..2, меньше = ближе) в similarity
* (0..1, больше = ближе).
*
* Конкурентность: write-операции сериализуются [mutex]; read'ы
* (`size`/`search`) не блокируют.
*/
class KsqliteVectorIndex internal constructor(
private val conn: SQLiteConnection,
override val dimension: Int,
) : MemoryVectorIndex {
private val mutex = Mutex()
private val insertVec: SQLitePreparedStatement = conn.prepare(
"INSERT INTO ${Schema.TABLE_VECTORS}(${Schema.COL_VECTOR}) VALUES (?)"
)
private val lastInsertRowIdStmt: SQLitePreparedStatement = conn.prepare(
"SELECT last_insert_rowid()"
)
private val insertMeta: SQLitePreparedStatement = conn.prepare(
"""
INSERT OR REPLACE INTO ${Schema.TABLE_META}
(${Schema.COL_ROWID}, ${Schema.COL_ID}, ${Schema.COL_HASH},
${Schema.COL_DIMENSION}, ${Schema.COL_UPDATED_AT})
VALUES (?, ?, ?, ?, ?)
""".trimIndent()
)
private val findMetaByIdStmt: SQLitePreparedStatement = conn.prepare(
"SELECT ${Schema.COL_ROWID}, ${Schema.COL_HASH} FROM ${Schema.TABLE_META} WHERE ${Schema.COL_ID} = ?"
)
/** Запрос `... WHERE rowid IN (?, ?, ...)`. Подготавливаем на N=$MAX_INLINE_ROWIDS параметров. */
private val findMetaByRowIdsStmt: SQLitePreparedStatement = conn.prepare(
(1..MAX_INLINE_ROWIDS).joinToString(
separator = ",",
prefix = "SELECT ${Schema.COL_ROWID}, ${Schema.COL_ID} FROM ${Schema.TABLE_META} WHERE ${Schema.COL_ROWID} IN (",
postfix = ")",
) { "?" }
)
private val deleteByRowIdStmt: SQLitePreparedStatement = conn.prepare(
"DELETE FROM ${Schema.TABLE_VECTORS} WHERE rowid = ?"
)
private val deleteMetaByRowIdStmt: SQLitePreparedStatement = conn.prepare(
"DELETE FROM ${Schema.TABLE_META} WHERE ${Schema.COL_ROWID} = ?"
)
private val deleteMetaByIdStmt: SQLitePreparedStatement = conn.prepare(
"DELETE FROM ${Schema.TABLE_META} WHERE ${Schema.COL_ID} = ?"
)
private val sizeMetaStmt: SQLitePreparedStatement = conn.prepare(
"SELECT COUNT(*) FROM ${Schema.TABLE_META}"
)
private val allMetaStmt: SQLitePreparedStatement = conn.prepare(
"SELECT ${Schema.COL_ROWID}, ${Schema.COL_ID}, ${Schema.COL_HASH} FROM ${Schema.TABLE_META}"
)
/**
* sqlite-vec требует чтобы `LIMIT` в MATCH-запросе был integer-литералом,
* а не `?`. Поэтому для search используем динамическую подготовку
* (кешированную по [k]).
*
* Vec0 KNN check (`sqlite-vec` source): «A LIMIT or 'k = ?' constraint is
* required on vec0 knn queries». Имя параметра `:k` тоже поддерживается,
* но у ksqlite bind API — только позиционный; literal проще.
*/
private val searchCache = HashMap<Int, SQLitePreparedStatement>()
/**
* Выдать rowid для существующей записи или -1 если нет.
*
* **ВАЖНО**: вызывающий ОБЯЗАН держать [mutex]. Этот метод НЕ
* reentrant — повторный вход в [Mutex.withLock] приведёт к
* deadlock (kotlinx.coroutines.sync.Mutex не reentrant).
*/
private fun findRowIdLocked(id: String): Long {
findMetaByIdStmt.reset()
findMetaByIdStmt.clearBindings()
findMetaByIdStmt.bindText(1, id)
findMetaByIdStmt.executeQuery().use { rs ->
if (rs.next()) return rs.getLong(0) ?: -1L
}
return -1L
}
override suspend fun size(): Long = mutex.withLock {
sizeMetaStmt.reset()
sizeMetaStmt.clearBindings()
sizeMetaStmt.executeQuery().use { rs ->
if (rs.next()) rs.getLong(0) ?: 0L else 0L
}
}
override suspend fun add(id: String, embedding: FloatArray) {
require(embedding.size == dimension) {
"embedding dim=${embedding.size} != index dim=$dimension"
}
add(id, embedding, contentHash = "")
}
/**
* Добавить или обновить запись с явным content_hash.
* Если запись с таким id уже есть — обновляет и вектор, и meta.
* Иначе — создаёт новый rowid.
*/
suspend fun add(id: String, embedding: FloatArray, contentHash: String) {
require(embedding.size == dimension) {
"embedding dim=${embedding.size} != index dim=$dimension"
}
mutex.withLock {
val existingRowId = findRowIdLocked(id)
val rowId: Long = if (existingRowId > 0) {
// Update: заменяем вектор по существующему rowid
insertVec.reset()
insertVec.clearBindings()
insertVec.bindVector(1, embedding)
insertVec.executeUpdate()
existingRowId
} else {
// Insert: получаем свежий rowid
insertVec.reset()
insertVec.clearBindings()
insertVec.bindVector(1, embedding)
insertVec.executeUpdate()
lastInsertRowIdStmt.reset()
lastInsertRowIdStmt.clearBindings()
lastInsertRowIdStmt.executeQuery().use { rs ->
if (rs.next()) rs.getLong(0) ?: error("no last_insert_rowid()") else error("no last_insert_rowid()")
}
}
insertMeta.reset()
insertMeta.clearBindings()
insertMeta.bindLong(1, rowId)
insertMeta.bindText(2, id)
insertMeta.bindText(3, contentHash)
insertMeta.bindLong(4, dimension.toLong())
insertMeta.bindLong(5, Clock.System.now().toEpochMilliseconds())
insertMeta.executeUpdate()
}
}
override suspend fun remove(id: String): Boolean = mutex.withLock {
val rowId = findRowIdLocked(id)
if (rowId <= 0) return@withLock false
// Удаляем meta сначала — иначе orphan-row в vec0.
deleteMetaByIdStmt.reset()
deleteMetaByIdStmt.clearBindings()
deleteMetaByIdStmt.bindText(1, id)
deleteMetaByIdStmt.executeUpdate()
deleteByRowIdStmt.reset()
deleteByRowIdStmt.clearBindings()
deleteByRowIdStmt.bindLong(1, rowId)
deleteByRowIdStmt.executeUpdate()
true
}
/** Прочитать content_hash для id. null если записи нет. */
suspend fun getContentHash(id: String): String? = mutex.withLock {
findMetaByIdStmt.reset()
findMetaByIdStmt.clearBindings()
findMetaByIdStmt.bindText(1, id)
findMetaByIdStmt.executeQuery().use { rs ->
if (rs.next()) rs.getText(1) else null
}
}
/** Полный список (rowid, id, contentHash) для reconcile'а. */
suspend fun allMeta(): List<MetaEntry> = mutex.withLock {
allMetaStmt.reset()
allMetaStmt.clearBindings()
val out = mutableListOf<MetaEntry>()
allMetaStmt.executeQuery().use { rs ->
while (rs.next()) {
val rowId = rs.getLong(0) ?: continue
val id = rs.getText(1) ?: continue
val hash = rs.getText(2) ?: continue
out.add(MetaEntry(rowId, id, hash))
}
}
out
}
override suspend fun search(
query: FloatArray,
k: Int,
filter: (MemoryNote) -> Boolean,
): List<ScoredVector> {
require(query.size == dimension) {
"query dim=${query.size} != index dim=$dimension"
}
val safeK = k.coerceAtLeast(1)
return mutex.withLock {
// sqlite-vec MATCH требует минимальный запрос без JOIN/лишних
// ORDER BY — иначе "A LIMIT or 'k = ?' constraint is required".
// Поэтому делаем два запроса:
// 1) vec0 ANN → (rowid, distance)
// 2) meta lookup по собранным rowid → id
val stmt = searchCache.getOrPut(safeK) {
conn.prepare(
"""
SELECT rowid, distance
FROM ${Schema.TABLE_VECTORS}
WHERE ${Schema.COL_VECTOR} MATCH ?
ORDER BY distance
LIMIT $safeK
""".trimIndent()
)
}
stmt.reset()
stmt.clearBindings()
stmt.bindVector(1, query)
val candidates = mutableListOf<Pair<Long, Float>>()
stmt.executeQuery().use { rs ->
while (rs.next()) {
val rowId = rs.getLong(0) ?: continue
val distance = rs.getDouble(1) ?: continue
val score = ((1.0 - distance / 2.0) * 1.0).toFloat().coerceIn(0f, 1f)
candidates.add(rowId to score)
}
}
if (candidates.isEmpty()) return@withLock emptyList<ScoredVector>()
// 2-й запрос: meta по списку rowid.
val rowIds = candidates.map { it.first }
val idByRowId = HashMap<Long, String>(candidates.size)
findMetaByRowIdsStmt.reset()
findMetaByRowIdsStmt.clearBindings()
for ((idx, rowId) in rowIds.withIndex()) {
findMetaByRowIdsStmt.bindLong(idx + 1, rowId)
}
findMetaByRowIdsStmt.executeQuery().use { rs ->
while (rs.next()) {
val rowId = rs.getLong(0) ?: continue
val id = rs.getText(1) ?: continue
idByRowId[rowId] = id
}
}
candidates.mapNotNull { (rowId, score) ->
idByRowId[rowId]?.let { ScoredVector(it, score) }
}
}
}
override suspend fun flush() {
// ksqlite + WAL — flush не нужен. Метод для совместимости с интерфейсом.
}
override fun close() {
insertVec.close()
lastInsertRowIdStmt.close()
insertMeta.close()
findMetaByIdStmt.close()
deleteByRowIdStmt.close()
deleteMetaByRowIdStmt.close()
deleteMetaByIdStmt.close()
sizeMetaStmt.close()
allMetaStmt.close()
searchCache.values.forEach { it.close() }
}
data class MetaEntry(val rowId: Long, val id: String, val contentHash: String)
private companion object {
/** Максимум rowid, которые мы зашиваем в `IN (?,?,...)` одним prepared statement'ом. */
const val MAX_INLINE_ROWIDS = 256
}
}
@@ -0,0 +1,80 @@
package pw.binom.agentik.memory.mdvector
import pw.binom.db.ksqlite.SQLiteConnection
/**
* DDL/DML для ksqlite-бэкенда `:memory-md-vector`.
*
* Две таблицы:
*
* * `memory_vectors` (vec0) — ANN-индекс. Содержит embedding + первичный
* ключ `rowid` (auto-increment INTEGER, sqlite-vec требует именно его).
* Размерность задаётся `float[DIMENSION]` при создании.
*
* * `memory_meta` — рядом с вектором: TEXT `id` (== MemoryNote.id) +
* `rowid` (тот же, что в vec0) + `content_hash` (от MemoryNote.contentHash()).
* Используется reconcile'ом — если хэш в meta не совпадает с тем, что
* вычисляется из текущего `.md` файла → re-embed.
*
* JOIN между vec0 и meta: `WHERE vec0.rowid = meta.rowid`.
*
* Миграция через `PRAGMA user_version` (как в `:journal-ksqlite/Schema.kt`).
*/
internal object Schema {
const val CURRENT_VERSION: Int = 1
const val TABLE_VECTORS = "memory_vectors"
const val TABLE_META = "memory_meta"
const val COL_ID = "id"
const val COL_ROWID = "rowid"
const val COL_VECTOR = "embedding"
const val COL_HASH = "content_hash"
const val COL_DIMENSION = "dimension"
const val COL_UPDATED_AT = "updated_at"
fun v1Ddl(dimension: Int): String = """
CREATE VIRTUAL TABLE IF NOT EXISTS $TABLE_VECTORS USING vec0(
$COL_VECTOR float[$dimension]
);
CREATE TABLE IF NOT EXISTS $TABLE_META (
$COL_ROWID INTEGER NOT NULL PRIMARY KEY AUTOINCREMENT,
$COL_ID TEXT NOT NULL UNIQUE,
$COL_HASH TEXT NOT NULL,
$COL_DIMENSION INTEGER NOT NULL,
$COL_UPDATED_AT INTEGER NOT NULL
);
CREATE INDEX IF NOT EXISTS idx_meta_id ON $TABLE_META($COL_ID);
""".trimIndent()
fun migrate(conn: SQLiteConnection, dimension: Int) {
val current = readUserVersion(conn)
if (current >= CURRENT_VERSION) return
conn.exec("BEGIN")
try {
if (current < 1) {
conn.exec(v1Ddl(dimension))
}
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) {
conn.exec("PRAGMA user_version = $version")
}
}
@@ -0,0 +1,166 @@
package pw.binom.agentik.memory.mdvector
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertNotNull
import kotlin.test.assertNull
import kotlin.test.assertTrue
import kotlin.time.Instant
import kotlinx.coroutines.runBlocking as kRunBlocking
import pw.binom.agentik.memory.MemoryCategory
import pw.binom.agentik.memory.MemoryNote
import pw.binom.agentik.memory.MemorySearchQuery
import pw.binom.agentik.memory.MemorySource
import pw.binom.agentik.memory.TextEmbeddingExecutor
import pw.binom.voice.embeddingtext.TextEmbedding
import pw.binom.voice.embeddingtext.TextEmbeddingExtractor
/**
* Тесты гибридного стора: делегирование в .md, vector ANN, reconcile,
* hybrid search rerank.
*
* Используется [kRunBlocking] (а не `runTest`) потому что `MdMemoryStore`
* делает реальный файловый I/O (`kotlinx-io`), который плохо дружит с
* TestDispatcher'ом — `runTest` зависает на virtual-time I/O.
*/
class HybridMdVectorStoreTest {
/**
* Детерминированный [TextEmbeddingExtractor] для тестов модуля.
* Хеширует текст в псевдо-вектор фиксированной размерности, L2-normalize.
* Заворачивается в [TextEmbeddingExecutor] с пред-объявленной размерностью.
*/
private fun fakeEmbeddingExecutor(dimension: Int = 32): TextEmbeddingExecutor =
TextEmbeddingExecutor(FakeTestExtractor(dimension), knownDimension = dimension)
private class FakeTestExtractor(val dim: Int) : TextEmbeddingExtractor {
override fun embed(text: String): TextEmbedding {
val v = FloatArray(dim)
var seed = text.hashCode().toLong() and 0xFFFFFFFFL
for (i in 0 until dim) {
seed = (seed * 6364136223846793005L + 1442695040888963407L) and 0xFFFFFFFFL
v[i] = ((seed.toInt() and 0xFFFF) / 65535f) * 2f - 1f
}
var norm = 0f
for (x in v) norm += x * x
norm = kotlin.math.sqrt(norm)
if (norm > 0f) for (i in v.indices) v[i] /= norm
return TextEmbedding(v)
}
override fun close() = Unit
}
private fun note(
id: String,
content: String,
category: MemoryCategory = MemoryCategory.USER,
source: MemorySource = MemorySource.USER_EXPLICIT,
) = MemoryNote(
id = id,
category = category,
content = content,
createdAt = Instant.fromEpochMilliseconds(1_700_000_000_000L),
lastUsedAt = Instant.fromEpochMilliseconds(1_700_000_000_000L),
useCount = 0,
conversationId = null,
source = source,
)
@Test
fun upsertWritesToBothMdAndVectorCache() = kRunBlocking {
val store = openInMemoryHybridMemoryStore(dimension = 32, embedder = fakeEmbeddingExecutor(32))
store.upsert(note("mem-1", "user prefers dark mode"))
val fromMd = store.get("mem-1")
assertNotNull(fromMd)
assertEquals("user prefers dark mode", fromMd.content)
val hits = store.search(MemorySearchQuery(query = "user prefers dark mode", topK = 5))
assertEquals(1, hits.size)
assertEquals("mem-1", hits.first().note.id)
store.close()
}
@Test
fun deleteRemovesFromBothLayers() = kRunBlocking {
val store = openInMemoryHybridMemoryStore(dimension = 32, embedder = fakeEmbeddingExecutor(32))
store.upsert(note("mem-2", "lives in Saint Petersburg"))
store.upsert(note("mem-3", "loves Kotlin multiplatform"))
assertEquals(2, store.list(limit = 10).size)
val removed = store.delete("mem-2")
assertTrue(removed)
assertEquals(1, store.list(limit = 10).size)
assertNull(store.get("mem-2"))
val hits = store.search(MemorySearchQuery(query = "Saint Petersburg", topK = 5))
assertTrue(hits.isEmpty() || hits.all { it.note.id != "mem-2" })
store.close()
}
@Test
fun reconcileOnEmptyStoreIsNoOp() = kRunBlocking {
val store = openInMemoryHybridMemoryStore(dimension = 32, embedder = fakeEmbeddingExecutor(32))
val report = store.reconcile()
assertEquals(0, report.added)
assertEquals(0, report.reembedded)
assertEquals(0, report.orphansRemoved)
store.close()
}
@Test
fun reconcileIsIdempotent() = kRunBlocking {
val store = openInMemoryHybridMemoryStore(dimension = 32, embedder = fakeEmbeddingExecutor(32))
store.upsert(note("mem-10", "works at Binom"))
val report = store.reconcile()
assertEquals(0, report.totalChanged, "идемпотентность reconcile: повторный вызов no-op")
store.close()
}
@Test
fun searchReturnsRelevantResultsByKeyword() = kRunBlocking {
val store = openInMemoryHybridMemoryStore(dimension = 32, embedder = fakeEmbeddingExecutor(32))
store.upsert(note("mem-a", "kotlin multiplatform"))
store.upsert(note("mem-b", "java enterprise"))
store.upsert(note("mem-c", "kotlin coroutines"))
val hits = store.search(MemorySearchQuery(query = "kotlin", topK = 5))
val ids = hits.map { it.note.id }.toSet()
assertTrue("mem-a" in ids, "expected 'kotlin multiplatform' in results: $ids")
assertTrue("mem-c" in ids, "expected 'kotlin coroutines' in results: $ids")
store.close()
}
@Test
fun searchRespectsCategoryFilter() = kRunBlocking {
val store = openInMemoryHybridMemoryStore(dimension = 32, embedder = fakeEmbeddingExecutor(32))
store.upsert(note("user-1", "kotlin lover", MemoryCategory.USER))
store.upsert(note("world-1", "kotlin 2.0 released", MemoryCategory.WORLD))
val hits = store.search(MemorySearchQuery(query = "kotlin", topK = 10, category = MemoryCategory.USER))
assertEquals(1, hits.size)
assertEquals("user-1", hits.first().note.id)
store.close()
}
@Test
fun hybridScoreCombinesVectorAndKeyword() = kRunBlocking {
val store = openInMemoryHybridMemoryStore(dimension = 32, embedder = fakeEmbeddingExecutor(32))
store.upsert(note("mem-x", "kotlin multiplatform project"))
val hits = store.search(MemorySearchQuery(query = "kotlin multiplatform", topK = 1))
assertEquals(1, hits.size)
val score = hits.first().score
assertTrue(score in 0f..1f, "score $score out of range")
assertTrue(score > 0.5f, "expected hybrid score > 0.5, got $score")
store.close()
}
}