From 9e12b22e8505d119a1ad1ce82aa14b9600929216 Mon Sep 17 00:00:00 2001 From: subochev Date: Tue, 15 Sep 2026 04:14:03 +0300 Subject: [PATCH] =?UTF-8?q?memory:=20Curator=20(Phase=204)=20=E2=80=94=20?= =?UTF-8?q?=D1=84=D0=BE=D0=BD=D0=BE=D0=B2=D0=B0=D1=8F=20=D0=B0=D1=80=D1=85?= =?UTF-8?q?=D0=B8=D0=B2=D0=B0=D1=86=D0=B8=D1=8F=20=D1=81=D1=82=D0=B0=D1=80?= =?UTF-8?q?=D1=8B=D1=85=20=D0=BD=D0=B5=D0=B8=D1=81=D0=BF=D0=BE=D0=BB=D1=8C?= =?UTF-8?q?=D0=B7=D1=83=D0=B5=D0=BC=D1=8B=D1=85=20=D0=B7=D0=B0=D0=BC=D0=B5?= =?UTF-8?q?=D1=82=D0=BE=D0=BA?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * :memory-api — MemoryStore.archiveStale(maxAge, maxUseCount, now): default имплементация через list + delete (бэкенды могут переопределить). * :standalone — Curator: фоновая корутина на Dispatchers.IO, раз в сутки дёргает archiveStale(90d, 0). start()/stop(), runPass() — однократный прогон для тестов. * :standalone/Main — Curator стартует автоматически если memory включён, stop() в shutdown hook. * :standalone — публичный TestInMemoryMemoryStore вынесен из CompactionTest, переиспользуется в CuratorTest. * README — раздел "Куратор памяти" с описанием семантики для md и vector бэкендов. Tests: 288 total (+4 CuratorTest). Fatjar smoke-tested, curator стартует на app boot, выводит `[Curator] started` в лог. Defaults: interval=1d, maxAge=90d, maxUseCount=0. Override через новые config-флаги отложен. --- .../pw/binom/agentik/memory/MemoryStore.kt | 32 +++++++ standalone/README.md | 14 +++ .../pw/binom/agentik/standalone/Main.kt | 6 ++ .../standalone/agent/memory/Curator.kt | 70 +++++++++++++++ .../standalone/agent/CompactionTest.kt | 40 +-------- .../standalone/agent/TestMemoryStores.kt | 55 ++++++++++++ .../standalone/agent/memory/CuratorTest.kt | 88 +++++++++++++++++++ 7 files changed, 267 insertions(+), 38 deletions(-) create mode 100644 standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/memory/Curator.kt create mode 100644 standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/TestMemoryStores.kt create mode 100644 standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/memory/CuratorTest.kt diff --git a/memory-api/src/commonMain/kotlin/pw/binom/agentik/memory/MemoryStore.kt b/memory-api/src/commonMain/kotlin/pw/binom/agentik/memory/MemoryStore.kt index b27510c..8787aeb 100644 --- a/memory-api/src/commonMain/kotlin/pw/binom/agentik/memory/MemoryStore.kt +++ b/memory-api/src/commonMain/kotlin/pw/binom/agentik/memory/MemoryStore.kt @@ -40,6 +40,38 @@ interface MemoryStore : AutoCloseable { suspend fun delete(id: String): Boolean suspend fun markUsed(id: String, at: Instant = Clock.System.now()) + /** + * Архивирует заметки, которые: + * - не использовались дольше [maxAge] (считая от `now`); + * - имеют `useCount <= [maxUseCount]` (по умолчанию 0, т.е. только никогда + * не выданные в prefetch). + * + * Семантика архивации зависит от бэкенда: + * - `:memory-md` — переименовывает -файл с суффиксом `.archived.{ts}`; + * - `:memory-vector` — удаляет из SQLite и JVector (там данные + * пересоздаются из `agentik.db` при старте). + * + * Default-имплементация использует [list] + [delete]; бэкенды могут + * переопределить для более чистой семантики (особенно MD). + * + * @return количество архивированных заметок. + */ + suspend fun archiveStale( + maxAge: kotlin.time.Duration, + maxUseCount: Int = 0, + now: Instant = Clock.System.now(), + ): Int { + val all = list(limit = Int.MAX_VALUE) + val cutoff = now - maxAge + var archived = 0 + for (n in all) { + if (n.lastUsedAt < cutoff && n.useCount <= maxUseCount) { + if (delete(n.id)) archived++ + } + } + return archived + } + fun events(): Flow = emptyFlow() override fun close() diff --git a/standalone/README.md b/standalone/README.md index d2e1872..4133d6f 100644 --- a/standalone/README.md +++ b/standalone/README.md @@ -175,6 +175,19 @@ export AGENTIK_COMPRESSION_THRESHOLD=0.7 # сжимаем раньше нового compaction не запускается (защита от зацикливания). Решение — поднять `OPENAI_CONTEXT_WINDOW` или понизить threshold. +## Куратор памяти (Curator) + +Фоновая корутина (запускается автоматически, если `AGENTIK_MEMORY_DIR != off`): +раз в сутки архивирует заметки, которые **не выдавались в prefetch дольше 90 +дней** и **имеют `useCount == 0`**. Семантика архивации зависит от бэкенда — +`:memory-md` переименовывает §-файл в `.archived.{ts}`, `:memory-vector` +удаляет из SQLite и JVector. + +Параметры пока захардкожены в `Curator.DEFAULT_INTERVAL` и +`Curator.DEFAULT_MAX_AGE` (1 день и 90 дней); для override нужен новый +config-флаг. На каждом проходе выводится `[Curator] archived N stale notes`, +если N > 0. + ## Память `AGENTIK_MEMORY_DIR` указывает на каталог, в котором лежат три §-файла: @@ -317,6 +330,7 @@ agentik standalone listening on http://localhost:8080 soul: /etc/agentik/SOUL.md (842 chars) memory: /var/lib/agentik/memory (md-backend) compaction: enabled, threshold=0.8, window=128000 tokens + curator: enabled (interval=1d, maxAge=90d) ``` ## Остановка diff --git a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/Main.kt b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/Main.kt index 0968294..c76285b 100644 --- a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/Main.kt +++ b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/Main.kt @@ -168,6 +168,12 @@ fun main() { } else { println(" compaction: disabled (OPENAI_CONTEXT_WINDOW not set)") } + if (memorySystem != null) { + val curator = pw.binom.agentik.standalone.agent.memory.Curator(memorySystem.store) + curator.start() + Runtime.getRuntime().addShutdownHook(Thread { curator.stop() }) + println(" curator: enabled (interval=${pw.binom.agentik.standalone.agent.memory.Curator.DEFAULT_INTERVAL}, maxAge=${pw.binom.agentik.standalone.agent.memory.Curator.DEFAULT_MAX_AGE})") + } Runtime.getRuntime().addShutdownHook(Thread { agent.close() mcpRegistry.close() diff --git a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/memory/Curator.kt b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/memory/Curator.kt new file mode 100644 index 0000000..a9ef444 --- /dev/null +++ b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/memory/Curator.kt @@ -0,0 +1,70 @@ +package pw.binom.agentik.standalone.agent.memory + +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.Job +import kotlinx.coroutines.delay +import kotlinx.coroutines.isActive +import kotlinx.coroutines.launch +import pw.binom.agentik.memory.MemoryStore +import kotlin.time.Clock +import kotlin.time.Duration + +/** + * Периодическая фоновая архивация заметок долговременной памяти. + * + * Семантика: раз в [interval] (по умолчанию сутки) проходим по всем заметкам и + * архивируем те, что: + * - не выдавались в prefetch дольше [maxAge] (по умолчанию 90 дней); + * - имеют `useCount <= [maxUseCount]` (по умолчанию 0). + * + * Сама политика архивации (file-rename / SQLite-delete) живёт в + * `MemoryStore.archiveStale(...)`; этот класс только дёргает её по таймеру. + * + * Запускается из [start] — фоновая корутина на [Dispatchers.IO]. Остановить — + * [stop] или закрыть скоуп. + */ +class Curator( + private val store: MemoryStore, + private val interval: Duration = DEFAULT_INTERVAL, + private val maxAge: Duration = DEFAULT_MAX_AGE, + private val maxUseCount: Int = 0, + private val clock: Clock = Clock.System, + private val scope: CoroutineScope = CoroutineScope(Dispatchers.IO), +) { + private var job: Job? = null + + fun start() { + if (job != null) return + job = scope.launch { + println("[Curator] started (interval=$interval, maxAge=$maxAge, maxUseCount=$maxUseCount)") + while (isActive) { + runPass() + delay(interval) + } + } + } + + /** Однократный проход — полезно для тестов и для немедленного прохода после старта. */ + suspend fun runPass(): Int { + val archived = store.archiveStale( + maxAge = maxAge, + maxUseCount = maxUseCount, + now = clock.now(), + ) + if (archived > 0) { + println("[Curator] archived $archived stale notes (maxAge=$maxAge, maxUseCount=$maxUseCount)") + } + return archived + } + + fun stop() { + job?.cancel() + job = null + } + + companion object { + val DEFAULT_INTERVAL: Duration = Duration.parse("24h") + val DEFAULT_MAX_AGE: Duration = Duration.parse("90d") + } +} diff --git a/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/CompactionTest.kt b/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/CompactionTest.kt index c107738..1ede88a 100644 --- a/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/CompactionTest.kt +++ b/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/CompactionTest.kt @@ -128,7 +128,7 @@ class CompactionTest { @Test fun `compaction calls memoryReviewer reviewPreCompaction`() = runTest { fakeLlm.reply = "ok" - val memStore = InMemoryMemoryStore() + val memStore = TestInMemoryMemoryStore() val compactor = RecordingCompactor("compacted summary") val agent = newAgent( contextWindow = 30, @@ -178,41 +178,5 @@ private class RecordingCompactor(private val result: String) : ContextCompactor } /** - * Простой in-memory MemoryStore для тестов compaction'а (триггер памяти). + * (InMemoryMemoryStore вынесен в [TestInMemoryMemoryStore].) */ -private class InMemoryMemoryStore : MemoryStore { - private val notes = mutableMapOf() - private val ev = MutableSharedFlow(extraBufferCapacity = 16) - - override suspend fun upsert(note: MemoryNote) { - notes[note.id] = note - ev.tryEmit(MemoryStoreEvent.Upserted(note)) - } - override suspend fun get(id: String): MemoryNote? = notes[id] - override suspend fun list( - category: MemoryCategory?, - conversationId: String?, - limit: Int, - offset: Int, - ): List = - notes.values - .filter { category == null || it.category == category } - .drop(offset) - .take(limit) - override suspend fun search(query: pw.binom.agentik.memory.MemorySearchQuery): List = - notes.values - .filter { query.category == null || it.category == query.category } - .map { pw.binom.agentik.memory.MemorySearchResult(it, 1.0f) } - .take(query.topK) - override suspend fun delete(id: String): Boolean { - val ok = notes.remove(id) != null - if (ok) ev.tryEmit(MemoryStoreEvent.Deleted(id)) - return ok - } - override suspend fun markUsed(id: String, at: Instant) { - notes[id]?.let { notes[id] = it.copy(lastUsedAt = at, useCount = it.useCount + 1) } - } - override fun events(): Flow = ev.asSharedFlow() - override fun close() {} - fun list() = notes.values.toList() -} diff --git a/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/TestMemoryStores.kt b/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/TestMemoryStores.kt new file mode 100644 index 0000000..97c83d4 --- /dev/null +++ b/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/TestMemoryStores.kt @@ -0,0 +1,55 @@ +package pw.binom.agentik.standalone.agent + +import kotlinx.coroutines.flow.Flow +import kotlinx.coroutines.flow.MutableSharedFlow +import kotlinx.coroutines.flow.asSharedFlow +import pw.binom.agentik.memory.MemoryCategory +import pw.binom.agentik.memory.MemoryNote +import pw.binom.agentik.memory.MemorySearchQuery +import pw.binom.agentik.memory.MemorySearchResult +import pw.binom.agentik.memory.MemoryStore +import pw.binom.agentik.memory.MemoryStoreEvent +import kotlin.time.Instant + +/** + * Простой in-memory MemoryStore для тестов. + * + * Не зависит от :memory-md / :memory-vector; идёт через `archiveStale` из + * default-имплементации MemoryStore. + */ +class TestInMemoryMemoryStore : MemoryStore { + private val notes = mutableMapOf() + private val ev = MutableSharedFlow(extraBufferCapacity = 16) + + override suspend fun upsert(note: MemoryNote) { + notes[note.id] = note + ev.tryEmit(MemoryStoreEvent.Upserted(note)) + } + override suspend fun get(id: String): MemoryNote? = notes[id] + override suspend fun list( + category: MemoryCategory?, + conversationId: String?, + limit: Int, + offset: Int, + ): List = + notes.values + .filter { category == null || it.category == category } + .drop(offset) + .take(limit) + override suspend fun search(query: MemorySearchQuery): List = + notes.values + .filter { query.category == null || it.category == query.category } + .map { MemorySearchResult(it, 1.0f) } + .take(query.topK) + override suspend fun delete(id: String): Boolean { + val ok = notes.remove(id) != null + if (ok) ev.tryEmit(MemoryStoreEvent.Deleted(id)) + return ok + } + override suspend fun markUsed(id: String, at: Instant) { + notes[id]?.let { notes[id] = it.copy(lastUsedAt = at, useCount = it.useCount + 1) } + } + override fun events(): Flow = ev.asSharedFlow() + override fun close() {} + fun allIds(): List = notes.keys.sorted() +} diff --git a/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/memory/CuratorTest.kt b/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/memory/CuratorTest.kt new file mode 100644 index 0000000..839b85f --- /dev/null +++ b/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/memory/CuratorTest.kt @@ -0,0 +1,88 @@ +package pw.binom.agentik.standalone.agent.memory + +import kotlinx.coroutines.test.runTest +import pw.binom.agentik.standalone.agent.TestInMemoryMemoryStore +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 pw.binom.agentik.memory.MemoryCategory +import pw.binom.agentik.memory.MemoryNote +import pw.binom.agentik.memory.MemorySource + +class CuratorTest { + + @Test + fun `archives notes older than maxAge with zero useCount`() = runTest { + val store = TestInMemoryMemoryStore() + val now = Instant.parse("2026-09-15T00:00:00Z") + store.upsert(note("old", lastUsedOffsetDays = 100, useCount = 0)) + store.upsert(note("fresh", lastUsedOffsetDays = 1, useCount = 0)) + val curator = Curator(store, maxAge = Duration.parse("90d"), clock = FakeClock(now)) + assertEquals(1, curator.runPass()) + assertEquals(listOf("fresh"), store.allIds()) + } + + @Test + fun `keeps notes that were used recently even if old`() = runTest { + val store = TestInMemoryMemoryStore() + val now = Instant.parse("2026-09-15T00:00:00Z") + store.upsert(note("frequently-used", lastUsedOffsetDays = 1, useCount = 50)) + val curator = Curator(store, maxAge = Duration.parse("90d"), clock = FakeClock(now)) + assertEquals(0, curator.runPass()) + assertEquals(listOf("frequently-used"), store.allIds()) + } + + @Test + fun `archives notes with useCount above default when configured`() = runTest { + val store = TestInMemoryMemoryStore() + val now = Instant.parse("2026-09-15T00:00:00Z") + store.upsert(note("twice-used", lastUsedOffsetDays = 100, useCount = 2)) + val curator = Curator( + store, + maxAge = Duration.parse("90d"), + maxUseCount = 5, + clock = FakeClock(now), + ) + assertEquals(1, curator.runPass()) + assertTrue(store.allIds().isEmpty()) + } + + @Test + fun `start and stop launch and cancel the background loop`() = runTest { + val store = TestInMemoryMemoryStore() + val curator = Curator( + store, + interval = Duration.parse("10ms"), + maxAge = Duration.parse("90d"), + ) + curator.start() + Thread.sleep(50) + curator.stop() + // ничего не падает, корутина отменена + } + + private fun note( + id: String, + lastUsedOffsetDays: Long, + useCount: Int, + ): MemoryNote { + val now = Instant.parse("2026-09-15T00:00:00Z") + return MemoryNote( + id = id, + category = MemoryCategory.WORLD, + content = "fact $id", + createdAt = now - Duration.parse("${lastUsedOffsetDays}d"), + lastUsedAt = now - Duration.parse("${lastUsedOffsetDays}d"), + useCount = useCount, + conversationId = null, + source = MemorySource.AGENT_SAVE, + ) + } + + private class FakeClock(private val now: Instant) : Clock { + override fun now(): Instant = now + } +}