Compare commits
1 Commits
7
...
68543357c2
| Author | SHA1 | Date | |
|---|---|---|---|
| 68543357c2 |
@@ -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). ---
|
||||||
|
|||||||
@@ -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)
|
||||||
|
}
|
||||||
|
|||||||
@@ -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)
|
|
||||||
}
|
|
||||||
|
|||||||
-39
@@ -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
|
|
||||||
}
|
|
||||||
}
|
|
||||||
+26
-55
@@ -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
|
|
||||||
}
|
|
||||||
|
|||||||
+8
-6
@@ -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] отрежет лишних.
|
||||||
|
|||||||
+7
-5
@@ -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,
|
||||||
),
|
),
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|||||||
+15
-11
@@ -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) {
|
||||||
|
|||||||
+22
-26
@@ -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
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
+4
-4
@@ -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" })
|
||||||
|
|||||||
+11
-5
@@ -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() }
|
||||||
|
|||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
+160
@@ -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
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
+228
@@ -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")
|
||||||
@@ -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")
|
||||||
|
|||||||
@@ -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,
|
||||||
|
|||||||
Reference in New Issue
Block a user