memory: Curator (Phase 4) — фоновая архивация старых неиспользуемых заметок
* :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-флаги отложен.
This commit is contained in:
@@ -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<MemoryStoreEvent> = emptyFlow()
|
||||
|
||||
override fun close()
|
||||
|
||||
@@ -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)
|
||||
```
|
||||
|
||||
## Остановка
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
@@ -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<String, MemoryNote>()
|
||||
private val ev = MutableSharedFlow<MemoryStoreEvent>(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<MemoryNote> =
|
||||
notes.values
|
||||
.filter { category == null || it.category == category }
|
||||
.drop(offset)
|
||||
.take(limit)
|
||||
override suspend fun search(query: pw.binom.agentik.memory.MemorySearchQuery): List<pw.binom.agentik.memory.MemorySearchResult> =
|
||||
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<MemoryStoreEvent> = ev.asSharedFlow()
|
||||
override fun close() {}
|
||||
fun list() = notes.values.toList()
|
||||
}
|
||||
|
||||
@@ -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<String, MemoryNote>()
|
||||
private val ev = MutableSharedFlow<MemoryStoreEvent>(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<MemoryNote> =
|
||||
notes.values
|
||||
.filter { category == null || it.category == category }
|
||||
.drop(offset)
|
||||
.take(limit)
|
||||
override suspend fun search(query: MemorySearchQuery): List<MemorySearchResult> =
|
||||
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<MemoryStoreEvent> = ev.asSharedFlow()
|
||||
override fun close() {}
|
||||
fun allIds(): List<String> = notes.keys.sorted()
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user