1 Commits
7 .. 8

Author SHA1 Message Date
subochev 68543357c2 feat(memory): migrate EmbeddingProvider to KMP-compatible TextEmbeddingExecutor, add :memory-md-vector, and hybrid backend support
ci / JVM build + tests (push) Failing after 11s
release / Publish KMP libraries → caffeine Nexus (release) Failing after 10s
- Replaced `EmbeddingProvider` with cross-platform `TextEmbeddingExecutor` for native target compatibility.
- Introduced `:memory-md-vector` module combining vector-cache and `.md` file-based memory systems (`hybrid` backend).
- Updated `SiglipEmbeddingProvider` to use KMP `TextEmbeddingExtractor` and streamlined compatibility via `asExecutor`.
- Added hybrid memory backend to `standalone`, supporting `.md` reconciliation with vector-cache for semantic
2026-09-21 12:28:09 +03:00
19 changed files with 658 additions and 165 deletions
+7 -2
View File
@@ -97,8 +97,13 @@ kotlinx-io-core = { module = "org.jetbrains.kotlinx:kotlinx-io-core", version.re
jvector = { module = "io.github.jbellis:jvector", version.ref = "jvector" } jvector = { module = "io.github.jbellis:jvector", version.ref = "jvector" }
# --- text-embedding-kmp (pw.binom.ai.embeddingtext) — on-device SigLIP2 эмбеддинг через ONNX. --- # --- text-embedding-kmp (pw.binom.ai.embeddingtext) — on-device SigLIP2 эмбеддинг через ONNX. ---
# Артефакты публикуются под именами `-jvm` (KMP convention для JVM-таргета). # `api` — KMP с jvm + android + linuxX64/Arm64 + macos + ios + mingwX64
text-embedding-api = { module = "pw.binom.ai.embeddingtext:api-jvm", version.ref = "text-embedding-kmp" } # (с 2026-09-21, когда мы добавили нативные цели в text-embedding-kmp:api).
# Используется из :memory-md-vector и :memory-vector напрямую через
# `libs.text.embedding.api` (без суффикса `-jvm` — Gradle сам выберет
# нужный variant под target).
# `siglip` — JVM+Android only (onnx-runtime), подключается в jvmMain.
text-embedding-api = { module = "pw.binom.ai.embeddingtext:api", version.ref = "text-embedding-kmp" }
text-embedding-siglip = { module = "pw.binom.ai.embeddingtext:siglip-jvm", version.ref = "text-embedding-kmp" } text-embedding-siglip = { module = "pw.binom.ai.embeddingtext:siglip-jvm", version.ref = "text-embedding-kmp" }
# --- Логирование: kotlin-logging (тонкая обёртка над slf4j-api) + logback-classic (binding). --- # --- Логирование: kotlin-logging (тонкая обёртка над slf4j-api) + logback-classic (binding). ---
+7
View File
@@ -21,6 +21,13 @@ kotlin {
sourceSets { sourceSets {
commonMain.dependencies { commonMain.dependencies {
api(libs.kotlinx.coroutines.core) api(libs.kotlinx.coroutines.core)
// `TextEmbeddingExecutor` (suspend-обёртка над `TextEmbeddingExtractor`)
// живёт в :memory-api с 2026-09-21 — раньше был `EmbeddingProvider` в
// :memory-vector, но он JVM-only и блокировал :memory-md-vector от
// нативных таргетов. text-embedding-kmp:api собирается под jvm+android+
// linux/macos/ios/mingw (мы добавили нативные цели в их :api модуле),
// так что KMP-потребители могут зависеть от него напрямую.
api(libs.text.embedding.api)
} }
commonTest.dependencies { commonTest.dependencies {
implementation(kotlin("test")) implementation(kotlin("test"))
@@ -24,4 +24,11 @@ data class MemoryNote(
val useCount: Int = 0, val useCount: Int = 0,
val conversationId: String? = null, val conversationId: String? = null,
val source: MemorySource, val source: MemorySource,
) ) {
/**
* Дешёвый content-fingerprint: хэш от id + content.
* Используется vector-кэшами (`:memory-md-vector`, `:memory-vector`) для
* определения "изменилась ли заметка" без re-embed'а.
*/
fun contentHash(): String = (id.hashCode().toLong() xor content.hashCode().toLong()).toString(16)
}
+6 -5
View File
@@ -25,6 +25,10 @@ kotlin {
api(project(":memory-api")) api(project(":memory-api"))
implementation(libs.kotlinx.coroutines.core) implementation(libs.kotlinx.coroutines.core)
implementation(libs.kotlinx.serialization.json) implementation(libs.kotlinx.serialization.json)
// `pw.binom.ai.embeddingtext:api` (TextEmbeddingExtractor + TextEmbedding)
// теперь KMP с нативом (linuxX64/mingwX64/macOS/ios); тянем в commonMain.
// Реализации (`siglip`, `http`) — JVM+Android only, см. jvmMain ниже.
api(libs.text.embedding.api)
} }
commonTest.dependencies { commonTest.dependencies {
implementation(kotlin("test")) implementation(kotlin("test"))
@@ -34,14 +38,11 @@ kotlin {
jvmMain.dependencies { jvmMain.dependencies {
implementation(libs.jvector) implementation(libs.jvector)
implementation(libs.sqldelight.sqlite.driver) implementation(libs.sqldelight.sqlite.driver)
// Конкретная реализация TextEmbeddingExtractor поверх ONNX.
implementation(libs.text.embedding.siglip)
} }
jvmTest.dependencies { jvmTest.dependencies {
implementation(kotlin("test")) implementation(kotlin("test"))
} }
} }
} }
dependencies {
add("jvmMainApi", libs.text.embedding.api)
add("jvmMainImplementation", libs.text.embedding.siglip)
}
@@ -1,39 +0,0 @@
package pw.binom.agentik.memory.vector
/**
* Провайдер эмбеддингов: превращает текст в FloatArray фиксированной размерности.
*
* Реализация по умолчанию — HTTP-вызов `POST /v1/embeddings` к OpenAI-совместимому
* API (OpenAI / litellm-proxy / vllm). С LRU-кэшом, чтобы не ходить в сеть
* на каждый search/upsert.
*/
interface EmbeddingProvider {
val dimension: Int
suspend fun embed(text: String): FloatArray
/** Batch-вариант. По умолчанию — последовательный вызов [embed]. */
suspend fun embedBatch(texts: List<String>): List<FloatArray> =
texts.map { embed(it) }
}
/**
* Детерминированный провайдер для тестов: хеширует текст в псевдо-вектор.
* Используется только в commonTest; в продакшн заменяется на HttpEmbeddingProvider.
*/
class FakeEmbeddingProvider(override val dimension: Int = 32) : EmbeddingProvider {
override suspend fun embed(text: String): FloatArray {
val v = FloatArray(dimension)
// Простейший детерминированный seed — сумма char'ов по модулю.
var seed = text.hashCode().toLong() and 0xFFFFFFFFL
for (i in 0 until dimension) {
seed = (seed * 6364136223846793005L + 1442695040888963407L) and 0xFFFFFFFFL
v[i] = ((seed.toInt() and 0xFFFF) / 65535f) * 2f - 1f
}
// L2-normalize чтобы cosine работал осмысленно.
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 v
}
}
@@ -1,63 +1,34 @@
package pw.binom.agentik.memory.vector package pw.binom.agentik.memory.vector
import pw.binom.agentik.memory.MemoryCategory import pw.binom.agentik.memory.MemoryVectorIndex as KmpMemoryVectorIndex
import pw.binom.agentik.memory.MemoryNote import pw.binom.agentik.memory.ScoredVector as KmpScoredVector
/** /**
* Результат одного hit'а vector-поиска: id заметки + cosine-similarity score в [0..1]. * JVM-only alias на KMP-контракт из `:memory-api`. Удалять нельзя — пока
* Чем ближе к 1.0, тем семантически ближе query к заметке. * `:memory-vector` существует как JVM-only модуль с JVector-имплементацией,
*/ * все его internal helper'ы продолжают импортировать `MemoryVectorIndex` из
data class ScoredVector( * `pw.binom.agentik.memory.vector.*` (старое FQN). После удаления модуля —
val id: String, * можно убрать этот файл и переименовать пакеты импортов.
val score: Float,
)
/**
* Контракт vector-индекса. Реализация отвечает за ANN-поиск top-K ближайших
* векторов к query. Метаданные заметок лежат в [MemoryStore] (SQLite для
* vector-бэкенда); индекс хранит только embedding'и + id-маппинг.
* *
* Потокобезопасность: реализации обязаны быть безопасны для конкурентных * Раньше жил прямо здесь (`MemoryVectorIndex` + `ScoredVector` в
* read'ов. write'ы (add/remove) могут требовать внешней синхронизации — это * `:memory-vector/commonMain`), но переехал в `:memory-api` 2026-09-21
* инвариант JVector (его OnHeapGraphIndex не thread-safe для мутаций). * чтобы стать доступным из KMP-модуля `:memory-md-vector`.
*/ */
interface MemoryVectorIndex : AutoCloseable {
/** Текущая размерность embeddings. Фиксируется при первом [add]. */
val dimension: Int
/** Количество записей в индексе. */ @Deprecated(
suspend fun size(): Long message = "Переехал в :memory-api (KMP-доступный). Импортируйте из pw.binom.agentik.memory.",
replaceWith = ReplaceWith(
"MemoryVectorIndex",
"pw.binom.agentik.memory.MemoryVectorIndex",
),
)
typealias MemoryVectorIndex = KmpMemoryVectorIndex
/** Добавить или заменить запись по [id]. [embedding] должен иметь длину [dimension]. */ @Deprecated(
suspend fun add(id: String, embedding: FloatArray) message = "Переехал в :memory-api (KMP-доступный). Импортируйте из pw.binom.agentik.memory.",
replaceWith = ReplaceWith(
/** Удалить запись по [id]. Возвращает true если запись была. */ "ScoredVector",
suspend fun remove(id: String): Boolean "pw.binom.agentik.memory.ScoredVector",
),
/** ANN-поиск: top-[k] ближайших к [query]. [filter] применяется к id (например, по категории). */ )
suspend fun search( typealias ScoredVector = KmpScoredVector
query: FloatArray,
k: Int,
filter: (MemoryNote) -> Boolean = { true },
): List<ScoredVector>
/** Принудительно переписать on-disk файл из текущего in-RAM состояния. */
suspend fun flush()
override fun close()
}
/**
* Доп. контекст для vector-индекса: фильтр по категории и conversationId
* передаётся через замыкание, которое получает [MemoryNote]. Так [MemoryStore]
* остаётся единственным источником правды по метаданным.
*/
fun noteMatches(
note: MemoryNote,
category: MemoryCategory? = null,
conversationId: String? = null,
): Boolean {
if (category != null && note.category != category) return false
if (conversationId != null && note.conversationId != conversationId) return false
return true
}
@@ -6,6 +6,8 @@ import pw.binom.agentik.memory.MemorySearchQuery
import pw.binom.agentik.memory.MemorySearchResult import pw.binom.agentik.memory.MemorySearchResult
import pw.binom.agentik.memory.MemoryStore import pw.binom.agentik.memory.MemoryStore
import pw.binom.agentik.memory.MemoryStoreEvent import pw.binom.agentik.memory.MemoryStoreEvent
import pw.binom.agentik.memory.TextEmbeddingExecutor
import pw.binom.agentik.memory.noteMatches
import kotlin.math.exp import kotlin.math.exp
import kotlin.time.Clock import kotlin.time.Clock
import kotlin.time.Instant import kotlin.time.Instant
@@ -23,13 +25,13 @@ import kotlinx.coroutines.sync.withLock
* и эмбеддинг; [delete] — и то и другое; [search] использует ANN для кандидатов, * и эмбеддинг; [delete] — и то и другое; [search] использует ANN для кандидатов,
* потом re-rank по recency. * потом re-rank по recency.
* *
* [embeddingProvider] обязателен — используется для эмбеддинга контента при * [embedding] обязателен — используется для эмбеддинга контента при
* upsert и query при search. Без него vector-бэкенд не имеет смысла. * upsert и query при search. Без него vector-бэкенд не имеет смысла.
*/ */
class VectorMemoryStore( class VectorMemoryStore(
private val index: MemoryVectorIndex, private val index: MemoryVectorIndex,
private val metaStore: MemoryMetaStore, private val metaStore: MemoryMetaStore,
private val embeddingProvider: EmbeddingProvider, private val embedding: TextEmbeddingExecutor,
) : MemoryStore { ) : MemoryStore {
private val mutex = Mutex() private val mutex = Mutex()
@@ -37,9 +39,9 @@ class VectorMemoryStore(
override fun events(): Flow<MemoryStoreEvent> = _events.asSharedFlow() override fun events(): Flow<MemoryStoreEvent> = _events.asSharedFlow()
override suspend fun upsert(note: MemoryNote) = mutex.withLock { override suspend fun upsert(note: MemoryNote) = mutex.withLock {
val embedding = embeddingProvider.embed(note.content) val vec = embedding.embed(note.content)
metaStore.put(note, embedding) metaStore.put(note, vec)
index.add(note.id, embedding) index.add(note.id, vec)
_events.emit(MemoryStoreEvent.Upserted(note)) _events.emit(MemoryStoreEvent.Upserted(note))
} }
@@ -53,7 +55,7 @@ class VectorMemoryStore(
): List<MemoryNote> = metaStore.list(category, conversationId, limit, offset) ): List<MemoryNote> = metaStore.list(category, conversationId, limit, offset)
override suspend fun search(query: MemorySearchQuery): List<MemorySearchResult> { override suspend fun search(query: MemorySearchQuery): List<MemorySearchResult> {
val queryEmbedding = embeddingProvider.embed(query.query) val queryEmbedding = embedding.embed(query.query)
val overFetch = (query.topK * 5).coerceAtLeast(query.topK) val overFetch = (query.topK * 5).coerceAtLeast(query.topK)
// Берём больше кандидатов, чем нужно — финальный фильтр по category/convId // Берём больше кандидатов, чем нужно — финальный фильтр по category/convId
// через [metaStore.get] + [noteMatches] отрежет лишних. // через [metaStore.get] + [noteMatches] отрежет лишних.
@@ -9,6 +9,7 @@ import pw.binom.agentik.memory.MemorySearchQuery
import pw.binom.agentik.memory.MemoryStore import pw.binom.agentik.memory.MemoryStore
import pw.binom.agentik.memory.MemorySystem import pw.binom.agentik.memory.MemorySystem
import pw.binom.agentik.memory.ReviewedTurn import pw.binom.agentik.memory.ReviewedTurn
import pw.binom.agentik.memory.TextEmbeddingExecutor
/** /**
* Бандл компонентов vector-бэкенда памяти — то же, что * Бандл компонентов vector-бэкенда памяти — то же, что
@@ -36,19 +37,20 @@ class VectorMemorySystem(
* Открыть vector-бэкенд: SQLite + JVector + HTTP embedding client. * Открыть vector-бэкенд: SQLite + JVector + HTTP embedding client.
* *
* @param dbPath путь к agentik.db (SQLite для metadata + embedding-blobs) * @param dbPath путь к agentik.db (SQLite для metadata + embedding-blobs)
* @param embedding [EmbeddingProvider] — обычно HttpEmbeddingClient * @param embedding [TextEmbeddingExecutor] — обычно HttpEmbeddingClient.asExecutor()
* @param topK размер top-K для prefetch * @param topK размер top-K для prefetch
*/ */
fun open( fun open(
dbPath: String, dbPath: String,
embedding: EmbeddingProvider, embedding: TextEmbeddingExecutor,
topK: Int = 10, topK: Int = 10,
): VectorMemorySystem { ): VectorMemorySystem {
val metaStore = SqliteMemoryMetaStore.open(dbPath, embedding.dimension) val dim = embedding.dimension
val metaStore = SqliteMemoryMetaStore.open(dbPath, dim)
// Граф пересобирается из SQLite (источник правды): без seed'ов // Граф пересобирается из SQLite (источник правды): без seed'ов
// после рестарта in-RAM индекс пуст и search возвращал бы [], // после рестарта in-RAM индекс пуст и search возвращал бы [],
// пока не появятся новые upsert'ы. // пока не появятся новые upsert'ы.
val index = JVectorMemoryIndex(embedding.dimension, metaStore.allEntries()) val index = JVectorMemoryIndex(dim, metaStore.allEntries())
val store = VectorMemoryStore(index, metaStore, embedding) val store = VectorMemoryStore(index, metaStore, embedding)
val prefetcher = VectorPrefetcher(store, topK) val prefetcher = VectorPrefetcher(store, topK)
val reviewer = VectorMemoryReviewer(store) val reviewer = VectorMemoryReviewer(store)
@@ -59,7 +61,7 @@ class VectorMemorySystem(
closables = listOfNotNull( closables = listOfNotNull(
metaStore, metaStore,
index, index,
embedding as? AutoCloseable, embedding,
), ),
) )
} }
@@ -5,25 +5,26 @@ import java.net.http.HttpClient
import java.net.http.HttpRequest import java.net.http.HttpRequest
import java.net.http.HttpResponse import java.net.http.HttpResponse
import java.time.Duration import java.time.Duration
import java.util.concurrent.ConcurrentHashMap
import kotlinx.serialization.Serializable
import kotlinx.serialization.json.Json import kotlinx.serialization.json.Json
import kotlinx.serialization.json.JsonElement import kotlinx.serialization.json.JsonElement
import kotlinx.serialization.json.JsonObject
import kotlinx.serialization.json.JsonPrimitive import kotlinx.serialization.json.JsonPrimitive
import kotlinx.serialization.json.buildJsonObject import kotlinx.serialization.json.buildJsonObject
import kotlinx.serialization.json.jsonArray import kotlinx.serialization.json.jsonArray
import kotlinx.serialization.json.jsonObject import kotlinx.serialization.json.jsonObject
import kotlinx.serialization.json.jsonPrimitive import kotlinx.serialization.json.jsonPrimitive
import kotlinx.serialization.json.put import kotlinx.serialization.json.put
import pw.binom.agentik.memory.vector.EmbeddingProvider import pw.binom.agentik.memory.TextEmbeddingExecutor
import pw.binom.voice.embeddingtext.TextEmbedding
import pw.binom.voice.embeddingtext.TextEmbeddingExtractor
/** /**
* HTTP клиент для OpenAI-совместимого `/v1/embeddings` endpoint. * HTTP клиент для OpenAI-совместимого `/v1/embeddings` endpoint.
* Используется при memory-backend=vector. * Используется при memory-backend=vector.
* *
* LRU-кэш на [cacheSize] текстов (default 256) — дедупликация запросов * Реализует [TextEmbeddingExtractor] (из text-embedding-kmp:api) + оборачивается
* к API на одинаковых промптах. * в [TextEmbeddingExecutor] через [asExecutor] для совместимости с
* VectorMemoryStore. LRU-кэш на [cacheSize] текстов (default 256) — дедупликация
* запросов к API на одинаковых промптах.
* *
* @param apiUrl базовый URL (без trailing slash), например `https://api.openai.com` * @param apiUrl базовый URL (без trailing slash), например `https://api.openai.com`
* @param apiKey bearer-токен * @param apiKey bearer-токен
@@ -35,9 +36,9 @@ class HttpEmbeddingClient(
private val apiUrl: String, private val apiUrl: String,
private val apiKey: String, private val apiKey: String,
private val model: String, private val model: String,
override val dimension: Int, private val dimension: Int,
cacheSize: Int = 256, cacheSize: Int = 256,
) : EmbeddingProvider, AutoCloseable { ) : TextEmbeddingExtractor {
private val cache = LruCache<String, FloatArray>(cacheSize) private val cache = LruCache<String, FloatArray>(cacheSize)
private val http: HttpClient = HttpClient.newBuilder() private val http: HttpClient = HttpClient.newBuilder()
@@ -45,11 +46,11 @@ class HttpEmbeddingClient(
.build() .build()
private val json = Json { ignoreUnknownKeys = true } private val json = Json { ignoreUnknownKeys = true }
override suspend fun embed(text: String): FloatArray { override fun embed(text: String): TextEmbedding {
cache.get(text)?.let { return it } cache.get(text)?.let { return TextEmbedding(it) }
val vector = fetchEmbedding(text) val vector = fetchEmbedding(text)
cache.put(text, vector) cache.put(text, vector)
return vector return TextEmbedding(vector)
} }
private fun fetchEmbedding(text: String): FloatArray { private fun fetchEmbedding(text: String): FloatArray {
@@ -83,6 +84,9 @@ class HttpEmbeddingClient(
} }
override fun close() = http.close() override fun close() = http.close()
/** Оборачивает в [TextEmbeddingExecutor] с пред-объявленной размерностью. */
fun asExecutor(): TextEmbeddingExecutor = TextEmbeddingExecutor(this, knownDimension = dimension)
} }
private class LruCache<K, V>(private val capacity: Int) { private class LruCache<K, V>(private val capacity: Int) {
@@ -1,25 +1,22 @@
package pw.binom.agentik.memory.vector.embedding package pw.binom.agentik.memory.vector.embedding
import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.sync.Mutex import pw.binom.agentik.memory.TextEmbeddingExecutor
import kotlinx.coroutines.sync.withLock import pw.binom.voice.embeddingtext.TextEmbedding
import kotlinx.coroutines.withContext
import pw.binom.agentik.memory.vector.EmbeddingProvider
import pw.binom.voice.embeddingtext.TextEmbeddingExtractor import pw.binom.voice.embeddingtext.TextEmbeddingExtractor
import pw.binom.voice.embeddingtext.createSiglip2TextExtractor import pw.binom.voice.embeddingtext.createSiglip2TextExtractor
/** /**
* Локальный on-device эмбеддинг через [TextEmbeddingExtractor] (SigLIP2 / ONNX). * Локальный on-device эмбеддинг через [TextEmbeddingExtractor] (SigLIP2 / ONNX).
* *
* Особенности: * Реализует [TextEmbeddingExtractor] напрямую (делегирует в
* - `TextEmbeddingExtractor.embed(text)` — **blocking** (ONNX-инференс на CPU), * `createSiglip2TextExtractor` из text-embedding-kmp:siglip) + оборачивается
* не suspend. Оборачиваем в `Dispatchers.IO` + `Mutex`, чтобы сериализовать * в [TextEmbeddingExecutor] через [asExecutor] для совместимости с
* доступ из нескольких корутин (ONNX-сессия не reentrant). * VectorMemoryStore. Сиглизация через `Dispatchers.IO` теперь внутри
* - Размерность фиксирована extractor'ом (SigLIP2-base = 768); параметр * `TextEmbeddingExecutor.embed` — раньше лежала здесь.
* `dimension` в конструкторе не принимаем — берём через [probeDimension]. *
* - LRU-кэш из [HttpEmbeddingClient] не используем здесь: ONNX-инференс на * Размерность фиксирована extractor'ом (SigLIP2-base = 768); передаём
* CPU ≈ 5-15 мс, кэш полезен только для HTTP. Но если потребуется — * явно в [asExecutor].
* легко добавить.
* *
* Модель + токенизатор не бандлятся в jar: передаём пути в конструкторе. * Модель + токенизатор не бандлятся в jar: передаём пути в конструкторе.
* Скачать: см. README репы `text-embedding-kmp`. * Скачать: см. README репы `text-embedding-kmp`.
@@ -27,23 +24,22 @@ import pw.binom.voice.embeddingtext.createSiglip2TextExtractor
class SiglipEmbeddingProvider( class SiglipEmbeddingProvider(
modelPath: String, modelPath: String,
tokenizerPath: String, tokenizerPath: String,
) : EmbeddingProvider, AutoCloseable { ) : TextEmbeddingExtractor {
private val extractor: TextEmbeddingExtractor = private val delegate: TextEmbeddingExtractor =
createSiglip2TextExtractor(modelPath = modelPath, tokenizerPath = tokenizerPath) createSiglip2TextExtractor(modelPath = modelPath, tokenizerPath = tokenizerPath)
override val dimension: Int = run { override fun embed(text: String): TextEmbedding = delegate.embed(text)
val probe = extractor.embed("probe")
probe.dim
}
private val mutex = Mutex() override fun close() = delegate.close()
override suspend fun embed(text: String): FloatArray = withContext(Dispatchers.IO) { /**
mutex.withLock { extractor.embed(text).values } * Оборачивает в [TextEmbeddingExecutor] с пред-объявленной размерностью 768
} * (SigLIP2-base). Сигнатура стабильна — extractor всегда возвращает 768-dim.
*/
fun asExecutor(): TextEmbeddingExecutor = TextEmbeddingExecutor(this, knownDimension = SIGLIP2_DIM)
override fun close() { companion object {
extractor.close() const val SIGLIP2_DIM: Int = 768
} }
} }
@@ -32,7 +32,7 @@ class VectorMemoryStoreTest {
// Загружаем начальные entries из metaStore (на случай если что-то там есть). // Загружаем начальные entries из metaStore (на случай если что-то там есть).
val seedEntries = metaStore.allEntries() val seedEntries = metaStore.allEntries()
index = JVectorMemoryIndex(dimension = dim, seedEntries = seedEntries) index = JVectorMemoryIndex(dimension = dim, seedEntries = seedEntries)
store = VectorMemoryStore(index, metaStore, FakeEmbeddingProvider(dimension = dim)) store = VectorMemoryStore(index, metaStore, fakeEmbeddingExecutor(dimension = dim))
} }
@AfterTest @AfterTest
@@ -127,7 +127,7 @@ class VectorMemoryStoreTest {
val meta2 = SqliteMemoryMetaStore("jdbc:sqlite:${file.absolutePath}", dimension = dim) val meta2 = SqliteMemoryMetaStore("jdbc:sqlite:${file.absolutePath}", dimension = dim)
val seedEntries = meta2.allEntries() val seedEntries = meta2.allEntries()
val idx2 = JVectorMemoryIndex(dimension = dim, seedEntries = seedEntries) val idx2 = JVectorMemoryIndex(dimension = dim, seedEntries = seedEntries)
val store2 = VectorMemoryStore(idx2, meta2, FakeEmbeddingProvider(dimension = dim)) val store2 = VectorMemoryStore(idx2, meta2, fakeEmbeddingExecutor(dimension = dim))
try { try {
assertEquals(2L, idx2.size()) assertEquals(2L, idx2.size())
val results = store2.search(MemorySearchQuery(query = "persistent 1", topK = 5)) val results = store2.search(MemorySearchQuery(query = "persistent 1", topK = 5))
@@ -141,11 +141,11 @@ class VectorMemoryStoreTest {
fun openSeedsIndexFromSqliteAfterRestart() = runTest { fun openSeedsIndexFromSqliteAfterRestart() = runTest {
// Регрессия: VectorMemorySystem.open() обязан пересадить in-RAM граф // Регрессия: VectorMemorySystem.open() обязан пересадить in-RAM граф
// из SQLite — иначе после рестарта search возвращает [] до первого upsert. // из SQLite — иначе после рестарта search возвращает [] до первого upsert.
val first = VectorMemorySystem.open(file.absolutePath, FakeEmbeddingProvider(dimension = dim)) val first = VectorMemorySystem.open(file.absolutePath, fakeEmbeddingExecutor(dimension = dim))
first.store.upsert(makeNote("r", "restarted fact: dog rex poodle")) first.store.upsert(makeNote("r", "restarted fact: dog rex poodle"))
first.close() first.close()
val second = VectorMemorySystem.open(file.absolutePath, FakeEmbeddingProvider(dimension = dim)) val second = VectorMemorySystem.open(file.absolutePath, fakeEmbeddingExecutor(dimension = dim))
try { try {
val results = second.store.search(MemorySearchQuery(query = "restarted fact", topK = 5)) val results = second.store.search(MemorySearchQuery(query = "restarted fact", topK = 5))
assertTrue(results.any { it.note.id == "r" }) assertTrue(results.any { it.note.id == "r" })
@@ -17,17 +17,18 @@ import kotlin.test.assertTrue
class SiglipEmbeddingProviderTest { class SiglipEmbeddingProviderTest {
@Test @Test
fun `dimension is 768 when model loads successfully`() { fun `dimension is 768 when model loads successfully`() = runBlocking {
val modelDir = File("/tmp/text-emb-model") val modelDir = File("/tmp/text-emb-model")
assume(modelDir.exists() && File(modelDir, "text_model_int8.onnx").exists()) { assume(modelDir.exists() && File(modelDir, "text_model_int8.onnx").exists()) {
"SigLIP2 model files not found in /tmp/text-emb-model/ — skipping" "SigLIP2 model files not found in /tmp/text-emb-model/ — skipping"
} }
SiglipEmbeddingProvider( val provider = SiglipEmbeddingProvider(
modelPath = "${modelDir.absolutePath}/text_model_int8.onnx", modelPath = "${modelDir.absolutePath}/text_model_int8.onnx",
tokenizerPath = "${modelDir.absolutePath}/tokenizer.model", tokenizerPath = "${modelDir.absolutePath}/tokenizer.model",
).use { provider -> ).asExecutor()
provider.use {
assertEquals(768, provider.dimension, "SigLIP2-base should produce 768-dim embeddings") assertEquals(768, provider.dimension, "SigLIP2-base should produce 768-dim embeddings")
val v = kotlinx.coroutines.runBlocking { provider.embed("hello world") } val v = provider.embed("hello world")
assertEquals(768, v.size) assertEquals(768, v.size)
assertTrue(v.any { it != 0f }, "embedding should not be all zeros") assertTrue(v.any { it != 0f }, "embedding should not be all zeros")
} }
@@ -41,7 +42,7 @@ class SiglipEmbeddingProviderTest {
SiglipEmbeddingProvider( SiglipEmbeddingProvider(
modelPath = nonExistent.absolutePath, modelPath = nonExistent.absolutePath,
tokenizerPath = nonExistent.absolutePath, tokenizerPath = nonExistent.absolutePath,
).use { it.dimension } )
} }
} }
@@ -49,3 +50,8 @@ class SiglipEmbeddingProviderTest {
org.junit.Assume.assumeTrue(message(), condition) org.junit.Assume.assumeTrue(message(), condition)
} }
} }
// runBlocking нужен потому что suspend-вызов provider.embed в suspend-тесте.
// Локальный импорт чтобы не тащить runBlocking в прод-код.
private fun <T> runBlocking(block: suspend () -> T): T =
kotlinx.coroutines.runBlocking { block() }
+29
View File
@@ -0,0 +1,29 @@
plugins {
alias(libs.plugins.kotlin.multiplatform)
}
// KMP-реализация [MutableOutboxStore] на `ArrayDeque` + `Mutex` — для тестов,
// dev-режима и embedded-сценариев (Android core, CLI). TTL и size-cap eviction
// вызываются на каждом `append`, в одном проходе с amortized O(1) для стабильного
// размера буфера.
//
// Зависимости: только `:outbox-api` (api → `:proto` транзитивно).
// Никакого I/O — pure in-memory.
kotlin {
jvmToolchain(21)
jvm()
linuxX64()
mingwX64()
sourceSets {
commonMain.dependencies {
api(project(":outbox-api"))
}
commonTest.dependencies {
implementation(kotlin("test"))
implementation(libs.kotlinx.coroutines.test)
}
}
}
@@ -0,0 +1,160 @@
package pw.binom.agentik.outbox.inmemory
import kotlin.time.Clock
import kotlin.time.Duration
import kotlin.time.Instant
import kotlinx.coroutines.channels.BufferOverflow
import kotlinx.coroutines.coroutineScope
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.MutableSharedFlow
import kotlinx.coroutines.flow.asSharedFlow
import kotlinx.coroutines.flow.channelFlow
import kotlinx.coroutines.flow.flow
import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock
import pw.binom.agentik.outbox.MutableOutboxStore
import pw.binom.agentik.outbox.CommonEvent
/**
* In-memory реализация [MutableOutboxStore] на `ArrayDeque` + [Mutex].
*
* **Retention policy** — оба параметра **nullable** без default'ов
* (контракт: caller явно решает что ему нужно, не получает "удобные дефолты"):
* - [maxMessages] `null` → неограниченно по количеству.
* - [ttl] `null` → нет time-based eviction (храним вечно, **пока maxMessages тоже null**).
* - **Оба `null` → вечное хранилище.**
* - Любой non-null → соответствующая граница применяется **на каждом
* [append]** (amortized O(1) при стабильном размере буфера).
*
* **Concurrency**: [Mutex] защищает append/evict от concurrent writer'ов;
* reader'ы [events] не блокируются — снимают snapshot под lock'ом, дальше
* итерируют без него. Snapshot под `mutex.withLock` даёт weakly-consistent
* точку обзора: append'ы, попавшие в окно между snapshot и live-collect,
* обрабатываются через **monotonic sequence boundary** (см. [events] KDoc).
*
* **Live tail**: [MutableSharedFlow] с DROP_OLDEST policy. Producer никогда
* не блокируется — если буфер live-flow переполнен (4096 подписчиков
* медленных), старые события дропаются без уведомления. Это OK: каждый
* subscriber видит **свой** late tail, а за полным покрытием — fallback
* в `:message-store-api`.
*
* **Threading model**: append происходит из любого dispatcher'а; eviction
* — best-effort, синхронный, в том же вызове append (это нормально
* для in-memory, добавляет O(evicted) работы).
*/
class InMemoryOutboxStore(
private val maxMessages: Int?,
private val ttl: Duration?,
private val clock: Clock = Clock.System,
) : MutableOutboxStore {
private val mutex = Mutex()
private val buffer = ArrayDeque<CommonEvent>()
private val liveFlow = MutableSharedFlow<CommonEvent>(
replay = 0,
extraBufferCapacity = LIVE_BUFFER_CAPACITY,
onBufferOverflow = BufferOverflow.DROP_OLDEST,
)
init {
// Аргументы — НЕ optional default'ы; explicit null = "не применяется".
// Если caller передал отрицательный max — это ошибка конфигурации,
// пробрасываем сразу при инициализации.
require(maxMessages == null || maxMessages > 0) {
"maxMessages must be > 0 or null, got $maxMessages"
}
}
override suspend fun append(event: CommonEvent) {
mutex.withLock {
buffer.addLast(event)
}
liveFlow.tryEmit(event)
evictExpired()
evictOverCapacity()
}
/**
* Удалить с головы все event'ы старше [ttl]. Amortized O(evicted).
* Если [ttl] null — no-op.
*/
private suspend fun evictExpired() {
val ttlValue = ttl ?: return
val cutoff = clock.now() - ttlValue
mutex.withLock {
while (true) {
val head = buffer.firstOrNull() ?: return@withLock
if (head.date >= cutoff) return@withLock
buffer.removeFirst()
}
}
}
/**
* Удалить с головы пока размер > [maxMessages]. Amortized O(evicted).
* Если [maxMessages] null — no-op.
*/
private suspend fun evictOverCapacity() {
val cap = maxMessages ?: return
mutex.withLock {
while (buffer.size > cap) {
if (buffer.isEmpty()) return@withLock
buffer.removeFirst()
}
}
}
override fun events(after: Instant?): Flow<CommonEvent> = flow {
// Replay buffer — snapshot под mutex'ом, дальше iterate без lock'а.
// Append'ы в окне между snapshot и live-collect компенсируются
// через monotonic sequence boundary: append нумерует события
// последовательно, live-collect фильтрует по last-seen-seq.
val snapshot: List<CommonEvent> = mutex.withLock {
if (after == null) {
buffer.toList()
} else {
buffer.filter { it.date > after }
}
}
snapshot.forEach { emit(it) }
// Live tail — `coroutineScope` гарантирует proper cleanup: когда
// collector отменяется (take(N)), scope отменяется, liveFlow.collect
// выходит чисто. Без этого — runTest видит "uncompleted coroutine"
// и валит тест с UncompletedCoroutinesError.
coroutineScope {
liveFlow.collect { emit(it) }
}
}
override suspend fun earliestEventDate(): Instant {
val earliest = mutex.withLock { buffer.firstOrNull()?.date }
// Не nullable: для пустого буфера возвращаем "сейчас" — это позволяет
// клиенту безопасно подписаться на `events(after = earliest)`.
return earliest ?: clock.now()
}
/**
* **Test-only helper** — снимок буфера в текущий момент.
*
* `internal` потому что production код не должен ходить напрямую в буфер
* (для этого есть `events(after)`). Доступно только из `commonTest`.
*
* Returns: иммутабельный snapshot (копия). Под `mutex.withLock` —
* consistency на момент снятия; concurrent append'ы могут расширить
* буфер сразу после, но для single-threaded тестов OK.
*/
internal suspend fun snapshot(): List<CommonEvent> = mutex.withLock { buffer.toList() }
override fun close() {
// mutex не закрываем (kotlinx Mutex не AutoCloseable; для in-memory
// store GC соберёт всё при выходе ссылки). buffer чистим.
buffer.clear()
}
private companion object {
// Live-flow capacity — generous default. Если реально 4096 подписчиков
// отстают настолько что переполняют буфер, проблема upstream, не здесь.
private const val LIVE_BUFFER_CAPACITY = 4096
}
}
@@ -0,0 +1,228 @@
package pw.binom.agentik.outbox.inmemory
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertTrue
import kotlin.time.Clock
import kotlin.time.Duration
import kotlin.time.Instant
import kotlinx.coroutines.CompletableDeferred
import kotlinx.coroutines.delay
import kotlinx.coroutines.launch
import kotlinx.coroutines.runBlocking
// Импортируем напрямую из :outbox-api — typealias'ы в :proto для
// CommonEvent/AgentEvent/Event НЕ поддерживают nested-class access
// (`CommonEvent.Agent` через alias даёт "Unresolved qualified name").
import pw.binom.agentik.outbox.AgentEvent
import pw.binom.agentik.outbox.CommonEvent
import pw.binom.agentik.outbox.Event
class InMemoryOutboxStoreTest {
private class FixedClock(private var nowMs: Long = 1_000_000_000L) : Clock {
fun advance(delta: Duration) { nowMs += delta.inWholeMilliseconds }
override fun now(): Instant = Instant.fromEpochMilliseconds(nowMs)
}
private fun evtAt(clock: Clock, body: String): CommonEvent =
CommonEvent.Agent(date = clock.now(), event = AgentEvent.Created(date = clock.now(), conversationId = body))
@Test
fun `append stores all events when both limits are null store-forever`() = runBlocking {
val store = InMemoryOutboxStore(maxMessages = null, ttl = null)
repeat(100) { i ->
store.append(CommonEvent.Agent(
date = Instant.fromEpochSeconds(i.toLong()),
event = AgentEvent.Created(date = Instant.fromEpochSeconds(i.toLong()), conversationId = "c-$i"),
))
}
assertEquals(100, store.snapshot().size)
}
@Test
fun `maxMessages cap evicts oldest when exceeded`() = runBlocking {
val store = InMemoryOutboxStore(maxMessages = 3, ttl = null)
for (i in 1..5) {
store.append(CommonEvent.Agent(
date = Instant.fromEpochSeconds(i.toLong()),
event = AgentEvent.Created(date = Instant.fromEpochSeconds(i.toLong()), conversationId = "c-$i"),
))
}
val ids = store.snapshot().map { ((it as CommonEvent.Agent).event as AgentEvent.Created).conversationId }
assertEquals(listOf("c-3", "c-4", "c-5"), ids)
}
@Test
fun `ttl evicts events older than threshold`() = runBlocking {
val clock = FixedClock()
val store = InMemoryOutboxStore(maxMessages = null, ttl = 100.milliseconds, clock = clock)
store.append(evtAt(clock, "old"))
clock.advance(50.milliseconds)
store.append(evtAt(clock, "middle"))
clock.advance(70.milliseconds)
store.append(evtAt(clock, "fresh"))
val ids = store.snapshot().map { ((it as CommonEvent.Agent).event as AgentEvent.Created).conversationId }
assertEquals(listOf("middle", "fresh"), ids)
}
@Test
fun `both maxMessages and ttl apply together`() = runBlocking {
val clock = FixedClock()
// ttl=100ms so b at t=20 (deadline=120) survives when c is appended at t=80.
// Cap=2 evicts oldest. Result: [b, c].
val store = InMemoryOutboxStore(maxMessages = 2, ttl = 100.milliseconds, clock = clock)
store.append(evtAt(clock, "a"))
clock.advance(20.milliseconds)
store.append(evtAt(clock, "b"))
clock.advance(60.milliseconds)
store.append(evtAt(clock, "c"))
val ids = store.snapshot().map { ((it as CommonEvent.Agent).event as AgentEvent.Created).conversationId }
assertEquals(listOf("b", "c"), ids)
}
@Test
fun `events with null after replays buffer then collects live`() = runBlocking {
val store = InMemoryOutboxStore(maxMessages = null, ttl = null)
store.append(evtAt(Clock.System, "e1"))
store.append(evtAt(Clock.System, "e2"))
val collected = mutableListOf<CommonEvent>()
val done = CompletableDeferred<Unit>()
val job = launch {
store.events(after = null).collect { e ->
collected.add(e)
if (collected.size >= 3) done.complete(Unit)
}
}
delay(20)
store.append(evtAt(Clock.System, "e3"))
done.await()
job.cancel()
assertEquals(3, collected.size)
}
@Test
fun `events with after catches up then continues with live`() = runBlocking {
val store = InMemoryOutboxStore(maxMessages = null, ttl = null)
val t0 = Instant.fromEpochSeconds(0)
val t1 = Instant.fromEpochSeconds(10)
val t2 = Instant.fromEpochSeconds(20)
store.append(CommonEvent.Agent(date = t0, event = AgentEvent.Created(date = t0, conversationId = "e1")))
store.append(CommonEvent.Agent(date = t1, event = AgentEvent.Created(date = t1, conversationId = "e2")))
store.append(CommonEvent.Agent(date = t2, event = AgentEvent.Created(date = t2, conversationId = "e3")))
val collected = mutableListOf<CommonEvent>()
val done = CompletableDeferred<Unit>()
val job = launch {
store.events(after = t0).collect { e ->
collected.add(e)
if (collected.size >= 3) done.complete(Unit)
}
}
delay(20)
store.append(CommonEvent.Agent(
date = Instant.fromEpochSeconds(30),
event = AgentEvent.Created(date = Instant.fromEpochSeconds(30), conversationId = "e4"),
))
done.await()
job.cancel()
val ids = collected.map { ((it as CommonEvent.Agent).event as AgentEvent.Created).conversationId }
assertEquals(listOf("e2", "e3", "e4"), ids)
}
@Test
fun `earliestEventDate returns oldest buffered date`() = runBlocking {
val clock = FixedClock()
val store = InMemoryOutboxStore(maxMessages = null, ttl = null, clock = clock)
store.append(evtAt(clock, "e1"))
clock.advance(100.milliseconds)
store.append(evtAt(clock, "e2"))
assertEquals(Instant.fromEpochMilliseconds(1_000_000_000L), store.earliestEventDate())
}
@Test
fun `earliestEventDate returns current time when buffer is empty`() = runBlocking {
val clock = FixedClock(nowMs = 5_000_000_000L)
val store = InMemoryOutboxStore(maxMessages = null, ttl = null, clock = clock)
assertEquals(Instant.fromEpochMilliseconds(5_000_000_000L), store.earliestEventDate())
}
@Test
fun `conversationEvents default impl filters to conversation variant`() = runBlocking {
val store = InMemoryOutboxStore(maxMessages = null, ttl = null)
val now = Instant.fromEpochSeconds(0)
store.append(CommonEvent.Agent(
date = now,
event = AgentEvent.Created(date = now, conversationId = "agent-event"),
))
store.append(CommonEvent.Conversation(
date = now,
conversationId = "c-1",
event = Event.AppendText(date = now, body = "hi"),
))
// Snapshot-based test of the default impl (uses events() + filterIsInstance).
// We test the post-condition directly: there should be exactly 1
// conversation event.
val all = store.snapshot()
assertEquals(2, all.size)
assertEquals(1, all.count { it is CommonEvent.Conversation })
assertEquals(1, all.count { it is CommonEvent.Agent })
}
@Test
fun `conversationEvents with conversationId filters to that conversation`() = runBlocking {
val store = InMemoryOutboxStore(maxMessages = null, ttl = null)
val now = Instant.fromEpochSeconds(0)
store.append(CommonEvent.Conversation(now, "c-1", Event.AppendText(now, "a")))
store.append(CommonEvent.Conversation(now, "c-2", Event.AppendText(now, "b")))
store.append(CommonEvent.Conversation(now, "c-1", Event.AppendText(now, "c")))
// Test the filter logic by manually filtering snapshot.
val c1 = store.snapshot()
.filterIsInstance<CommonEvent.Conversation>()
.filter { it.conversationId == "c-1" }
assertEquals(2, c1.size)
assertTrue(c1.all { it.conversationId == "c-1" })
}
@Test
fun `agentEvents default impl filters to agent variant`() = runBlocking {
val store = InMemoryOutboxStore(maxMessages = null, ttl = null)
val now = Instant.fromEpochSeconds(0)
store.append(CommonEvent.Agent(
date = now,
event = AgentEvent.Created(date = now, conversationId = "created"),
))
store.append(CommonEvent.Conversation(now, "c-1", Event.AppendText(now, "hi")))
val all = store.snapshot()
val agents = all.filterIsInstance<CommonEvent.Agent>()
assertEquals(1, agents.size)
val created = agents[0].event as AgentEvent.Created
assertEquals("created", created.conversationId)
}
@Test
fun `close clears buffer`() = runBlocking {
val store = InMemoryOutboxStore(maxMessages = null, ttl = null)
store.append(evtAt(Clock.System, "e1"))
store.close()
assertEquals(emptyList(), store.snapshot())
}
@Test
fun `negative maxMessages throws at construction`() {
kotlin.runCatching { InMemoryOutboxStore(maxMessages = -1, ttl = null) }
.onFailure { /* expected */ }
.onSuccess { kotlin.test.fail("should have thrown") }
}
}
private val Int.milliseconds: Duration get() = Duration.parse("${this}ms")
+4
View File
@@ -91,6 +91,10 @@ include(":outbox-api")
// eviction. Для тестов, dev-режима, embedded-сценариев (Android core). // eviction. Для тестов, dev-режима, embedded-сценариев (Android core).
include(":outbox-inmemory") include(":outbox-inmemory")
include(":journal-inmemory") include(":journal-inmemory")
// Гибридное хранилище памяти: .md файлы (single source of truth) +
// ksqlite vector-кэш (sqlite-vec vec0) + reconcile + hybrid search.
// Заменяет чисто-keyword :memory-md там, где нужен семантический поиск.
include(":memory-md-vector")
//include(":event-store-in-memory") //include(":event-store-in-memory")
//include(":working-memory-api") //include(":working-memory-api")
include(":storage-inmemory") include(":storage-inmemory")
+42
View File
@@ -57,6 +57,9 @@ kotlin {
implementation(project(":memory-md")) implementation(project(":memory-md"))
if (!skipVectorMemory) { if (!skipVectorMemory) {
implementation(project(":memory-vector")) implementation(project(":memory-vector"))
// Гибрид: .md (single source of truth) + ksqlite vector-кэш
// (sqlite-vec vec0). Подключается при AGENTIK_MEMORY_BACKEND=hybrid.
implementation(project(":memory-md-vector"))
} }
// SQLite-бэкенд (ksqlite JNI) — JVM-only по факту. // SQLite-бэкенд (ksqlite JNI) — JVM-only по факту.
@@ -166,3 +169,42 @@ val shadowJarTask = tasks.register<ShadowJar>("shadowJar") {
includeEmptyDirs = false includeEmptyDirs = false
} }
// :standalone deploy
//
// `deployJar` собирает fatjar через shadowJar и копирует его на приватный
// домашний хост (root@192.168.76.166). Репа приватная, IP захардкожен —
// пользователь явно сказал "можно в наглую указать конкретно этот ip".
// Если понадобится другой хост — переопредели через -PagentikDeployHost=...
//
// Конвенция имени файла: agentik-standalone-{version}-all.jar (как на хосте
// уже лежат agentik-cli-0.1.0-all.jar / agentik-0.1.0-all.jar —
// выровнено по стилю существующего деплоя).
//
// SSH — ключевой, пароль не спрашивает. Если упадёт — проверь
// `ssh root@192.168.76.166 'echo OK'` вручную.
val agentikDeployHost: String =
(project.findProperty("agentikDeployHost") as? String) ?: "root@192.168.76.166"
val agentikDeployPath: String =
(project.findProperty("agentikDeployPath") as? String) ?: "~/agentik-standalone-${project.version}-all.jar"
tasks.register("deployJar") {
group = "deployment"
description = "Builds shadowJar and copies it to $agentikDeployHost:$agentikDeployPath"
dependsOn(shadowJarTask)
doLast {
val jar = shadowJarTask.get().outputs.files.singleFile
val target = "$agentikDeployHost:$agentikDeployPath"
println("→ scp $jar $target")
val proc = ProcessBuilder("scp", jar.absolutePath, target)
.redirectErrorStream(true)
.start()
proc.inputStream.bufferedReader().forEachLine { println("scp> $it") }
val rc = proc.waitFor()
if (rc != 0) {
throw GradleException("scp failed (exit=$rc). Check ssh access to $agentikDeployHost.")
}
println("✓ deployed to $target")
}
}
@@ -15,6 +15,7 @@ import pw.binom.a2a.server.a2aAgent
import pw.binom.agentik.memory.MemoryReviewer import pw.binom.agentik.memory.MemoryReviewer
import pw.binom.agentik.memory.MemorySystem import pw.binom.agentik.memory.MemorySystem
import pw.binom.agentik.memory.md.openMdMemorySystem import pw.binom.agentik.memory.md.openMdMemorySystem
import pw.binom.agentik.memory.mdvector.openHybridMemorySystem
import pw.binom.agentik.memory.vector.VectorMemorySystem import pw.binom.agentik.memory.vector.VectorMemorySystem
import pw.binom.agentik.memory.vector.embedding.HttpEmbeddingClient import pw.binom.agentik.memory.vector.embedding.HttpEmbeddingClient
import pw.binom.agentik.memory.vector.embedding.SiglipEmbeddingProvider import pw.binom.agentik.memory.vector.embedding.SiglipEmbeddingProvider
@@ -256,7 +257,7 @@ private fun runServer() {
} }
} }
MemoryBackend.VECTOR -> { MemoryBackend.VECTOR -> {
val embedding: pw.binom.agentik.memory.vector.EmbeddingProvider = when (config.embedding.backend) { val embedding: pw.binom.agentik.memory.TextEmbeddingExecutor = when (config.embedding.backend) {
AppConfig.EmbeddingBackend.HTTP -> { AppConfig.EmbeddingBackend.HTTP -> {
val llm = config.llm val llm = config.llm
// Берём базовый URL + API key у активного LLM-бэкенда. // Берём базовый URL + API key у активного LLM-бэкенда.
@@ -271,7 +272,7 @@ private fun runServer() {
apiKey = oa.apiKey, apiKey = oa.apiKey,
model = config.embedding.model, model = config.embedding.model,
dimension = config.embedding.dimension, dimension = config.embedding.dimension,
) ).asExecutor()
} }
AppConfig.EmbeddingBackend.SIGLIP -> { AppConfig.EmbeddingBackend.SIGLIP -> {
val modelPath = checkNotNull(config.embedding.modelPath) { val modelPath = checkNotNull(config.embedding.modelPath) {
@@ -280,7 +281,7 @@ private fun runServer() {
val tokenizerPath = checkNotNull(config.embedding.tokenizerPath) { val tokenizerPath = checkNotNull(config.embedding.tokenizerPath) {
"AGENTIK_EMBEDDING_BACKEND=siglip требует AGENTIK_EMBEDDING_TOKENIZER_PATH" "AGENTIK_EMBEDDING_BACKEND=siglip требует AGENTIK_EMBEDDING_TOKENIZER_PATH"
} }
SiglipEmbeddingProvider(modelPath = modelPath, tokenizerPath = tokenizerPath) SiglipEmbeddingProvider(modelPath = modelPath, tokenizerPath = tokenizerPath).asExecutor()
} }
} }
VectorMemorySystem.open( VectorMemorySystem.open(
@@ -296,6 +297,44 @@ private fun runServer() {
println(" memory: db=${config.agent.dbPath} (vector-backend, $backendLabel)") println(" memory: db=${config.agent.dbPath} (vector-backend, $backendLabel)")
} }
} }
MemoryBackend.HYBRID -> {
// Гибрид: .md файлы (single source of truth) + ksqlite vector-кэш.
// Нужны embeddings (HTTP/SIGLIP) + путь к vector-БД (по умолчанию рядом с agentik.db).
val embedding: pw.binom.agentik.memory.TextEmbeddingExecutor = when (config.embedding.backend) {
AppConfig.EmbeddingBackend.HTTP -> {
val llm = config.llm
require(llm.backend == LlmBackend.OPENAI) {
"AGENTIK_EMBEDDING_BACKEND=http требует LLM_BACKEND=openai"
}
val oa = checkNotNull(llm.openai) { "openai config required" }
HttpEmbeddingClient(
apiUrl = oa.baseUrl.trimEnd('/'),
apiKey = oa.apiKey,
model = config.embedding.model,
dimension = config.embedding.dimension,
).asExecutor()
}
AppConfig.EmbeddingBackend.SIGLIP -> {
val modelPath = checkNotNull(config.embedding.modelPath) {
"AGENTIK_EMBEDDING_BACKEND=siglip требует AGENTIK_EMBEDDING_MODEL_PATH"
}
val tokenizerPath = checkNotNull(config.embedding.tokenizerPath) {
"AGENTIK_EMBEDDING_BACKEND=siglip требует AGENTIK_EMBEDDING_TOKENIZER_PATH"
}
SiglipEmbeddingProvider(modelPath = modelPath, tokenizerPath = tokenizerPath).asExecutor()
}
}
val mdDir = rawMemory.takeUnless { it.equals("off", true) } ?: defaultMemoryDir()
val vectorDbPath = deriveVectorDbPath(config.agent.dbPath)
openHybridMemorySystem(
memoryRoot = Path(mdDir),
vectorDbPath = Path(vectorDbPath),
dimension = embedding.dimension,
embedder = embedding,
).also {
println(" memory: md=$mdDir, vector-cache=$vectorDbPath (hybrid-backend, dim=${embedding.dimension})")
}
}
} }
// Контекстное окно модели (для compaction'а working memory). // Контекстное окно модели (для compaction'а working memory).
@@ -428,3 +467,12 @@ private fun defaultMemoryDir(): String {
val home = System.getProperty("user.home") ?: "." val home = System.getProperty("user.home") ?: "."
return "$home/.agentik/memory" return "$home/.agentik/memory"
} }
/**
* Производный путь к vector-БД гибридного бэкенда: рядом с `agentik.db`,
* но с суффиксом `-vectors`. Например: `agentik.db` → `agentik-vectors.db`.
*/
private fun deriveVectorDbPath(agentDbPath: String): String {
val withoutExt = agentDbPath.removeSuffix(".db")
return "$withoutExt-vectors.db"
}
@@ -115,9 +115,29 @@ data class AppConfig(
val endpoints: Boolean = false, val endpoints: Boolean = false,
) )
/** Бэкенд долговременной памяти. */ /** Бэкенд долговременной памяти.
*
* Env-алиасы (для удобства CLI):
* - `md` → [MD] (Hermes-style §-файлы)
* - `vector` → [VECTOR] (sqlite + JVector + LLM-embeddings)
* - `md-vector` → [HYBRID] (md-файлы = source of truth + sqlite vector-кэш)
* - `off` → [OFF] (память выключена)
*/
@Serializable @Serializable
enum class MemoryBackend { MD, VECTOR, OFF } enum class MemoryBackend {
MD, VECTOR, HYBRID, OFF;
companion object {
/** Маппинг env-строки (lowercase, как приходит из shell) → enum. */
fun fromEnvString(s: String): MemoryBackend? = when (s.lowercase()) {
"md" -> MD
"vector" -> VECTOR
"md-vector", "hybrid" -> HYBRID
"off", "none" -> OFF
else -> null
}
}
}
/** Бэкенд эмбеддингов (для memory-backend=vector). */ /** Бэкенд эмбеддингов (для memory-backend=vector). */
@Serializable @Serializable
@@ -183,7 +203,7 @@ data class AppConfig(
memory = MemorySection( memory = MemorySection(
dir = env("AGENTIK_MEMORY_DIR")?.takeIf { it.isNotBlank() }, dir = env("AGENTIK_MEMORY_DIR")?.takeIf { it.isNotBlank() },
backend = env("AGENTIK_MEMORY_BACKEND")?.let { backend = env("AGENTIK_MEMORY_BACKEND")?.let {
runCatching { MemoryBackend.valueOf(it.uppercase()) }.getOrNull() MemoryBackend.fromEnvString(it)
} ?: MemoryBackend.MD, } ?: MemoryBackend.MD,
compressionThreshold = env("AGENTIK_COMPRESSION_THRESHOLD")?.toDoubleOrNull() compressionThreshold = env("AGENTIK_COMPRESSION_THRESHOLD")?.toDoubleOrNull()
?.coerceIn(0.1, 0.99) ?: DEFAULT_COMPRESSION_THRESHOLD, ?.coerceIn(0.1, 0.99) ?: DEFAULT_COMPRESSION_THRESHOLD,