skills: SkillMiner — фоновое авто-создание скилов + debug-эндпоинты
SkillMiner (сетка безопасности skill self-improvement): каждые
AGENTIK_SKILL_MINING_INTERVAL user-ходов (default 15) LLM смотрит
последние AGENTIK_SKILL_MINING_MAX_TURNS ходы (default 30) + каталог
существующих скилов и возвращает structured JSON {"skills":[...]}.
Найденное upsert-ится в SkillStore — модель "забыла" вызвать
skill_save в ходе разговора, минер добирает её постфактум.
- SkillMiner.kt: короткий LiteConversation (one-shot), blocking-инференс
на Dispatchers.IO, defensive парсинг (кривой ответ -> пустой список).
- SkillMiningPrompts/SkillMiningParser: тот же подход, что
ReflectionParser (structured-output вместо tool-calling).
- ChatConversation.scheduleSkillMining() — хук после каждого хода
(рядом со scheduleReflection); ChatAgent/Main — прокидывание.
- DebugRoutes.kt: AGENTIK_DEBUG_ENDPOINTS=1 включает POST
/debug/reflect, /debug/skill-mine, /debug/curate, /debug/compact и
GET /debug/tokens для ручного триггерирования фоновых фич.
- ChatConversation.forceCompactNow(): принудительный compaction
без проверки порога (для /debug/compact).
- Тесты: SkillMinerTest (5) + SkillMiningParserTest (9); FakeLiteLlm
теперь записывает send()/sendContents() в lastContents.
179 jvm-тестов :standalone зелёные, README обновлён.
This commit is contained in:
@@ -55,6 +55,9 @@ java -jar standalone/build/libs/standalone-all.jar
|
|||||||
| `AGENTIK_COMPRESSION_THRESHOLD` | `0.8` | Доля лимита, при которой запускается compaction |
|
| `AGENTIK_COMPRESSION_THRESHOLD` | `0.8` | Доля лимита, при которой запускается compaction |
|
||||||
| `AGENTIK_REFLECTION_INTERVAL` | `10` | Self-reflection: каждый N-й пользовательский ход агент оценивает себя (LiteLlm) и сохраняет рефлексию. `0` = выключено. |
|
| `AGENTIK_REFLECTION_INTERVAL` | `10` | Self-reflection: каждый N-й пользовательский ход агент оценивает себя (LiteLlm) и сохраняет рефлексию. `0` = выключено. |
|
||||||
| `AGENTIK_REFLECTION_TOP_K` | `3` | Сколько последних рефлексий подмешивать в system prompt как «слабые места». `0` = не подмешивать. |
|
| `AGENTIK_REFLECTION_TOP_K` | `3` | Сколько последних рефлексий подмешивать в system prompt как «слабые места». `0` = не подмешивать. |
|
||||||
|
| `AGENTIK_SKILL_MINING_INTERVAL` | `15` | Skill mining: через сколько user-ходов запускать фоновый прогон SkillMiner. `0` = выключено. |
|
||||||
|
| `AGENTIK_SKILL_MINING_MAX_TURNS` | `30` | Сколько последних ходов передавать SkillMiner'у за один прогон. |
|
||||||
|
| `AGENTIK_DEBUG_ENDPOINTS` | `0` | `1` включает debug-эндпоинты (`/debug/reflect`, `/debug/skill-mine`, `/debug/curate`, `/debug/compact`, `/debug/tokens`) |
|
||||||
|
|
||||||
### OpenAI backend
|
### OpenAI backend
|
||||||
|
|
||||||
@@ -354,10 +357,45 @@ AGENTIK_SKILLS_DIR=/etc/agentik/skills
|
|||||||
`skill_save`/`skill_delete` не требуют рестарта агента — изменения видны
|
`skill_save`/`skill_delete` не требуют рестарта агента — изменения видны
|
||||||
на ближайшем вызове `read_skill` (включая в этом же диалоге).
|
на ближайшем вызове `read_skill` (включая в этом же диалоге).
|
||||||
|
|
||||||
|
### Skill mining (автонавыки)
|
||||||
|
|
||||||
|
Модель может "протупить" и не вызвать `skill_save`, хотя приём был
|
||||||
|
переиспользуемым. Сетка безопасности — фоновый [SkillMiner]: каждые
|
||||||
|
`AGENTIK_SKILL_MINING_INTERVAL` пользовательских ходов (default 15, `0` =
|
||||||
|
выключено) LLM смотрит последние `AGENTIK_SKILL_MINING_MAX_TURNS` ходов
|
||||||
|
(default 30) + каталог существующих скилов и возвращает structured JSON
|
||||||
|
`{"skills": [{name, description, body}]}`. Найденные скилы upsert-ятся в
|
||||||
|
`AGENTIK_SKILLS_DIR` — агент становится умнее между сессиями. Обновления
|
||||||
|
существующих скилов (то же имя) поддерживаются, дубли — нет.
|
||||||
|
|
||||||
|
| Переменная | Default | Что делает |
|
||||||
|
|---|---|---|
|
||||||
|
| `AGENTIK_SKILL_MINING_INTERVAL` | `15` | Через сколько user-ходов запускать mining. `0` — выкл. |
|
||||||
|
| `AGENTIK_SKILL_MINING_MAX_TURNS` | `30` | Сколько последних ходов показывать минеру |
|
||||||
|
|
||||||
|
### Debug-эндпоинты
|
||||||
|
|
||||||
|
`AGENTIK_DEBUG_ENDPOINTS=1` включает эндпоинты для ручного триггерирования
|
||||||
|
фоновых фич (не ждать интервалов). Только локальная отладка: без
|
||||||
|
авторизации, в проде не включать.
|
||||||
|
|
||||||
|
| Эндпоинт | Действие |
|
||||||
|
|---|---|
|
||||||
|
| `POST /debug/reflect?conversationId=...` | прогон LlmReflector прямо сейчас, результат в БД |
|
||||||
|
| `POST /debug/skill-mine?conversationId=...` | прогон SkillMiner прямо сейчас, найденное в `AGENTIK_SKILLS_DIR` |
|
||||||
|
| `POST /debug/curate` | прогон Curator.runPass (архивация stale-заметок памяти) |
|
||||||
|
| `POST /debug/compact?conversationId=...` | принудительный compaction working memory диалога |
|
||||||
|
| `GET /debug/tokens?conversationId=...` | token-статистика диалога из БД (turns/in/out/total) |
|
||||||
|
|
||||||
|
Каждый возвращает JSON с результатом (что сохранил/нашёл/сжал), чтобы было
|
||||||
|
видно не только "триггер сработал", а что именно LLM намайнила.
|
||||||
|
|
||||||
```bash
|
```bash
|
||||||
AGENTIK_SKILLS_DIR=/etc/agentik/skills
|
AGENTIK_SKILLS_DIR=/etc/agentik/skills
|
||||||
|
AGENTIK_DEBUG_ENDPOINTS=1
|
||||||
```
|
```
|
||||||
|
|
||||||
|
|
||||||
## MCP-инструменты
|
## MCP-инструменты
|
||||||
|
|
||||||
`AGENTIK_MCP_CONFIG` — путь к JSON-файлу со списком MCP-серверов
|
`AGENTIK_MCP_CONFIG` — путь к JSON-файлу со списком MCP-серверов
|
||||||
|
|||||||
@@ -0,0 +1,150 @@
|
|||||||
|
package pw.binom.agentik.standalone
|
||||||
|
|
||||||
|
import io.ktor.http.ContentType
|
||||||
|
import io.ktor.http.HttpStatusCode
|
||||||
|
import io.ktor.server.response.respondText
|
||||||
|
import io.ktor.server.routing.Route
|
||||||
|
import io.ktor.server.routing.get
|
||||||
|
import io.ktor.server.routing.post
|
||||||
|
import kotlinx.serialization.json.buildJsonObject
|
||||||
|
import kotlinx.serialization.json.put
|
||||||
|
import pw.binom.agentik.memory.ConversationTurn
|
||||||
|
import pw.binom.agentik.proto.Agent
|
||||||
|
import pw.binom.agentik.standalone.agent.ChatConversation
|
||||||
|
import pw.binom.agentik.standalone.agent.LlmReflector
|
||||||
|
import pw.binom.agentik.standalone.agent.SkillMiner
|
||||||
|
import pw.binom.agentik.standalone.agent.memory.Curator
|
||||||
|
import pw.binom.agentik.standalone.persistence.sqlite.SqliteStores
|
||||||
|
import pw.binom.agentik.skills.SkillStore
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Debug-эндпоинты для ручного триггерирования фоновых фич (без ожидания
|
||||||
|
* интервалов). Подключаются только при `AGENTIK_DEBUG_ENDPOINTS=1`:
|
||||||
|
*
|
||||||
|
* - `POST /debug/reflect?conversationId=...` — прогон [LlmReflector] прямо сейчас
|
||||||
|
* - `POST /debug/skill-mine?conversationId=...` — прогон [SkillMiner] прямо сейчас
|
||||||
|
* - `POST /debug/curate` — прогон [Curator.runPass] прямо сейчас
|
||||||
|
* - `POST /debug/compact?conversationId=...` — принудительный compaction
|
||||||
|
* - `GET /debug/tokens?conversationId=...` — token-статистика диалога из БД
|
||||||
|
*
|
||||||
|
* Каждый возвращает JSON с результатом (что сохранил / нашёл / сжал), чтобы в
|
||||||
|
* тестах было видно не только "триггер сработал", а что именно LLM намайнила.
|
||||||
|
* Только локальная отладка: endpoint'ы без авторизации, в проде не включать.
|
||||||
|
*/
|
||||||
|
internal fun Route.debugRoutes(
|
||||||
|
agent: Agent,
|
||||||
|
stores: SqliteStores,
|
||||||
|
reflector: LlmReflector?,
|
||||||
|
skillMiner: SkillMiner?,
|
||||||
|
skillStore: SkillStore?,
|
||||||
|
curator: Curator?,
|
||||||
|
) {
|
||||||
|
post("/debug/reflect") {
|
||||||
|
val convId = call.parameters["conversationId"]
|
||||||
|
?: return@post call.respondText("conversationId required", status = HttpStatusCode.BadRequest)
|
||||||
|
val minerReflector = reflector
|
||||||
|
if (minerReflector == null) return@post call.respondText("reflection disabled", status = HttpStatusCode.NotFound)
|
||||||
|
val turns = recentTurns(stores, convId, maxTurns = 6)
|
||||||
|
if (turns.isEmpty()) return@post call.respondText("no turns in conversation", status = HttpStatusCode.NotFound)
|
||||||
|
val reflection = minerReflector.reflect(turns)
|
||||||
|
if (reflection == null) {
|
||||||
|
call.respondText("""{"reflected":false,"reason":"unparseable LLM reply"}""", contentType = ContentType.Application.Json)
|
||||||
|
} else {
|
||||||
|
val stamped = reflection.copy(conversationId = convId)
|
||||||
|
stores.reflections.insert(stamped)
|
||||||
|
call.respondText(
|
||||||
|
buildJsonObject {
|
||||||
|
put("reflected", true)
|
||||||
|
put("id", stamped.id)
|
||||||
|
put("score", stamped.score.toString())
|
||||||
|
put("summary", stamped.summary)
|
||||||
|
put("weakSpots", kotlinx.serialization.json.JsonArray(stamped.weakSpots.map { kotlinx.serialization.json.JsonPrimitive(it) }))
|
||||||
|
}.toString(),
|
||||||
|
contentType = ContentType.Application.Json,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
post("/debug/skill-mine") {
|
||||||
|
val convId = call.parameters["conversationId"]
|
||||||
|
?: return@post call.respondText("conversationId required", status = HttpStatusCode.BadRequest)
|
||||||
|
val miner = skillMiner
|
||||||
|
val store = skillStore
|
||||||
|
if (miner == null || store == null) return@post call.respondText("skill mining disabled", status = HttpStatusCode.NotFound)
|
||||||
|
val turns = recentTurns(stores, convId, maxTurns = miner.maxTurns)
|
||||||
|
if (turns.isEmpty()) return@post call.respondText("no turns in conversation", status = HttpStatusCode.NotFound)
|
||||||
|
val mined = miner.mine(turns, store.catalog.skills)
|
||||||
|
for (s in mined) store.upsert(s)
|
||||||
|
val json = buildJsonObject {
|
||||||
|
put("mined", mined.size.toString())
|
||||||
|
put(
|
||||||
|
"skills",
|
||||||
|
kotlinx.serialization.json.JsonArray(
|
||||||
|
mined.map { s ->
|
||||||
|
kotlinx.serialization.json.buildJsonObject {
|
||||||
|
put("name", s.name)
|
||||||
|
put("description", s.description)
|
||||||
|
}
|
||||||
|
},
|
||||||
|
),
|
||||||
|
)
|
||||||
|
}.toString()
|
||||||
|
call.respondText(json, contentType = ContentType.Application.Json)
|
||||||
|
}
|
||||||
|
|
||||||
|
post("/debug/curate") {
|
||||||
|
val c = curator ?: return@post call.respondText("curator disabled (memory off?)", status = HttpStatusCode.NotFound)
|
||||||
|
val archived = c.runPass()
|
||||||
|
call.respondText("""{"archived":$archived}""", contentType = ContentType.Application.Json)
|
||||||
|
}
|
||||||
|
|
||||||
|
post("/debug/compact") {
|
||||||
|
val convId = call.parameters["conversationId"]
|
||||||
|
?: return@post call.respondText("conversationId required", status = HttpStatusCode.BadRequest)
|
||||||
|
val conv = agent.getConversation(convId) ?: return@post call.respondText("conversation not found", status = HttpStatusCode.NotFound)
|
||||||
|
val chatConv = conv as? ChatConversation ?: return@post call.respondText("not a ChatConversation", status = HttpStatusCode.InternalServerError)
|
||||||
|
val ok = chatConv.forceCompactNow()
|
||||||
|
call.respondText("""{"compacted":$ok}""", contentType = ContentType.Application.Json)
|
||||||
|
}
|
||||||
|
|
||||||
|
get("/debug/tokens") {
|
||||||
|
val convId = call.parameters["conversationId"]
|
||||||
|
?: return@get call.respondText("conversationId required", status = HttpStatusCode.BadRequest)
|
||||||
|
val stats = stores.messages.tokenStats(convId)
|
||||||
|
val json = buildJsonObject {
|
||||||
|
put("conversationId", convId)
|
||||||
|
put("turns", stats.turns.toString())
|
||||||
|
put("inputTokens", stats.inputTokens.toString())
|
||||||
|
put("outputTokens", stats.outputTokens.toString())
|
||||||
|
put("totalTokens", (stats.inputTokens + stats.outputTokens).toString())
|
||||||
|
}.toString()
|
||||||
|
call.respondText(json, contentType = ContentType.Application.Json)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Последние [maxTurns] пар user/assistant из working memory диалога
|
||||||
|
* (для debug-триггеров reflector/miner; та же логика, что у хуков
|
||||||
|
* [ChatConversation]).
|
||||||
|
*/
|
||||||
|
internal suspend fun recentTurns(stores: SqliteStores, conversationId: String, maxTurns: Int): List<ConversationTurn> {
|
||||||
|
val rows = stores.workingMemory.list(conversationId)
|
||||||
|
val pairs = mutableListOf<ConversationTurn>()
|
||||||
|
var pendingUser: String? = null
|
||||||
|
for (row in rows) {
|
||||||
|
when (val e = row.entry) {
|
||||||
|
is pw.binom.agentik.standalone.persistence.WorkingMemoryEntry.User -> pendingUser = e.content.text()
|
||||||
|
is pw.binom.agentik.standalone.persistence.WorkingMemoryEntry.Assistant -> {
|
||||||
|
val user = pendingUser ?: ""
|
||||||
|
pendingUser = null
|
||||||
|
pairs += ConversationTurn(userMessage = user, assistantMessage = e.content.text())
|
||||||
|
}
|
||||||
|
else -> {}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return pairs.takeLast(maxTurns)
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Текстовое содержимое записей working memory (Text-контент, без картинок). */
|
||||||
|
internal fun List<pw.binom.agentik.standalone.persistence.Content>.text(): String =
|
||||||
|
filterIsInstance<pw.binom.agentik.standalone.persistence.Content.Text>().joinToString("\n") { it.body }
|
||||||
@@ -67,6 +67,17 @@ fun main() {
|
|||||||
}
|
}
|
||||||
val skills = skillStore?.catalog ?: SkillCatalog.EMPTY
|
val skills = skillStore?.catalog ?: SkillCatalog.EMPTY
|
||||||
|
|
||||||
|
// Skill mining: фоновый LLM-прогон, который находит переиспользуемые скилы,
|
||||||
|
// которые модель забыла сохранить через `skill_save`. Работает только когда
|
||||||
|
// есть куда писать (skillStore) и интервал > 0.
|
||||||
|
val skillMiner: pw.binom.agentik.standalone.agent.SkillMiner? =
|
||||||
|
if (skillStore != null && config.skillMiningInterval > 0) {
|
||||||
|
pw.binom.agentik.standalone.agent.SkillMiner(
|
||||||
|
llm = llm,
|
||||||
|
maxTurns = config.skillMiningMaxTurns,
|
||||||
|
)
|
||||||
|
} else null
|
||||||
|
|
||||||
// SOUL.md — файл персоны. Если задан — читается как plain text/markdown,
|
// SOUL.md — файл персоны. Если задан — читается как plain text/markdown,
|
||||||
// вставляется в самое начало systemInstruction. Если отсутствует — exit-code != 0
|
// вставляется в самое начало systemInstruction. Если отсутствует — exit-code != 0
|
||||||
// (на старте агента это фатально: нечего показывать LLM).
|
// (на старте агента это фатально: нечего показывать LLM).
|
||||||
@@ -186,12 +197,34 @@ fun main() {
|
|||||||
recentReflections = recentReflections,
|
recentReflections = recentReflections,
|
||||||
reflector = reflector,
|
reflector = reflector,
|
||||||
reflectionInterval = config.reflectionInterval,
|
reflectionInterval = config.reflectionInterval,
|
||||||
|
skillMiner = skillMiner,
|
||||||
|
skillMiningInterval = config.skillMiningInterval,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
// Куратор памяти: фоновая архивация stale-заметок. Поднимается до server'а,
|
||||||
|
// чтобы debug-эндпоинты могли его триггерить вручную.
|
||||||
|
val curator: pw.binom.agentik.standalone.agent.memory.Curator? =
|
||||||
|
if (memorySystem != null) {
|
||||||
|
val c = pw.binom.agentik.standalone.agent.memory.Curator(memorySystem.store)
|
||||||
|
c.start()
|
||||||
|
Runtime.getRuntime().addShutdownHook(Thread { c.stop() })
|
||||||
|
c
|
||||||
|
} else null
|
||||||
|
|
||||||
val server = embeddedServer(CIO, port = config.port) {
|
val server = embeddedServer(CIO, port = config.port) {
|
||||||
routing {
|
routing {
|
||||||
get("/health") { call.respondText("ok") }
|
get("/health") { call.respondText("ok") }
|
||||||
agentikAgent(agent, path = "/agentik")
|
agentikAgent(agent, path = "/agentik")
|
||||||
|
if (config.debugEndpoints) {
|
||||||
|
debugRoutes(
|
||||||
|
agent = agent,
|
||||||
|
stores = stores,
|
||||||
|
reflector = reflector,
|
||||||
|
skillMiner = skillMiner,
|
||||||
|
skillStore = skillStore,
|
||||||
|
curator = curator,
|
||||||
|
)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
println("agentik standalone listening on http://localhost:${config.port}")
|
println("agentik standalone listening on http://localhost:${config.port}")
|
||||||
@@ -209,12 +242,15 @@ fun main() {
|
|||||||
} else {
|
} else {
|
||||||
println(" compaction: disabled (OPENAI_CONTEXT_WINDOW not set)")
|
println(" compaction: disabled (OPENAI_CONTEXT_WINDOW not set)")
|
||||||
}
|
}
|
||||||
if (memorySystem != null) {
|
if (curator != 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})")
|
println(" curator: enabled (interval=${pw.binom.agentik.standalone.agent.memory.Curator.DEFAULT_INTERVAL}, maxAge=${pw.binom.agentik.standalone.agent.memory.Curator.DEFAULT_MAX_AGE})")
|
||||||
}
|
}
|
||||||
|
if (skillMiner != null) {
|
||||||
|
println(" skill-mining: enabled (interval=${config.skillMiningInterval} turns, maxTurns=${config.skillMiningMaxTurns})")
|
||||||
|
}
|
||||||
|
if (config.debugEndpoints) {
|
||||||
|
println(" debug endpoints: enabled (/debug/reflect, /debug/skill-mine, /debug/curate, /debug/compact, /debug/tokens)")
|
||||||
|
}
|
||||||
// Token stats по существующим диалогам (агрегат на старте — каждая запись
|
// Token stats по существующим диалогам (агрегат на старте — каждая запись
|
||||||
// парсится из payload_json, ну >100 turns и БД приличная — но в рамках
|
// парсится из payload_json, ну >100 turns и БД приличная — но в рамках
|
||||||
// стартапа это терпимо).
|
// стартапа это терпимо).
|
||||||
|
|||||||
@@ -92,6 +92,16 @@ class ChatAgent(
|
|||||||
* Через сколько пользовательских ходов запускать рефлексию. `0` = выключено.
|
* Через сколько пользовательских ходов запускать рефлексию. `0` = выключено.
|
||||||
*/
|
*/
|
||||||
private val reflectionInterval: Int = 0,
|
private val reflectionInterval: Int = 0,
|
||||||
|
/**
|
||||||
|
* Фоновый минер скилов: каждые N ходов LLM смотрит последние ходы и
|
||||||
|
* upsert-ит переиспользуемые скилы в [skillStore]. `null` = mining выключен.
|
||||||
|
* Сетка безопасности, если модель забыла вызвать `skill_save` сама.
|
||||||
|
*/
|
||||||
|
private val skillMiner: SkillMiner? = null,
|
||||||
|
/**
|
||||||
|
* Через сколько пользовательских ходов запускать skill mining. `0` = выключено.
|
||||||
|
*/
|
||||||
|
private val skillMiningInterval: Int = 0,
|
||||||
) : ProtoAgent, AutoCloseable {
|
) : ProtoAgent, AutoCloseable {
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -173,6 +183,9 @@ class ChatAgent(
|
|||||||
reflectionStore = stores.reflections,
|
reflectionStore = stores.reflections,
|
||||||
reflector = reflector,
|
reflector = reflector,
|
||||||
reflectionInterval = reflectionInterval,
|
reflectionInterval = reflectionInterval,
|
||||||
|
skillMiner = skillMiner,
|
||||||
|
skillMiningStore = skillStore,
|
||||||
|
skillMiningInterval = skillMiningInterval,
|
||||||
)
|
)
|
||||||
runBlocking {
|
runBlocking {
|
||||||
liveLock.withLock { live[conv.id] = conv }
|
liveLock.withLock { live[conv.id] = conv }
|
||||||
@@ -220,6 +233,9 @@ class ChatAgent(
|
|||||||
reflectionStore = stores.reflections,
|
reflectionStore = stores.reflections,
|
||||||
reflector = reflector,
|
reflector = reflector,
|
||||||
reflectionInterval = reflectionInterval,
|
reflectionInterval = reflectionInterval,
|
||||||
|
skillMiner = skillMiner,
|
||||||
|
skillMiningStore = skillStore,
|
||||||
|
skillMiningInterval = skillMiningInterval,
|
||||||
)
|
)
|
||||||
|
|
||||||
override fun close() {
|
override fun close() {
|
||||||
|
|||||||
+98
-8
@@ -16,6 +16,7 @@ import kotlinx.coroutines.runBlocking
|
|||||||
import kotlinx.coroutines.sync.Mutex
|
import kotlinx.coroutines.sync.Mutex
|
||||||
import kotlinx.coroutines.sync.withLock
|
import kotlinx.coroutines.sync.withLock
|
||||||
import pw.binom.agentik.memory.MemoryNote
|
import pw.binom.agentik.memory.MemoryNote
|
||||||
|
import pw.binom.agentik.memory.ConversationTurn
|
||||||
import pw.binom.agentik.memory.MemoryPrefetcher
|
import pw.binom.agentik.memory.MemoryPrefetcher
|
||||||
import pw.binom.agentik.memory.MemoryReviewDecision
|
import pw.binom.agentik.memory.MemoryReviewDecision
|
||||||
import pw.binom.agentik.memory.MemoryReviewer
|
import pw.binom.agentik.memory.MemoryReviewer
|
||||||
@@ -121,6 +122,17 @@ class ChatConversation(
|
|||||||
* Через сколько пользовательских ходов запускать рефлексию. `0` = выключено.
|
* Через сколько пользовательских ходов запускать рефлексию. `0` = выключено.
|
||||||
*/
|
*/
|
||||||
private val reflectionInterval: Int = 0,
|
private val reflectionInterval: Int = 0,
|
||||||
|
/**
|
||||||
|
* Фоновый минер скилов (сетка безопасности для skill self-improvement):
|
||||||
|
* каждые [skillMiningInterval] пользовательских ходов LLM смотрит
|
||||||
|
* последние ходы и upsert-ит переиспользуемые скилы в [skillMiningStore].
|
||||||
|
* `null` = mining выключен.
|
||||||
|
*/
|
||||||
|
private val skillMiner: SkillMiner? = null,
|
||||||
|
/** Хранилище, куда mining upsert-ит найденные скилы. */
|
||||||
|
private val skillMiningStore: pw.binom.agentik.skills.SkillStore? = null,
|
||||||
|
/** Через сколько пользовательских ходов запускать mining. `0` = выключено. */
|
||||||
|
private val skillMiningInterval: Int = 0,
|
||||||
) : ProtoConversation, AutoCloseable {
|
) : ProtoConversation, AutoCloseable {
|
||||||
|
|
||||||
private var record: ConversationRecord = record
|
private var record: ConversationRecord = record
|
||||||
@@ -352,6 +364,7 @@ class ChatConversation(
|
|||||||
|
|
||||||
scheduleReview(userRecord, assistantContent)
|
scheduleReview(userRecord, assistantContent)
|
||||||
scheduleReflection(userRecord, assistantContent)
|
scheduleReflection(userRecord, assistantContent)
|
||||||
|
scheduleSkillMining(userRecord, assistantContent)
|
||||||
|
|
||||||
emitEvent(ProtoEvent.End(date = assistantAt))
|
emitEvent(ProtoEvent.End(date = assistantAt))
|
||||||
}
|
}
|
||||||
@@ -402,10 +415,25 @@ class ChatConversation(
|
|||||||
* No-op когда [contextWindow] или [contextCompactor] == null.
|
* No-op когда [contextWindow] или [contextCompactor] == null.
|
||||||
*/
|
*/
|
||||||
private suspend fun compactPreTurnIfNeeded() {
|
private suspend fun compactPreTurnIfNeeded() {
|
||||||
val window = contextWindow ?: return
|
compactPreTurn(force = false)
|
||||||
val compactor = contextCompactor ?: return
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Принудительный compaction без проверки порога — для debug-эндпоинта
|
||||||
|
* `POST /debug/compact`. Сжимает working memory независимо от текущей
|
||||||
|
* загрузки контекста. Возвращает `true` если суммаризация выполнена.
|
||||||
|
*/
|
||||||
|
suspend fun forceCompactNow(): Boolean = compactPreTurn(force = true)
|
||||||
|
|
||||||
|
/**
|
||||||
|
* [compactPreTurnIfNeeded] с явным флагом [force]: при `force = true`
|
||||||
|
* порог [compressionThreshold] не проверяется (debug-триггер).
|
||||||
|
*/
|
||||||
|
private suspend fun compactPreTurn(force: Boolean): Boolean {
|
||||||
|
val window = contextWindow ?: return false
|
||||||
|
val compactor = contextCompactor ?: return false
|
||||||
val wm = workingMemory.list(id)
|
val wm = workingMemory.list(id)
|
||||||
if (wm.isEmpty()) return
|
if (wm.isEmpty()) return false
|
||||||
|
|
||||||
val systemText = wm.firstOrNull { it.entry is WorkingMemoryEntry.System }
|
val systemText = wm.firstOrNull { it.entry is WorkingMemoryEntry.System }
|
||||||
?.let { (it.entry as WorkingMemoryEntry.System).text }
|
?.let { (it.entry as WorkingMemoryEntry.System).text }
|
||||||
@@ -419,7 +447,7 @@ class ChatConversation(
|
|||||||
toolsChars = toolsChars,
|
toolsChars = toolsChars,
|
||||||
)
|
)
|
||||||
|
|
||||||
if (estimated.toDouble() / window < compressionThreshold) return
|
if (!force && estimated.toDouble() / window < compressionThreshold) return false
|
||||||
|
|
||||||
// Берём для свёртки старые ходы, последние KEEP_RECENT_TURNS оставляем
|
// Берём для свёртки старые ходы, последние KEEP_RECENT_TURNS оставляем
|
||||||
// как есть — это самая свежая часть контекста, которая нужна модели для
|
// как есть — это самая свежая часть контекста, которая нужна модели для
|
||||||
@@ -430,7 +458,7 @@ class ChatConversation(
|
|||||||
} else {
|
} else {
|
||||||
history
|
history
|
||||||
}
|
}
|
||||||
if (toCompact.isEmpty()) return
|
if (toCompact.isEmpty()) return false
|
||||||
|
|
||||||
val turns = toCompact.mapNotNull { row ->
|
val turns = toCompact.mapNotNull { row ->
|
||||||
when (val e = row.entry) {
|
when (val e = row.entry) {
|
||||||
@@ -464,7 +492,7 @@ class ChatConversation(
|
|||||||
if (pendingUser != null) paired.add(pendingUser)
|
if (pendingUser != null) paired.add(pendingUser)
|
||||||
if (paired.isEmpty()) {
|
if (paired.isEmpty()) {
|
||||||
log.info { "compactPreTurn: nothing to compact for $id" }
|
log.info { "compactPreTurn: nothing to compact for $id" }
|
||||||
return
|
return false
|
||||||
}
|
}
|
||||||
|
|
||||||
val summaryText = try {
|
val summaryText = try {
|
||||||
@@ -473,9 +501,9 @@ class ChatConversation(
|
|||||||
throw e
|
throw e
|
||||||
} catch (e: Throwable) {
|
} catch (e: Throwable) {
|
||||||
log.warn(e) { "context summarization failed for $id: ${e.message}" }
|
log.warn(e) { "context summarization failed for $id: ${e.message}" }
|
||||||
return
|
return false
|
||||||
}
|
}
|
||||||
if (summaryText.isBlank()) return
|
if (summaryText.isBlank()) return false
|
||||||
|
|
||||||
// Триггер памяти: до удаления ходов даём ревьюеру шанс вытащить факты.
|
// Триггер памяти: до удаления ходов даём ревьюеру шанс вытащить факты.
|
||||||
val reviewer = memoryReviewer
|
val reviewer = memoryReviewer
|
||||||
@@ -523,6 +551,7 @@ class ChatConversation(
|
|||||||
if (after.toDouble() / window >= compressionThreshold) {
|
if (after.toDouble() / window >= compressionThreshold) {
|
||||||
log.warn { "compactPreTurn: still over threshold for $id (estimated=$after, window=$window, threshold=$compressionThreshold). Consider raising contextWindow or lowering threshold." }
|
log.warn { "compactPreTurn: still over threshold for $id (estimated=$after, window=$window, threshold=$compressionThreshold). Consider raising contextWindow or lowering threshold." }
|
||||||
}
|
}
|
||||||
|
return true
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -643,6 +672,67 @@ class ChatConversation(
|
|||||||
count
|
count
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Skill-mining триггер: каждые [skillMiningInterval] пользовательских ходов
|
||||||
|
* запускает фоновый [SkillMiner.mine] по последним ходам диалога (из
|
||||||
|
* working memory, не только текущий ход — минеру нужен контекст паттерна)
|
||||||
|
* и upsert-ит найденные скилы в [skillMiningStore]. Не блокирует turn.
|
||||||
|
*
|
||||||
|
* Это сетка безопасности для skill self-improvement: если модель в ходе
|
||||||
|
* разговора "протупила" и не вызвала `skill_save`, минер добирает её
|
||||||
|
* постфактум. Для temp-бесед и без miner'а — no-op.
|
||||||
|
*/
|
||||||
|
private fun scheduleSkillMining(
|
||||||
|
userRecord: MessageRecord.UserMessage,
|
||||||
|
assistantContent: List<Content>,
|
||||||
|
) {
|
||||||
|
if (skillMiningInterval <= 0) return
|
||||||
|
val miner = skillMiner ?: return
|
||||||
|
val store = skillMiningStore ?: return
|
||||||
|
if (record.isTemporal) return
|
||||||
|
val userTurnCount = countUserTurnsBlocking()
|
||||||
|
if (userTurnCount % skillMiningInterval != 0) return
|
||||||
|
val convId = id
|
||||||
|
scope.launch {
|
||||||
|
try {
|
||||||
|
val turns = recentTurnsFromWorkingMemory(miner.maxTurns)
|
||||||
|
if (turns.isEmpty()) return@launch
|
||||||
|
val existing = store.catalog.skills
|
||||||
|
val mined = miner.mine(turns, existing)
|
||||||
|
for (s in mined) {
|
||||||
|
runCatching { store.upsert(s) }
|
||||||
|
.onFailure { log.warn(it) { "skill-mine upsert '${s.name}' failed: ${it.message}" } }
|
||||||
|
}
|
||||||
|
log.info { "skill-mine: conv=$convId turns=${turns.size} existing=${existing.size} mined=${mined.size}" }
|
||||||
|
} catch (e: Throwable) {
|
||||||
|
log.warn(e) { "skill-mine failed for $convId: ${e.message}" }
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Собирает последние [maxTurns] пар user/assistant из working memory,
|
||||||
|
* упорядоченных по хронологии (для [SkillMiner]). System/Compacted/System
|
||||||
|
* записи пропускаются: минеру нужен именно разговор.
|
||||||
|
*/
|
||||||
|
private suspend fun recentTurnsFromWorkingMemory(maxTurns: Int): List<ConversationTurn> {
|
||||||
|
val rows = workingMemory.list(id)
|
||||||
|
val pairs = mutableListOf<ConversationTurn>()
|
||||||
|
var pendingUser: String? = null
|
||||||
|
for (row in rows) {
|
||||||
|
when (val e = row.entry) {
|
||||||
|
is WorkingMemoryEntry.User -> pendingUser = e.content.text()
|
||||||
|
is WorkingMemoryEntry.Assistant -> {
|
||||||
|
val user = pendingUser ?: ""
|
||||||
|
pendingUser = null
|
||||||
|
pairs += ConversationTurn(userMessage = user, assistantMessage = e.content.text())
|
||||||
|
}
|
||||||
|
else -> {}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return pairs.takeLast(maxTurns)
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Один tool-call: эмитим Event.ToolCall, выполняем tool (MCP), эмитим Event.ToolResult,
|
* Один tool-call: эмитим Event.ToolCall, выполняем tool (MCP), эмитим Event.ToolResult,
|
||||||
* пишем в audit + working memory, подаём результат в LiteConversation.
|
* пишем в audit + working memory, подаём результат в LiteConversation.
|
||||||
|
|||||||
@@ -0,0 +1,78 @@
|
|||||||
|
package pw.binom.agentik.standalone.agent
|
||||||
|
|
||||||
|
import kotlinx.coroutines.CoroutineDispatcher
|
||||||
|
import kotlinx.coroutines.Dispatchers
|
||||||
|
import kotlinx.coroutines.withContext
|
||||||
|
import mu.KotlinLogging
|
||||||
|
import pw.binom.agentik.memory.ConversationTurn
|
||||||
|
import pw.binom.agentik.skills.SkillFile
|
||||||
|
import pw.binom.litert.LiteConversationConfig
|
||||||
|
import pw.binom.litert.LiteLlm
|
||||||
|
|
||||||
|
private val log = mu.KotlinLogging.logger {}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Фоновый минер скилов (skill mining): сетка безопасности для skill self-improvement.
|
||||||
|
*
|
||||||
|
* Модель в ходе разговора может "протупить" и не вызвать `skill_save`, хотя приём
|
||||||
|
* был действительно переиспользуемым. [SkillMiner] периодически (каждые N ходов,
|
||||||
|
* см. `AGENTIK_SKILL_MINING_INTERVAL`) берёт последние ходы диалога, показывает их
|
||||||
|
* LLM вместе с текущим каталогом скилов и просит structured-output JSON:
|
||||||
|
* {"skills": [{"name","description","body"}]}. Найденные скилы upsert-ятся в
|
||||||
|
* [pw.binom.agentik.skills.SkillStore] — агент становится умнее между сессиями
|
||||||
|
* даже без прямого tool-call в ходе разговора.
|
||||||
|
*
|
||||||
|
* Архитектурно — точный аналог [LlmReflector]: короткоживущий LiteConversation
|
||||||
|
* (один прогон = один LLM-вызов), blocking-инференс на [dispatcher], defensive
|
||||||
|
* парсинг [SkillMiningParser] (кривой ответ → пустой список, не ломает agent loop).
|
||||||
|
*
|
||||||
|
* @param llm LLM-бэкенд
|
||||||
|
* @param maxTurns сколько последних ходов передавать модели (default 30)
|
||||||
|
* @param maxTokens потолок ответа модели (default 1536 — body скила бывает длинным)
|
||||||
|
* @param dispatcher диспетчер для блокирующего LLM-вызова
|
||||||
|
*/
|
||||||
|
class SkillMiner(
|
||||||
|
private val llm: LiteLlm,
|
||||||
|
val maxTurns: Int = 30,
|
||||||
|
private val maxTokens: Int = 1536,
|
||||||
|
private val dispatcher: CoroutineDispatcher = Dispatchers.IO,
|
||||||
|
) {
|
||||||
|
/**
|
||||||
|
* Mine по последним [turns] с учётом текущего каталога [existing].
|
||||||
|
*
|
||||||
|
* @return найденные/обновлённые скилы; пустой список — нечего сохранять
|
||||||
|
* или модель ответила мусором (defensive: прогон просто пропускается).
|
||||||
|
*
|
||||||
|
* Вызов блокирующий: ~1-3 с на CPU для on-device LiteRT-LM. Всегда вызывать
|
||||||
|
* из background scope (хук [pw.binom.agentik.standalone.agent.ChatConversation]
|
||||||
|
* или debug-эндпоинт).
|
||||||
|
*/
|
||||||
|
suspend fun mine(turns: List<ConversationTurn>, existing: List<SkillFile>): List<SkillFile> =
|
||||||
|
withContext(dispatcher) {
|
||||||
|
val batch = turns.takeLast(maxTurns)
|
||||||
|
if (batch.isEmpty()) return@withContext emptyList()
|
||||||
|
val conversation = llm.createConversation(
|
||||||
|
config = LiteConversationConfig(
|
||||||
|
systemInstruction = SkillMiningPrompts.SYSTEM_PROMPT,
|
||||||
|
initialMessages = emptyList(),
|
||||||
|
tools = emptyList(),
|
||||||
|
temperature = 0.2f,
|
||||||
|
maxTokens = maxTokens,
|
||||||
|
)
|
||||||
|
)
|
||||||
|
try {
|
||||||
|
val userPrompt = SkillMiningPrompts.buildUserPrompt(batch, existing)
|
||||||
|
val raw = conversation.send(userPrompt)
|
||||||
|
val mined = SkillMiningParser.parse(raw)
|
||||||
|
if (mined.isNotEmpty()) {
|
||||||
|
log.info { "skill-mine: found ${mined.size} skill(s) from ${batch.size} turns: ${mined.map { it.name }}" }
|
||||||
|
}
|
||||||
|
mined
|
||||||
|
} catch (e: Throwable) {
|
||||||
|
log.warn(e) { "skill-mine: LLM call failed, skipping pass" }
|
||||||
|
emptyList()
|
||||||
|
} finally {
|
||||||
|
conversation.close()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,207 @@
|
|||||||
|
package pw.binom.agentik.standalone.agent
|
||||||
|
|
||||||
|
import pw.binom.agentik.skills.SkillFile
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Минимальный парсер JSON-ответа [SkillMiner] (в том же defensive-стиле, что
|
||||||
|
* [ReflectionParser] и [ReviewDecisionParser]: без kotlinx-serialization, чтобы
|
||||||
|
* кривой ответ локальной модели не ронял agent loop).
|
||||||
|
*
|
||||||
|
* Ожидаемая форма:
|
||||||
|
* ```
|
||||||
|
* {"skills": [{"name": "x", "description": "...", "body": "..."}]}
|
||||||
|
* ```
|
||||||
|
* Допуски:
|
||||||
|
* - ответ может быть обёрнут в ```json ... ``` fences;
|
||||||
|
* - может быть текст до/после JSON — ищем первый `{` (или `[`) и балансируем;
|
||||||
|
* - допустим и голый массив `[{...}, {...}]` без ключа `"skills"`;
|
||||||
|
* - `body` может содержать `\n`, `\"`, `\\` — деэскейпим;
|
||||||
|
* - `description`/`body` могут отсутствовать (тогда пустые);
|
||||||
|
* - любой мусор (невалидные кавычки, незакрытые скобки) → пустой список,
|
||||||
|
* mining-прогон просто не сохранит ничего.
|
||||||
|
*/
|
||||||
|
object SkillMiningParser {
|
||||||
|
|
||||||
|
fun parse(raw: String): List<SkillFile> {
|
||||||
|
val blob = extractJsonBlob(raw) ?: return emptyList()
|
||||||
|
val array = extractSkillsArray(blob) ?: return emptyList()
|
||||||
|
return splitTopLevelObjects(array).mapNotNull { obj ->
|
||||||
|
val name = extractStringField(obj, "name")?.trim().orEmpty()
|
||||||
|
if (name.isEmpty()) return@mapNotNull null
|
||||||
|
val description = extractStringField(obj, "description")?.trim().orEmpty()
|
||||||
|
val body = extractStringField(obj, "body")?.trim().orEmpty()
|
||||||
|
SkillFile(name = name, description = description, body = body)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Обрезает ``` fences, находит первый `{` или `[` и возвращает подстроку
|
||||||
|
* до парной закрывающей (со знанием состояний string/escape).
|
||||||
|
*/
|
||||||
|
internal fun extractJsonBlob(raw: String): String? {
|
||||||
|
var s = raw.trim()
|
||||||
|
if (s.startsWith("```")) {
|
||||||
|
val nl = s.indexOf('\n')
|
||||||
|
if (nl > 0) s = s.substring(nl + 1)
|
||||||
|
if (s.endsWith("```")) s = s.substring(0, s.length - 3)
|
||||||
|
}
|
||||||
|
val open = minOf(
|
||||||
|
s.indexOf('{').takeIf { it >= 0 } ?: Int.MAX_VALUE,
|
||||||
|
s.indexOf('[').takeIf { it >= 0 } ?: Int.MAX_VALUE,
|
||||||
|
)
|
||||||
|
if (open == Int.MAX_VALUE) return null
|
||||||
|
val opener = s[open]
|
||||||
|
val closer = if (opener == '{') '}' else ']'
|
||||||
|
var depth = 0
|
||||||
|
var inString = false
|
||||||
|
var escape = false
|
||||||
|
for (i in open until s.length) {
|
||||||
|
val c = s[i]
|
||||||
|
if (escape) { escape = false; continue }
|
||||||
|
if (inString) {
|
||||||
|
when (c) {
|
||||||
|
'\\' -> escape = true
|
||||||
|
'"' -> inString = false
|
||||||
|
}
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
when (c) {
|
||||||
|
'"' -> inString = true
|
||||||
|
opener, '{', '[' -> depth++
|
||||||
|
'}', ']' -> {
|
||||||
|
depth--
|
||||||
|
if (depth == 0 && ((c == closer) || (opener == '{' && c == '}') || (opener == '[' && c == ']'))) {
|
||||||
|
return s.substring(open, i + 1)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return null
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Находит массив скилов: если в blob есть ключ `"skills"` — массив после него,
|
||||||
|
* иначе сам blob (если начинается с `[`).
|
||||||
|
*/
|
||||||
|
internal fun extractSkillsArray(blob: String): String? {
|
||||||
|
val keyIdx = blob.indexOf("\"skills\"")
|
||||||
|
if (keyIdx >= 0) {
|
||||||
|
val colon = blob.indexOf(':', keyIdx + "skills".length + 2)
|
||||||
|
if (colon < 0) return null
|
||||||
|
val open = blob.indexOf('[', colon + 1)
|
||||||
|
if (open < 0) return null
|
||||||
|
return balanceArray(blob, open)
|
||||||
|
}
|
||||||
|
if (blob.startsWith('[')) return blob.drop(1).dropLast(1)
|
||||||
|
return null
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Балансирует `[...]` от [start] (включительно). Возвращает содержимое без скобок. */
|
||||||
|
private fun balanceArray(s: String, start: Int): String? {
|
||||||
|
var depth = 0
|
||||||
|
var inString = false
|
||||||
|
var escape = false
|
||||||
|
for (i in start until s.length) {
|
||||||
|
val c = s[i]
|
||||||
|
if (escape) { escape = false; continue }
|
||||||
|
if (inString) {
|
||||||
|
when (c) {
|
||||||
|
'\\' -> escape = true
|
||||||
|
'"' -> inString = false
|
||||||
|
}
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
when (c) {
|
||||||
|
'"' -> inString = true
|
||||||
|
'[' -> depth++
|
||||||
|
']' -> {
|
||||||
|
depth--
|
||||||
|
if (depth == 0) return s.substring(start + 1, i)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return null
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Разбивает содержимое массива на топ-уровневые `{...}` объекты (string-aware). */
|
||||||
|
internal fun splitTopLevelObjects(arrayContent: String): List<String> {
|
||||||
|
val out = mutableListOf<String>()
|
||||||
|
var i = 0
|
||||||
|
while (i < arrayContent.length) {
|
||||||
|
if (arrayContent[i] == '{') {
|
||||||
|
var depth = 0
|
||||||
|
var inString = false
|
||||||
|
var escape = false
|
||||||
|
var j = i
|
||||||
|
while (j < arrayContent.length) {
|
||||||
|
val c = arrayContent[j]
|
||||||
|
if (escape) { escape = false; j++; continue }
|
||||||
|
if (inString) {
|
||||||
|
when (c) {
|
||||||
|
'\\' -> escape = true
|
||||||
|
'"' -> inString = false
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
when (c) {
|
||||||
|
'"' -> inString = true
|
||||||
|
'{' -> depth++
|
||||||
|
'}' -> {
|
||||||
|
depth--
|
||||||
|
if (depth == 0) {
|
||||||
|
out.add(arrayContent.substring(i, j + 1))
|
||||||
|
i = j + 1
|
||||||
|
break
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
j++
|
||||||
|
}
|
||||||
|
if (j >= arrayContent.length) break
|
||||||
|
} else {
|
||||||
|
i++
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return out
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Достаёт первое строковое поле `"<name>": "..."` из JSON-объекта с
|
||||||
|
* поддержкой escapes (`\"`, `\\`, `\n`, `\t`). `null` если поля нет.
|
||||||
|
*/
|
||||||
|
internal fun extractStringField(obj: String, name: String): String? {
|
||||||
|
val keyRe = Regex(""""$name"\s*:""")
|
||||||
|
val keyMatch = keyRe.find(obj) ?: return null
|
||||||
|
val colonIdx = keyMatch.range.last
|
||||||
|
// Пропускаем пробельные символы после :
|
||||||
|
var i = colonIdx + 1
|
||||||
|
while (i < obj.length && (obj[i] == ' ' || obj[i] == '\n' || obj[i] == '\r' || obj[i] == '\t')) i++
|
||||||
|
if (i >= obj.length || obj[i] != '"') {
|
||||||
|
// Значение не строка (null/число) — не поддерживаем.
|
||||||
|
return null
|
||||||
|
}
|
||||||
|
i++ // открывающая кавычка
|
||||||
|
val sb = StringBuilder()
|
||||||
|
while (i < obj.length) {
|
||||||
|
val c = obj[i]
|
||||||
|
if (c == '\\' && i + 1 < obj.length) {
|
||||||
|
when (val esc = obj[i + 1]) {
|
||||||
|
'"' -> sb.append('"')
|
||||||
|
'\\' -> sb.append('\\')
|
||||||
|
'n' -> sb.append('\n')
|
||||||
|
't' -> sb.append('\t')
|
||||||
|
'r' -> sb.append('\r')
|
||||||
|
else -> sb.append(esc)
|
||||||
|
}
|
||||||
|
i += 2
|
||||||
|
} else if (c == '"') {
|
||||||
|
return sb.toString()
|
||||||
|
} else {
|
||||||
|
sb.append(c)
|
||||||
|
i++
|
||||||
|
}
|
||||||
|
}
|
||||||
|
// Незакрытая строка — мусор от модели, считаем null.
|
||||||
|
return null
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,80 @@
|
|||||||
|
package pw.binom.agentik.standalone.agent
|
||||||
|
|
||||||
|
import pw.binom.agentik.memory.ConversationTurn
|
||||||
|
import pw.binom.agentik.skills.SkillFile
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Промпты для [SkillMiner]: structured-output JSON.
|
||||||
|
*
|
||||||
|
* Тот же подход, что [ReflectionPrompts] и [LlmMemoryReviewer]: локальной модели
|
||||||
|
* (LiteRT-LM) не доверяем tool-calling, поэтому просим строго JSON и парсим руками
|
||||||
|
* ([SkillMiningParser]).
|
||||||
|
*/
|
||||||
|
object SkillMiningPrompts {
|
||||||
|
|
||||||
|
/**
|
||||||
|
* System prompt. Задаёт роль "минёра скилов": модель смотрит на последние
|
||||||
|
* ходы разговора и решает, есть ли в них переиспользуемый приём, который
|
||||||
|
* стоит закрепить в скиле, чтобы агент становился умнее между сессиями.
|
||||||
|
*
|
||||||
|
* Ключевое отличие от `skill_save` (в-ходе): mining — это сетка безопасности,
|
||||||
|
* если модель в ходе разговора забыла сохранить скил. Поэтому промпт жёсткий
|
||||||
|
* по критериям: только реально переиспользуемое, без дублей каталога.
|
||||||
|
*/
|
||||||
|
val SYSTEM_PROMPT: String = """
|
||||||
|
Ты — фоновый минер навыков (skill miner) для ИИ-агента.
|
||||||
|
|
||||||
|
Тебе показывают последние ходы разговора агента с пользователем и каталог
|
||||||
|
его текущих навыков (скилов). Твоя задача — найти в этих ходах переиспользуемый
|
||||||
|
приём, процедуру или паттерн, который агент применит в БУДУЩИХ, других
|
||||||
|
разговорах. Такой приём закрепляется как скил, и агент с ним становится
|
||||||
|
умнее.
|
||||||
|
|
||||||
|
Критерии, ЧТО сохращать:
|
||||||
|
- конкретная воспроизводимая процедура (шаги, команды, форматы, приёмы);
|
||||||
|
- решение проблемы, которое модель выработала и в другом разговоре повторит;
|
||||||
|
- проверка/валидация, о которой модель "забыла" и её стоит закрепить.
|
||||||
|
|
||||||
|
ЧТО НЕ сохраняй:
|
||||||
|
- разовые факты конкретного разговора (это не приём, а данные);
|
||||||
|
- тривиальность ("ответь кратко"), которую и так видно из контекста;
|
||||||
|
- дубли уже существующих скилов в каталоге — если приём уже есть,
|
||||||
|
предложи ОБНОВЛЕНИЕ (то же имя, улучшённый body), а не новый скил.
|
||||||
|
|
||||||
|
Вывод — строго JSON, без пояснений до и после:
|
||||||
|
{"skills": [{"name": "...", "description": "...", "body": "..."}]}
|
||||||
|
Если сохранять нечего — {"skills": []}
|
||||||
|
|
||||||
|
Правила по полям:
|
||||||
|
- name: kebab-case или иерархия через двоеточие (например, "backend:spring:db-base");
|
||||||
|
короткое, 2-5 слов. Если обновляешь существующий скил — ТОЧНО его имя.
|
||||||
|
- description: 1-2 предложения, когда/зачем применять скил.
|
||||||
|
- body: markdown-инструкция: краткое описание + нумерованные шаги + примеры.
|
||||||
|
Только то, что агент должен помнить; без воды.
|
||||||
|
""".trimIndent()
|
||||||
|
|
||||||
|
/**
|
||||||
|
* User prompt: последние [turns] разговора + каталог существующих скилов.
|
||||||
|
* Ходы нумеруются, чтобы модель понимала хронологию.
|
||||||
|
*/
|
||||||
|
fun buildUserPrompt(turns: List<ConversationTurn>, existing: List<SkillFile>): String = buildString {
|
||||||
|
appendLine("### Последние ходы разговора (по хронологии)")
|
||||||
|
turns.forEachIndexed { i, t ->
|
||||||
|
appendLine()
|
||||||
|
appendLine("--- Ход ${i + 1} ---")
|
||||||
|
appendLine("[user] ${t.userMessage.take(1500)}")
|
||||||
|
appendLine("[assistant] ${t.assistantMessage.take(2000)}")
|
||||||
|
}
|
||||||
|
appendLine()
|
||||||
|
appendLine("### Текущий каталог скилов (не дубли, обновляй при необходимости)")
|
||||||
|
if (existing.isEmpty()) {
|
||||||
|
appendLine("(пока нет)")
|
||||||
|
} else {
|
||||||
|
for (s in existing) {
|
||||||
|
appendLine("- ${s.name}: ${s.description.take(160)}")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
appendLine()
|
||||||
|
appendLine("Если есть что сохранить (или обновить существующий) — выведи JSON. Иначе {\"skills\": []}.")
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -96,6 +96,25 @@ data class AgentikConfig(
|
|||||||
* Default: 3. `0` — не подмешивать.
|
* Default: 3. `0` — не подмешивать.
|
||||||
*/
|
*/
|
||||||
val reflectionTopK: Int = DEFAULT_REFLECTION_TOP_K,
|
val reflectionTopK: Int = DEFAULT_REFLECTION_TOP_K,
|
||||||
|
/**
|
||||||
|
* Через сколько пользовательских ходов запускать skill mining (фоновый
|
||||||
|
* LLM-прогон, который находит переиспользуемые скилы, которые модель
|
||||||
|
* забыла сохранить через `skill_save`). `0` — mining выключен. Default: 15.
|
||||||
|
* См. [pw.binom.agentik.standalone.agent.SkillMiner].
|
||||||
|
*/
|
||||||
|
val skillMiningInterval: Int = DEFAULT_SKILL_MINING_INTERVAL,
|
||||||
|
/**
|
||||||
|
* Сколько последних ходов передавать skill-miner'у за один прогон.
|
||||||
|
* Default: 30.
|
||||||
|
*/
|
||||||
|
val skillMiningMaxTurns: Int = DEFAULT_SKILL_MINING_MAX_TURNS,
|
||||||
|
/**
|
||||||
|
* Включает debug-эндпоинты (`/debug/reflect`, `/debug/skill-mine`,
|
||||||
|
* `/debug/curate`, `/debug/compact`, `/debug/tokens`) для ручного
|
||||||
|
* триггерирования фоновых фич без ожидания интервалов. Только локальная
|
||||||
|
* отладка: `AGENTIK_DEBUG_ENDPOINTS=1`.
|
||||||
|
*/
|
||||||
|
val debugEndpoints: Boolean = false,
|
||||||
) {
|
) {
|
||||||
/** Бэкенд долговременной памяти. */
|
/** Бэкенд долговременной памяти. */
|
||||||
@Serializable
|
@Serializable
|
||||||
@@ -113,6 +132,8 @@ data class AgentikConfig(
|
|||||||
const val DEFAULT_EMBEDDING_DIMENSION: Int = 1536
|
const val DEFAULT_EMBEDDING_DIMENSION: Int = 1536
|
||||||
const val DEFAULT_REFLECTION_INTERVAL: Int = 10
|
const val DEFAULT_REFLECTION_INTERVAL: Int = 10
|
||||||
const val DEFAULT_REFLECTION_TOP_K: Int = 3
|
const val DEFAULT_REFLECTION_TOP_K: Int = 3
|
||||||
|
const val DEFAULT_SKILL_MINING_INTERVAL: Int = 15
|
||||||
|
const val DEFAULT_SKILL_MINING_MAX_TURNS: Int = 30
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Читает конфигурацию из переменных среды.
|
* Читает конфигурацию из переменных среды.
|
||||||
@@ -145,6 +166,13 @@ data class AgentikConfig(
|
|||||||
?.coerceIn(0, 1000) ?: DEFAULT_REFLECTION_INTERVAL,
|
?.coerceIn(0, 1000) ?: DEFAULT_REFLECTION_INTERVAL,
|
||||||
reflectionTopK = env("AGENTIK_REFLECTION_TOP_K")?.toIntOrNull()
|
reflectionTopK = env("AGENTIK_REFLECTION_TOP_K")?.toIntOrNull()
|
||||||
?.coerceIn(0, 20) ?: DEFAULT_REFLECTION_TOP_K,
|
?.coerceIn(0, 20) ?: DEFAULT_REFLECTION_TOP_K,
|
||||||
|
skillMiningInterval = env("AGENTIK_SKILL_MINING_INTERVAL")?.toIntOrNull()
|
||||||
|
?.coerceIn(0, 1000) ?: DEFAULT_SKILL_MINING_INTERVAL,
|
||||||
|
skillMiningMaxTurns = env("AGENTIK_SKILL_MINING_MAX_TURNS")?.toIntOrNull()
|
||||||
|
?.coerceIn(1, 1000) ?: DEFAULT_SKILL_MINING_MAX_TURNS,
|
||||||
|
debugEndpoints = env("AGENTIK_DEBUG_ENDPOINTS")?.let {
|
||||||
|
it.equals("1", ignoreCase = true) || it.equals("true", ignoreCase = true)
|
||||||
|
} ?: false,
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -77,8 +77,14 @@ internal class FakeLiteConversation(
|
|||||||
mutableHistory.add(LiteMessage.model(parent.reply))
|
mutableHistory.add(LiteMessage.model(parent.reply))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
override fun send(prompt: String): String = parent.reply
|
override fun send(prompt: String): String {
|
||||||
override fun sendContents(contents: List<LiteContentPart>): String = parent.reply
|
parent.lastContents = listOf(LiteContentPart.Text(prompt))
|
||||||
|
return parent.reply
|
||||||
|
}
|
||||||
|
override fun sendContents(contents: List<LiteContentPart>): String {
|
||||||
|
parent.lastContents = contents
|
||||||
|
return parent.reply
|
||||||
|
}
|
||||||
override fun cancel() {}
|
override fun cancel() {}
|
||||||
override fun tokenCount(): Int = history.size
|
override fun tokenCount(): Int = history.size
|
||||||
override fun addToolResult(callId: String?, name: String, result: String) { error("not used") }
|
override fun addToolResult(callId: String?, name: String, result: String) { error("not used") }
|
||||||
|
|||||||
@@ -0,0 +1,65 @@
|
|||||||
|
package pw.binom.agentik.standalone.agent
|
||||||
|
|
||||||
|
import kotlinx.coroutines.runBlocking
|
||||||
|
import pw.binom.agentik.memory.ConversationTurn
|
||||||
|
import kotlin.test.Test
|
||||||
|
import kotlin.test.assertEquals
|
||||||
|
import kotlin.test.assertTrue
|
||||||
|
|
||||||
|
class SkillMinerTest {
|
||||||
|
|
||||||
|
private fun turns(n: Int): List<ConversationTurn> = (1..n).map {
|
||||||
|
ConversationTurn(userMessage = "q$it", assistantMessage = "a$it")
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun `mine returns parsed skills from model JSON reply`() = runBlocking {
|
||||||
|
val llm = FakeLiteLlm()
|
||||||
|
llm.reply = """{"skills": [{"name": "n", "description": "d", "body": "b"}]}"""
|
||||||
|
val miner = SkillMiner(llm, maxTurns = 10, dispatcher = kotlinx.coroutines.Dispatchers.Unconfined)
|
||||||
|
val out = miner.mine(turns(5), existing = emptyList())
|
||||||
|
assertEquals(1, out.size)
|
||||||
|
assertEquals("n", out[0].name)
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun `mine with empty reply returns empty`() = runBlocking {
|
||||||
|
val llm = FakeLiteLlm()
|
||||||
|
llm.reply = "{}"
|
||||||
|
val miner = SkillMiner(llm, dispatcher = kotlinx.coroutines.Dispatchers.Unconfined)
|
||||||
|
assertTrue(miner.mine(turns(5), emptyList()).isEmpty())
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun `miner failure degrades to empty list`() = runBlocking {
|
||||||
|
val llm = FakeLiteLlm()
|
||||||
|
llm.failMessage = "onnx died"
|
||||||
|
val miner = SkillMiner(llm, dispatcher = kotlinx.coroutines.Dispatchers.Unconfined)
|
||||||
|
assertTrue(miner.mine(turns(5), emptyList()).isEmpty())
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun `miner respects maxTurns`() = runBlocking {
|
||||||
|
val llm = FakeLiteLlm()
|
||||||
|
llm.reply = """{"skills": []}"""
|
||||||
|
val miner = SkillMiner(llm, maxTurns = 2, dispatcher = kotlinx.coroutines.Dispatchers.Unconfined)
|
||||||
|
miner.mine(turns(30), emptyList())
|
||||||
|
val prompt = llm.lastContents!!.first().let {
|
||||||
|
val p = it as pw.binom.litert.LiteContentPart.Text
|
||||||
|
p.text
|
||||||
|
}
|
||||||
|
// В промпт попало только последние 2 хода из 30.
|
||||||
|
assertTrue(prompt.contains("q29"), "last turns missing: $prompt")
|
||||||
|
assertTrue(!prompt.contains("q1\n"), "old turn leaked: $prompt")
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun `empty turns short-circuit without LLM call`() = runBlocking {
|
||||||
|
val llm = FakeLiteLlm()
|
||||||
|
llm.failMessage = "should not be called"
|
||||||
|
val miner = SkillMiner(llm, dispatcher = kotlinx.coroutines.Dispatchers.Unconfined)
|
||||||
|
val out = miner.mine(emptyList(), emptyList())
|
||||||
|
assertTrue(out.isEmpty())
|
||||||
|
assertTrue(llm.conversations.isEmpty(), "LLM must not be called for empty input")
|
||||||
|
}
|
||||||
|
}
|
||||||
+81
@@ -0,0 +1,81 @@
|
|||||||
|
package pw.binom.agentik.standalone.agent
|
||||||
|
|
||||||
|
import kotlin.test.Test
|
||||||
|
import kotlin.test.assertEquals
|
||||||
|
import kotlin.test.assertTrue
|
||||||
|
|
||||||
|
class SkillMiningParserTest {
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun `parses clean JSON`() {
|
||||||
|
val raw = """{"skills": [{"name": "backend:spring:db", "description": "x", "body": "# step 1"}]}"""
|
||||||
|
val out = SkillMiningParser.parse(raw)
|
||||||
|
assertEquals(1, out.size)
|
||||||
|
assertEquals("backend:spring:db", out[0].name)
|
||||||
|
assertEquals("x", out[0].description)
|
||||||
|
assertEquals("# step 1", out[0].body)
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun `parses JSON wrapped in fences with prose around`() {
|
||||||
|
val raw = "Окей, вот что я нашёл:\n```json\n" +
|
||||||
|
"{\"skills\": [{\"name\": \"a\", \"description\": \"d\", \"body\": \"b\"}]}\n" +
|
||||||
|
"```\nНадеюсь, помогло."
|
||||||
|
val out = SkillMiningParser.parse(raw)
|
||||||
|
assertEquals(1, out.size)
|
||||||
|
assertEquals("a", out[0].name)
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun `parses bare array without skills key`() {
|
||||||
|
val raw = """[{"name": "x", "description": "d", "body": "b"}, {"name": "y"}]"""
|
||||||
|
val out = SkillMiningParser.parse(raw)
|
||||||
|
assertEquals(2, out.size)
|
||||||
|
assertEquals("x", out[0].name)
|
||||||
|
assertEquals("", out[1].body)
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun `unescapes newlines and quotes in body`() {
|
||||||
|
val raw = "{\"skills\": [{\"name\": \"n\", \"description\": \"\", \"body\": \"line1\\nline2\\n\\nwith \\\"quotes\\\" and backslash \\\\\\\"\"}]}"
|
||||||
|
val out = SkillMiningParser.parse(raw)
|
||||||
|
assertEquals(1, out.size)
|
||||||
|
val body = out[0].body
|
||||||
|
assertTrue(body.contains("line1\nline2"), "body: $body")
|
||||||
|
assertTrue(body.contains("\"quotes\""), "body: $body")
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun `empty skills array returns empty list`() {
|
||||||
|
val out = SkillMiningParser.parse("""{"skills": []}""")
|
||||||
|
assertTrue(out.isEmpty())
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun `model chatter with no JSON returns empty`() {
|
||||||
|
val out = SkillMiningParser.parse("Скилов не нашёл, всё чисто.")
|
||||||
|
assertTrue(out.isEmpty())
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun `truncated JSON returns empty`() {
|
||||||
|
val out = SkillMiningParser.parse("""{"skills": [{"name": "a", "description": "d", "body": """"")
|
||||||
|
assertTrue(out.isEmpty())
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun `skills without name are dropped`() {
|
||||||
|
val raw = """{"skills": [{"description": "no name"}, {"name": "ok"}]}"""
|
||||||
|
val out = SkillMiningParser.parse(raw)
|
||||||
|
assertEquals(1, out.size)
|
||||||
|
assertEquals("ok", out[0].name)
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun `nested braces inside strings do not break balance`() {
|
||||||
|
val raw = """{"skills": [{"name": "n", "description": "d", "body": "echo '{\"k\": 1}'"}]}"""
|
||||||
|
val out = SkillMiningParser.parse(raw)
|
||||||
|
assertEquals(1, out.size)
|
||||||
|
assertEquals("echo '{\"k\": 1}'", out[0].body)
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user