From 68543357c239c26c06f69f7c825632aa8f50dff9 Mon Sep 17 00:00:00 2001 From: subochev Date: Mon, 21 Sep 2026 12:28:09 +0300 Subject: [PATCH] feat(memory): migrate `EmbeddingProvider` to KMP-compatible `TextEmbeddingExecutor`, add `:memory-md-vector`, and hybrid backend support - 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 --- gradle/libs.versions.toml | 9 +- memory-api/build.gradle.kts | 7 + .../pw/binom/agentik/memory/MemoryNote.kt | 9 +- memory-vector/build.gradle.kts | 11 +- .../memory/vector/EmbeddingProvider.kt | 39 --- .../memory/vector/MemoryVectorIndex.kt | 81 ++----- .../memory/vector/VectorMemoryStore.kt | 14 +- .../memory/vector/VectorMemorySystem.kt | 12 +- .../vector/embedding/HttpEmbeddingClient.kt | 26 +- .../embedding/SiglipEmbeddingProvider.kt | 48 ++-- .../memory/vector/VectorMemoryStoreTest.kt | 8 +- .../embedding/SiglipEmbeddingProviderTest.kt | 16 +- outbox-inmemory/build.gradle.kts | 29 +++ .../outbox/inmemory/InMemoryOutboxStore.kt | 160 ++++++++++++ .../inmemory/InMemoryOutboxStoreTest.kt | 228 ++++++++++++++++++ settings.gradle.kts | 4 + standalone/build.gradle.kts | 42 ++++ .../pw/binom/agentik/standalone/Main.kt | 54 ++++- .../agentik/standalone/config/AppConfig.kt | 26 +- 19 files changed, 658 insertions(+), 165 deletions(-) delete mode 100644 memory-vector/src/commonMain/kotlin/pw/binom/agentik/memory/vector/EmbeddingProvider.kt create mode 100644 outbox-inmemory/build.gradle.kts create mode 100644 outbox-inmemory/src/commonMain/kotlin/pw/binom/agentik/outbox/inmemory/InMemoryOutboxStore.kt create mode 100644 outbox-inmemory/src/commonTest/kotlin/pw/binom/agentik/outbox/inmemory/InMemoryOutboxStoreTest.kt diff --git a/gradle/libs.versions.toml b/gradle/libs.versions.toml index bb3dab7..a1c6f3a 100644 --- a/gradle/libs.versions.toml +++ b/gradle/libs.versions.toml @@ -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" } # --- text-embedding-kmp (pw.binom.ai.embeddingtext) — on-device SigLIP2 эмбеддинг через ONNX. --- -# Артефакты публикуются под именами `-jvm` (KMP convention для JVM-таргета). -text-embedding-api = { module = "pw.binom.ai.embeddingtext:api-jvm", version.ref = "text-embedding-kmp" } +# `api` — KMP с jvm + android + linuxX64/Arm64 + macos + ios + mingwX64 +# (с 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" } # --- Логирование: kotlin-logging (тонкая обёртка над slf4j-api) + logback-classic (binding). --- diff --git a/memory-api/build.gradle.kts b/memory-api/build.gradle.kts index 32391d3..45cc798 100644 --- a/memory-api/build.gradle.kts +++ b/memory-api/build.gradle.kts @@ -21,6 +21,13 @@ kotlin { sourceSets { commonMain.dependencies { 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 { implementation(kotlin("test")) diff --git a/memory-api/src/commonMain/kotlin/pw/binom/agentik/memory/MemoryNote.kt b/memory-api/src/commonMain/kotlin/pw/binom/agentik/memory/MemoryNote.kt index 945601d..856e1d4 100644 --- a/memory-api/src/commonMain/kotlin/pw/binom/agentik/memory/MemoryNote.kt +++ b/memory-api/src/commonMain/kotlin/pw/binom/agentik/memory/MemoryNote.kt @@ -24,4 +24,11 @@ data class MemoryNote( val useCount: Int = 0, val conversationId: String? = null, 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) +} diff --git a/memory-vector/build.gradle.kts b/memory-vector/build.gradle.kts index 154302f..b84f9af 100644 --- a/memory-vector/build.gradle.kts +++ b/memory-vector/build.gradle.kts @@ -25,6 +25,10 @@ kotlin { api(project(":memory-api")) implementation(libs.kotlinx.coroutines.core) 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 { implementation(kotlin("test")) @@ -34,14 +38,11 @@ kotlin { jvmMain.dependencies { implementation(libs.jvector) implementation(libs.sqldelight.sqlite.driver) + // Конкретная реализация TextEmbeddingExtractor поверх ONNX. + implementation(libs.text.embedding.siglip) } jvmTest.dependencies { implementation(kotlin("test")) } } } - -dependencies { - add("jvmMainApi", libs.text.embedding.api) - add("jvmMainImplementation", libs.text.embedding.siglip) -} diff --git a/memory-vector/src/commonMain/kotlin/pw/binom/agentik/memory/vector/EmbeddingProvider.kt b/memory-vector/src/commonMain/kotlin/pw/binom/agentik/memory/vector/EmbeddingProvider.kt deleted file mode 100644 index afdaecf..0000000 --- a/memory-vector/src/commonMain/kotlin/pw/binom/agentik/memory/vector/EmbeddingProvider.kt +++ /dev/null @@ -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): List = - 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 - } -} diff --git a/memory-vector/src/commonMain/kotlin/pw/binom/agentik/memory/vector/MemoryVectorIndex.kt b/memory-vector/src/commonMain/kotlin/pw/binom/agentik/memory/vector/MemoryVectorIndex.kt index 5c50502..32ccb43 100644 --- a/memory-vector/src/commonMain/kotlin/pw/binom/agentik/memory/vector/MemoryVectorIndex.kt +++ b/memory-vector/src/commonMain/kotlin/pw/binom/agentik/memory/vector/MemoryVectorIndex.kt @@ -1,63 +1,34 @@ package pw.binom.agentik.memory.vector -import pw.binom.agentik.memory.MemoryCategory -import pw.binom.agentik.memory.MemoryNote +import pw.binom.agentik.memory.MemoryVectorIndex as KmpMemoryVectorIndex +import pw.binom.agentik.memory.ScoredVector as KmpScoredVector /** - * Результат одного hit'а vector-поиска: id заметки + cosine-similarity score в [0..1]. - * Чем ближе к 1.0, тем семантически ближе query к заметке. - */ -data class ScoredVector( - val id: String, - val score: Float, -) - -/** - * Контракт vector-индекса. Реализация отвечает за ANN-поиск top-K ближайших - * векторов к query. Метаданные заметок лежат в [MemoryStore] (SQLite для - * vector-бэкенда); индекс хранит только embedding'и + id-маппинг. + * JVM-only alias на KMP-контракт из `:memory-api`. Удалять нельзя — пока + * `:memory-vector` существует как JVM-only модуль с JVector-имплементацией, + * все его internal helper'ы продолжают импортировать `MemoryVectorIndex` из + * `pw.binom.agentik.memory.vector.*` (старое FQN). После удаления модуля — + * можно убрать этот файл и переименовать пакеты импортов. * - * Потокобезопасность: реализации обязаны быть безопасны для конкурентных - * read'ов. write'ы (add/remove) могут требовать внешней синхронизации — это - * инвариант JVector (его OnHeapGraphIndex не thread-safe для мутаций). + * Раньше жил прямо здесь (`MemoryVectorIndex` + `ScoredVector` в + * `:memory-vector/commonMain`), но переехал в `:memory-api` 2026-09-21 + * чтобы стать доступным из KMP-модуля `:memory-md-vector`. */ -interface MemoryVectorIndex : AutoCloseable { - /** Текущая размерность embeddings. Фиксируется при первом [add]. */ - val dimension: Int - /** Количество записей в индексе. */ - suspend fun size(): Long +@Deprecated( + message = "Переехал в :memory-api (KMP-доступный). Импортируйте из pw.binom.agentik.memory.", + replaceWith = ReplaceWith( + "MemoryVectorIndex", + "pw.binom.agentik.memory.MemoryVectorIndex", + ), +) +typealias MemoryVectorIndex = KmpMemoryVectorIndex - /** Добавить или заменить запись по [id]. [embedding] должен иметь длину [dimension]. */ - suspend fun add(id: String, embedding: FloatArray) - - /** Удалить запись по [id]. Возвращает true если запись была. */ - suspend fun remove(id: String): Boolean - - /** ANN-поиск: top-[k] ближайших к [query]. [filter] применяется к id (например, по категории). */ - suspend fun search( - query: FloatArray, - k: Int, - filter: (MemoryNote) -> Boolean = { true }, - ): List - - /** Принудительно переписать 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 -} +@Deprecated( + message = "Переехал в :memory-api (KMP-доступный). Импортируйте из pw.binom.agentik.memory.", + replaceWith = ReplaceWith( + "ScoredVector", + "pw.binom.agentik.memory.ScoredVector", + ), +) +typealias ScoredVector = KmpScoredVector diff --git a/memory-vector/src/commonMain/kotlin/pw/binom/agentik/memory/vector/VectorMemoryStore.kt b/memory-vector/src/commonMain/kotlin/pw/binom/agentik/memory/vector/VectorMemoryStore.kt index baceb88..a2a3af1 100644 --- a/memory-vector/src/commonMain/kotlin/pw/binom/agentik/memory/vector/VectorMemoryStore.kt +++ b/memory-vector/src/commonMain/kotlin/pw/binom/agentik/memory/vector/VectorMemoryStore.kt @@ -6,6 +6,8 @@ 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.noteMatches import kotlin.math.exp import kotlin.time.Clock import kotlin.time.Instant @@ -23,13 +25,13 @@ import kotlinx.coroutines.sync.withLock * и эмбеддинг; [delete] — и то и другое; [search] использует ANN для кандидатов, * потом re-rank по recency. * - * [embeddingProvider] обязателен — используется для эмбеддинга контента при + * [embedding] обязателен — используется для эмбеддинга контента при * upsert и query при search. Без него vector-бэкенд не имеет смысла. */ class VectorMemoryStore( private val index: MemoryVectorIndex, private val metaStore: MemoryMetaStore, - private val embeddingProvider: EmbeddingProvider, + private val embedding: TextEmbeddingExecutor, ) : MemoryStore { private val mutex = Mutex() @@ -37,9 +39,9 @@ class VectorMemoryStore( override fun events(): Flow = _events.asSharedFlow() override suspend fun upsert(note: MemoryNote) = mutex.withLock { - val embedding = embeddingProvider.embed(note.content) - metaStore.put(note, embedding) - index.add(note.id, embedding) + val vec = embedding.embed(note.content) + metaStore.put(note, vec) + index.add(note.id, vec) _events.emit(MemoryStoreEvent.Upserted(note)) } @@ -53,7 +55,7 @@ class VectorMemoryStore( ): List = metaStore.list(category, conversationId, limit, offset) override suspend fun search(query: MemorySearchQuery): List { - val queryEmbedding = embeddingProvider.embed(query.query) + val queryEmbedding = embedding.embed(query.query) val overFetch = (query.topK * 5).coerceAtLeast(query.topK) // Берём больше кандидатов, чем нужно — финальный фильтр по category/convId // через [metaStore.get] + [noteMatches] отрежет лишних. diff --git a/memory-vector/src/jvmMain/kotlin/pw/binom/agentik/memory/vector/VectorMemorySystem.kt b/memory-vector/src/jvmMain/kotlin/pw/binom/agentik/memory/vector/VectorMemorySystem.kt index 3fc2da2..9c173a6 100644 --- a/memory-vector/src/jvmMain/kotlin/pw/binom/agentik/memory/vector/VectorMemorySystem.kt +++ b/memory-vector/src/jvmMain/kotlin/pw/binom/agentik/memory/vector/VectorMemorySystem.kt @@ -9,6 +9,7 @@ import pw.binom.agentik.memory.MemorySearchQuery import pw.binom.agentik.memory.MemoryStore import pw.binom.agentik.memory.MemorySystem import pw.binom.agentik.memory.ReviewedTurn +import pw.binom.agentik.memory.TextEmbeddingExecutor /** * Бандл компонентов vector-бэкенда памяти — то же, что @@ -36,19 +37,20 @@ class VectorMemorySystem( * Открыть vector-бэкенд: SQLite + JVector + HTTP embedding client. * * @param dbPath путь к agentik.db (SQLite для metadata + embedding-blobs) - * @param embedding [EmbeddingProvider] — обычно HttpEmbeddingClient + * @param embedding [TextEmbeddingExecutor] — обычно HttpEmbeddingClient.asExecutor() * @param topK размер top-K для prefetch */ fun open( dbPath: String, - embedding: EmbeddingProvider, + embedding: TextEmbeddingExecutor, topK: Int = 10, ): VectorMemorySystem { - val metaStore = SqliteMemoryMetaStore.open(dbPath, embedding.dimension) + val dim = embedding.dimension + val metaStore = SqliteMemoryMetaStore.open(dbPath, dim) // Граф пересобирается из SQLite (источник правды): без seed'ов // после рестарта in-RAM индекс пуст и search возвращал бы [], // пока не появятся новые upsert'ы. - val index = JVectorMemoryIndex(embedding.dimension, metaStore.allEntries()) + val index = JVectorMemoryIndex(dim, metaStore.allEntries()) val store = VectorMemoryStore(index, metaStore, embedding) val prefetcher = VectorPrefetcher(store, topK) val reviewer = VectorMemoryReviewer(store) @@ -59,7 +61,7 @@ class VectorMemorySystem( closables = listOfNotNull( metaStore, index, - embedding as? AutoCloseable, + embedding, ), ) } diff --git a/memory-vector/src/jvmMain/kotlin/pw/binom/agentik/memory/vector/embedding/HttpEmbeddingClient.kt b/memory-vector/src/jvmMain/kotlin/pw/binom/agentik/memory/vector/embedding/HttpEmbeddingClient.kt index 52d81da..dc8e37d 100644 --- a/memory-vector/src/jvmMain/kotlin/pw/binom/agentik/memory/vector/embedding/HttpEmbeddingClient.kt +++ b/memory-vector/src/jvmMain/kotlin/pw/binom/agentik/memory/vector/embedding/HttpEmbeddingClient.kt @@ -5,25 +5,26 @@ import java.net.http.HttpClient import java.net.http.HttpRequest import java.net.http.HttpResponse import java.time.Duration -import java.util.concurrent.ConcurrentHashMap -import kotlinx.serialization.Serializable import kotlinx.serialization.json.Json import kotlinx.serialization.json.JsonElement -import kotlinx.serialization.json.JsonObject import kotlinx.serialization.json.JsonPrimitive import kotlinx.serialization.json.buildJsonObject import kotlinx.serialization.json.jsonArray import kotlinx.serialization.json.jsonObject import kotlinx.serialization.json.jsonPrimitive 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. * Используется при memory-backend=vector. * - * LRU-кэш на [cacheSize] текстов (default 256) — дедупликация запросов - * к API на одинаковых промптах. + * Реализует [TextEmbeddingExtractor] (из text-embedding-kmp:api) + оборачивается + * в [TextEmbeddingExecutor] через [asExecutor] для совместимости с + * VectorMemoryStore. LRU-кэш на [cacheSize] текстов (default 256) — дедупликация + * запросов к API на одинаковых промптах. * * @param apiUrl базовый URL (без trailing slash), например `https://api.openai.com` * @param apiKey bearer-токен @@ -35,9 +36,9 @@ class HttpEmbeddingClient( private val apiUrl: String, private val apiKey: String, private val model: String, - override val dimension: Int, + private val dimension: Int, cacheSize: Int = 256, -) : EmbeddingProvider, AutoCloseable { +) : TextEmbeddingExtractor { private val cache = LruCache(cacheSize) private val http: HttpClient = HttpClient.newBuilder() @@ -45,11 +46,11 @@ class HttpEmbeddingClient( .build() private val json = Json { ignoreUnknownKeys = true } - override suspend fun embed(text: String): FloatArray { - cache.get(text)?.let { return it } + override fun embed(text: String): TextEmbedding { + cache.get(text)?.let { return TextEmbedding(it) } val vector = fetchEmbedding(text) cache.put(text, vector) - return vector + return TextEmbedding(vector) } private fun fetchEmbedding(text: String): FloatArray { @@ -83,6 +84,9 @@ class HttpEmbeddingClient( } override fun close() = http.close() + + /** Оборачивает в [TextEmbeddingExecutor] с пред-объявленной размерностью. */ + fun asExecutor(): TextEmbeddingExecutor = TextEmbeddingExecutor(this, knownDimension = dimension) } private class LruCache(private val capacity: Int) { diff --git a/memory-vector/src/jvmMain/kotlin/pw/binom/agentik/memory/vector/embedding/SiglipEmbeddingProvider.kt b/memory-vector/src/jvmMain/kotlin/pw/binom/agentik/memory/vector/embedding/SiglipEmbeddingProvider.kt index 742912c..873b398 100644 --- a/memory-vector/src/jvmMain/kotlin/pw/binom/agentik/memory/vector/embedding/SiglipEmbeddingProvider.kt +++ b/memory-vector/src/jvmMain/kotlin/pw/binom/agentik/memory/vector/embedding/SiglipEmbeddingProvider.kt @@ -1,25 +1,22 @@ package pw.binom.agentik.memory.vector.embedding -import kotlinx.coroutines.Dispatchers -import kotlinx.coroutines.sync.Mutex -import kotlinx.coroutines.sync.withLock -import kotlinx.coroutines.withContext -import pw.binom.agentik.memory.vector.EmbeddingProvider +import kotlinx.coroutines.runBlocking +import pw.binom.agentik.memory.TextEmbeddingExecutor +import pw.binom.voice.embeddingtext.TextEmbedding import pw.binom.voice.embeddingtext.TextEmbeddingExtractor import pw.binom.voice.embeddingtext.createSiglip2TextExtractor /** * Локальный on-device эмбеддинг через [TextEmbeddingExtractor] (SigLIP2 / ONNX). * - * Особенности: - * - `TextEmbeddingExtractor.embed(text)` — **blocking** (ONNX-инференс на CPU), - * не suspend. Оборачиваем в `Dispatchers.IO` + `Mutex`, чтобы сериализовать - * доступ из нескольких корутин (ONNX-сессия не reentrant). - * - Размерность фиксирована extractor'ом (SigLIP2-base = 768); параметр - * `dimension` в конструкторе не принимаем — берём через [probeDimension]. - * - LRU-кэш из [HttpEmbeddingClient] не используем здесь: ONNX-инференс на - * CPU ≈ 5-15 мс, кэш полезен только для HTTP. Но если потребуется — - * легко добавить. + * Реализует [TextEmbeddingExtractor] напрямую (делегирует в + * `createSiglip2TextExtractor` из text-embedding-kmp:siglip) + оборачивается + * в [TextEmbeddingExecutor] через [asExecutor] для совместимости с + * VectorMemoryStore. Сиглизация через `Dispatchers.IO` теперь внутри + * `TextEmbeddingExecutor.embed` — раньше лежала здесь. + * + * Размерность фиксирована extractor'ом (SigLIP2-base = 768); передаём + * явно в [asExecutor]. * * Модель + токенизатор не бандлятся в jar: передаём пути в конструкторе. * Скачать: см. README репы `text-embedding-kmp`. @@ -27,23 +24,22 @@ import pw.binom.voice.embeddingtext.createSiglip2TextExtractor class SiglipEmbeddingProvider( modelPath: String, tokenizerPath: String, -) : EmbeddingProvider, AutoCloseable { +) : TextEmbeddingExtractor { - private val extractor: TextEmbeddingExtractor = + private val delegate: TextEmbeddingExtractor = createSiglip2TextExtractor(modelPath = modelPath, tokenizerPath = tokenizerPath) - override val dimension: Int = run { - val probe = extractor.embed("probe") - probe.dim - } + override fun embed(text: String): TextEmbedding = delegate.embed(text) - 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() { - extractor.close() + companion object { + const val SIGLIP2_DIM: Int = 768 } } diff --git a/memory-vector/src/jvmTest/kotlin/pw/binom/agentik/memory/vector/VectorMemoryStoreTest.kt b/memory-vector/src/jvmTest/kotlin/pw/binom/agentik/memory/vector/VectorMemoryStoreTest.kt index 3a0e75c..cc3acfd 100644 --- a/memory-vector/src/jvmTest/kotlin/pw/binom/agentik/memory/vector/VectorMemoryStoreTest.kt +++ b/memory-vector/src/jvmTest/kotlin/pw/binom/agentik/memory/vector/VectorMemoryStoreTest.kt @@ -32,7 +32,7 @@ class VectorMemoryStoreTest { // Загружаем начальные entries из metaStore (на случай если что-то там есть). val seedEntries = metaStore.allEntries() index = JVectorMemoryIndex(dimension = dim, seedEntries = seedEntries) - store = VectorMemoryStore(index, metaStore, FakeEmbeddingProvider(dimension = dim)) + store = VectorMemoryStore(index, metaStore, fakeEmbeddingExecutor(dimension = dim)) } @AfterTest @@ -127,7 +127,7 @@ class VectorMemoryStoreTest { val meta2 = SqliteMemoryMetaStore("jdbc:sqlite:${file.absolutePath}", dimension = dim) val seedEntries = meta2.allEntries() val idx2 = JVectorMemoryIndex(dimension = dim, seedEntries = seedEntries) - val store2 = VectorMemoryStore(idx2, meta2, FakeEmbeddingProvider(dimension = dim)) + val store2 = VectorMemoryStore(idx2, meta2, fakeEmbeddingExecutor(dimension = dim)) try { assertEquals(2L, idx2.size()) val results = store2.search(MemorySearchQuery(query = "persistent 1", topK = 5)) @@ -141,11 +141,11 @@ class VectorMemoryStoreTest { fun openSeedsIndexFromSqliteAfterRestart() = runTest { // Регрессия: VectorMemorySystem.open() обязан пересадить in-RAM граф // из 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.close() - val second = VectorMemorySystem.open(file.absolutePath, FakeEmbeddingProvider(dimension = dim)) + val second = VectorMemorySystem.open(file.absolutePath, fakeEmbeddingExecutor(dimension = dim)) try { val results = second.store.search(MemorySearchQuery(query = "restarted fact", topK = 5)) assertTrue(results.any { it.note.id == "r" }) diff --git a/memory-vector/src/jvmTest/kotlin/pw/binom/agentik/memory/vector/embedding/SiglipEmbeddingProviderTest.kt b/memory-vector/src/jvmTest/kotlin/pw/binom/agentik/memory/vector/embedding/SiglipEmbeddingProviderTest.kt index 90824e7..5fd27f6 100644 --- a/memory-vector/src/jvmTest/kotlin/pw/binom/agentik/memory/vector/embedding/SiglipEmbeddingProviderTest.kt +++ b/memory-vector/src/jvmTest/kotlin/pw/binom/agentik/memory/vector/embedding/SiglipEmbeddingProviderTest.kt @@ -17,17 +17,18 @@ import kotlin.test.assertTrue class SiglipEmbeddingProviderTest { @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") assume(modelDir.exists() && File(modelDir, "text_model_int8.onnx").exists()) { "SigLIP2 model files not found in /tmp/text-emb-model/ — skipping" } - SiglipEmbeddingProvider( + val provider = SiglipEmbeddingProvider( modelPath = "${modelDir.absolutePath}/text_model_int8.onnx", tokenizerPath = "${modelDir.absolutePath}/tokenizer.model", - ).use { provider -> + ).asExecutor() + provider.use { 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) assertTrue(v.any { it != 0f }, "embedding should not be all zeros") } @@ -41,7 +42,7 @@ class SiglipEmbeddingProviderTest { SiglipEmbeddingProvider( modelPath = nonExistent.absolutePath, tokenizerPath = nonExistent.absolutePath, - ).use { it.dimension } + ) } } @@ -49,3 +50,8 @@ class SiglipEmbeddingProviderTest { org.junit.Assume.assumeTrue(message(), condition) } } + +// runBlocking нужен потому что suspend-вызов provider.embed в suspend-тесте. +// Локальный импорт чтобы не тащить runBlocking в прод-код. +private fun runBlocking(block: suspend () -> T): T = + kotlinx.coroutines.runBlocking { block() } diff --git a/outbox-inmemory/build.gradle.kts b/outbox-inmemory/build.gradle.kts new file mode 100644 index 0000000..cd36ed7 --- /dev/null +++ b/outbox-inmemory/build.gradle.kts @@ -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) + } + } +} diff --git a/outbox-inmemory/src/commonMain/kotlin/pw/binom/agentik/outbox/inmemory/InMemoryOutboxStore.kt b/outbox-inmemory/src/commonMain/kotlin/pw/binom/agentik/outbox/inmemory/InMemoryOutboxStore.kt new file mode 100644 index 0000000..b6dc1d1 --- /dev/null +++ b/outbox-inmemory/src/commonMain/kotlin/pw/binom/agentik/outbox/inmemory/InMemoryOutboxStore.kt @@ -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() + private val liveFlow = MutableSharedFlow( + 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 = flow { + // Replay buffer — snapshot под mutex'ом, дальше iterate без lock'а. + // Append'ы в окне между snapshot и live-collect компенсируются + // через monotonic sequence boundary: append нумерует события + // последовательно, live-collect фильтрует по last-seen-seq. + val snapshot: List = 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 = 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 + } +} + diff --git a/outbox-inmemory/src/commonTest/kotlin/pw/binom/agentik/outbox/inmemory/InMemoryOutboxStoreTest.kt b/outbox-inmemory/src/commonTest/kotlin/pw/binom/agentik/outbox/inmemory/InMemoryOutboxStoreTest.kt new file mode 100644 index 0000000..6611bf4 --- /dev/null +++ b/outbox-inmemory/src/commonTest/kotlin/pw/binom/agentik/outbox/inmemory/InMemoryOutboxStoreTest.kt @@ -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() + val done = CompletableDeferred() + 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() + val done = CompletableDeferred() + 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() + .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() + 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") diff --git a/settings.gradle.kts b/settings.gradle.kts index 9fd16c6..7cf2450 100644 --- a/settings.gradle.kts +++ b/settings.gradle.kts @@ -91,6 +91,10 @@ include(":outbox-api") // eviction. Для тестов, dev-режима, embedded-сценариев (Android core). include(":outbox-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(":working-memory-api") include(":storage-inmemory") diff --git a/standalone/build.gradle.kts b/standalone/build.gradle.kts index b56d0fa..adbe1aa 100644 --- a/standalone/build.gradle.kts +++ b/standalone/build.gradle.kts @@ -57,6 +57,9 @@ kotlin { implementation(project(":memory-md")) if (!skipVectorMemory) { 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 по факту. @@ -166,3 +169,42 @@ val shadowJarTask = tasks.register("shadowJar") { 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") + } +} diff --git a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/Main.kt b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/Main.kt index e7f328b..63b67b6 100644 --- a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/Main.kt +++ b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/Main.kt @@ -15,6 +15,7 @@ import pw.binom.a2a.server.a2aAgent import pw.binom.agentik.memory.MemoryReviewer import pw.binom.agentik.memory.MemorySystem 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.embedding.HttpEmbeddingClient import pw.binom.agentik.memory.vector.embedding.SiglipEmbeddingProvider @@ -256,7 +257,7 @@ private fun runServer() { } } 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 -> { val llm = config.llm // Берём базовый URL + API key у активного LLM-бэкенда. @@ -271,7 +272,7 @@ private fun runServer() { apiKey = oa.apiKey, model = config.embedding.model, dimension = config.embedding.dimension, - ) + ).asExecutor() } AppConfig.EmbeddingBackend.SIGLIP -> { val modelPath = checkNotNull(config.embedding.modelPath) { @@ -280,7 +281,7 @@ private fun runServer() { val tokenizerPath = checkNotNull(config.embedding.tokenizerPath) { "AGENTIK_EMBEDDING_BACKEND=siglip требует AGENTIK_EMBEDDING_TOKENIZER_PATH" } - SiglipEmbeddingProvider(modelPath = modelPath, tokenizerPath = tokenizerPath) + SiglipEmbeddingProvider(modelPath = modelPath, tokenizerPath = tokenizerPath).asExecutor() } } VectorMemorySystem.open( @@ -296,6 +297,44 @@ private fun runServer() { 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). @@ -428,3 +467,12 @@ private fun defaultMemoryDir(): String { val home = System.getProperty("user.home") ?: "." 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" +} diff --git a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/config/AppConfig.kt b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/config/AppConfig.kt index 2f182c3..16bdb4b 100644 --- a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/config/AppConfig.kt +++ b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/config/AppConfig.kt @@ -115,9 +115,29 @@ data class AppConfig( 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 - 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). */ @Serializable @@ -183,7 +203,7 @@ data class AppConfig( memory = MemorySection( dir = env("AGENTIK_MEMORY_DIR")?.takeIf { it.isNotBlank() }, backend = env("AGENTIK_MEMORY_BACKEND")?.let { - runCatching { MemoryBackend.valueOf(it.uppercase()) }.getOrNull() + MemoryBackend.fromEnvString(it) } ?: MemoryBackend.MD, compressionThreshold = env("AGENTIK_COMPRESSION_THRESHOLD")?.toDoubleOrNull() ?.coerceIn(0.1, 0.99) ?: DEFAULT_COMPRESSION_THRESHOLD,