Refactor: replace BackgroundScheduler with ReflectionScheduler, migrate to outbox-driven event processing, and remove skill mining logic
ci / JVM build + tests (push) Has been cancelled
ci / JVM build + tests (push) Has been cancelled
This commit is contained in:
@@ -106,4 +106,60 @@ sealed interface Event {
|
||||
@Serializable
|
||||
@SerialName("error")
|
||||
data class Error(override val date: Instant, val message: String, val code: String? = null) : Event
|
||||
|
||||
/**
|
||||
* Конвейер вызова тула упал (handler кинул Throwable, args не парсятся,
|
||||
* kernel прибил таск). Отличается от [ToolResult]: там мы сообщаем LLM
|
||||
* результат (даже если LLM его не понравился), тут — сигнал о
|
||||
* внутренней ошибке **самого исполнения тула**. Используется
|
||||
* background-подписчиками (например [ReflectionScheduler]-like
|
||||
* компонентами) для накопления паттернов отказов. Клиенту
|
||||
* показывается для transparency, но в UI особо не нужен.
|
||||
*/
|
||||
@Serializable
|
||||
@SerialName("tool_failed")
|
||||
data class ToolFailed(
|
||||
override val date: Instant,
|
||||
val toolCallId: String,
|
||||
val toolName: String?,
|
||||
val message: String,
|
||||
val durationMs: Long,
|
||||
) : Event
|
||||
|
||||
/**
|
||||
* Диалог переходит в закрытое состояние ([Conversation.close] /
|
||||
* [ConversationLoop.close] / `agent.deleteConversation`). Эмитится
|
||||
* **до** освобождения ресурсов, чтобы background-подписчики
|
||||
* (skill mining, reflection) успели сделать final pass. После
|
||||
* `Closed` диалог уже удалён из `agent.getConversations()` и
|
||||
* `getMessages()` отдаст только то, что осталось в журнале.
|
||||
*
|
||||
* Парный `Opening` намеренно отсутствует — симметрия не нужна,
|
||||
* так как открытие тривиально (id уже известен с момента
|
||||
* `Agent.createConversation` → [Event.ConversationCreated]
|
||||
* / [AgentEvent.Created] в outbox'е).
|
||||
*/
|
||||
@Serializable
|
||||
@SerialName("conversation_closing")
|
||||
data class ConversationClosing(
|
||||
override val date: Instant,
|
||||
val conversationId: String,
|
||||
) : Event
|
||||
|
||||
/**
|
||||
* Compactor сжал старые turn'ы — [turnsCompacted] из истории исчезли.
|
||||
* Эмитится **до** deletion для background-подписчиков (skill mining),
|
||||
* чтобы они успели сделать pass на исчезающем контенте.
|
||||
*
|
||||
* В отличие от [ToolFailed]/[ConversationClosing], это событие
|
||||
* семантически "много контента ушло" — subscribers могут решать,
|
||||
* стоит ли тратить tokens на mining ([turnsCompacted] > N).
|
||||
*/
|
||||
@Serializable
|
||||
@SerialName("compaction_triggered")
|
||||
data class CompactionTriggered(
|
||||
override val date: Instant,
|
||||
val conversationId: String,
|
||||
val turnsCompacted: Int,
|
||||
) : Event
|
||||
}
|
||||
|
||||
+2
-2
@@ -137,6 +137,6 @@ include(":vector-index-ksqlite")
|
||||
// uninstall. KMP, все 9 целей; чистый API-слой, без реализаций.
|
||||
include(":agent-api")
|
||||
// Skill-подсистема как компонент: тулы skill_save/skill_delete/skill_read +
|
||||
// SkillMiningComponent, слушающий BackgroundEventBus и добывающий новые скилы
|
||||
// через SkillMiner на Closing/Compaction.
|
||||
// SkillMiningComponent, слушающий outbox (Event.ConversationClosing /
|
||||
// Event.CompactionTriggered) и добывающий новые скилы через SkillMiner.
|
||||
include(":skill-mining")
|
||||
|
||||
+5
-4
@@ -24,10 +24,11 @@ import pw.binom.litert.LiteTool
|
||||
/**
|
||||
* События, по которым SkillMiningComponent решает, что пора майнить новые скилы.
|
||||
*
|
||||
* Standalone-часть мэпит свой [pw.binom.agentik.standalone.agent.BackgroundEventBus]
|
||||
* на этот sealed interface и подаёт результат в [SkillMiningComponent.events].
|
||||
* Делаем так, чтобы модуль :skill-mining не зависел от :standalone и его
|
||||
* внутренних типов.
|
||||
* Standalone-часть мэпит свой [pw.binom.agentik.outbox.OutboxStore] (через
|
||||
* [pw.binom.agentik.outbox.Event.ConversationClosing] и
|
||||
* [pw.binom.agentik.outbox.Event.CompactionTriggered]) на этот sealed
|
||||
* interface и подаёт результат в [SkillMiningComponent.events]. Делаем так,
|
||||
* чтобы модуль :skill-mining не зависел от :standalone и его внутренних типов.
|
||||
*/
|
||||
sealed interface SkillMiningEvent {
|
||||
val conversationId: String
|
||||
|
||||
@@ -95,9 +95,9 @@ kotlin {
|
||||
// liteTool DSL (типизированные LiteTool через @Serializable args)
|
||||
implementation(libs.litert.tools.kotlinx.serialization)
|
||||
// litert-openai: JVM-реализация
|
||||
implementation(libs.litert.openai)
|
||||
api(libs.litert.openai)
|
||||
// litert-google: встроенный LiteRT-LM движок, нужен только на runtime
|
||||
runtimeOnly(libs.litert.google)
|
||||
api(libs.litert.google)
|
||||
|
||||
// Ktor server (для :server facade + a2aServer)
|
||||
// Используем CIO вместо Netty — он KMP (jvm + linuxX64/ios/...), нам нужен
|
||||
|
||||
@@ -380,7 +380,6 @@ private fun runServer() {
|
||||
contextCompactor = contextCompactor,
|
||||
recentReflections = recentReflections,
|
||||
reflector = reflector,
|
||||
skillMiner = skillMiner,
|
||||
).install(pw.binom.agentik.mcp.bridge.McpBridgeComponent(mcpRegistry))
|
||||
|
||||
// Куратор памяти: фоновая архивация stale-заметок. Поднимается до server'а,
|
||||
|
||||
-84
@@ -1,84 +0,0 @@
|
||||
package pw.binom.agentik.standalone.agent
|
||||
|
||||
import kotlinx.coroutines.channels.BufferOverflow
|
||||
import kotlinx.coroutines.flow.MutableSharedFlow
|
||||
import kotlinx.coroutines.flow.SharedFlow
|
||||
import kotlinx.coroutines.flow.asSharedFlow
|
||||
|
||||
/**
|
||||
* Internal event bus для background work — separate from [ConversationEvents]
|
||||
* (который это SSE-event stream для клиента).
|
||||
*
|
||||
* Background work fires на **structural events**, не на interval-polling:
|
||||
* - **Review** — триггерится в CompactionCoordinator ПРЯМО ПЕРЕД удалением ходов
|
||||
* из working memory (last chance вытащить факты). Это уже было сделано
|
||||
* через `reviewer.reviewPreCompaction()` — оставляем как есть.
|
||||
* - **Reflection** — на `ConversationLifecycleEvent.Closing` (финальная
|
||||
* рефлексия перед закрытием) ИЛИ накопление N tool-failures в окне
|
||||
* (что-то идёт не так).
|
||||
* - **Skill mining** — на `CompactionEvent.Triggered` если `turnsToDelete > N`
|
||||
* (есть контент для минера) ИЛИ на `ConversationLifecycleEvent.Closing`.
|
||||
*
|
||||
* Subscribers (BackgroundScheduler) решают, что делать. НЕ текстовая
|
||||
* инспекция, НЕ regex — только структурные события с явным семантическим
|
||||
* смыслом. См. STANDALONE-REVIEW раздел "event-driven background".
|
||||
*/
|
||||
|
||||
/** Эмитится из [ToolDispatcher] после каждого `runToolAndPersist` (success/failure). */
|
||||
sealed interface ToolCallEvent {
|
||||
val toolName: String
|
||||
|
||||
data class Succeeded(
|
||||
override val toolName: String,
|
||||
val durationMs: Long,
|
||||
) : ToolCallEvent
|
||||
|
||||
data class Failed(
|
||||
override val toolName: String,
|
||||
val error: String,
|
||||
) : ToolCallEvent
|
||||
}
|
||||
|
||||
/** Эмитится из [CompactionCoordinator] ПЕРЕД `workingMemory.compact(...)`. */
|
||||
sealed interface CompactionEvent {
|
||||
/**
|
||||
* Compaction сейчас удалит N ходов из working memory. Background work
|
||||
* имеет последний шанс вытащить оттуда данные.
|
||||
*/
|
||||
data class Triggered(
|
||||
val turnsToDelete: Int,
|
||||
val conversationId: String,
|
||||
) : CompactionEvent
|
||||
}
|
||||
|
||||
/** Эмитится из `ConversationLoop.close()` сразу ПЕРЕД `agentScope.cancel()`. */
|
||||
sealed interface ConversationLifecycleEvent {
|
||||
data class Closing(val conversationId: String) : ConversationLifecycleEvent
|
||||
}
|
||||
|
||||
internal class BackgroundEventBus {
|
||||
private val _toolCallEvents = MutableSharedFlow<ToolCallEvent>(
|
||||
replay = 0,
|
||||
extraBufferCapacity = 256,
|
||||
onBufferOverflow = BufferOverflow.DROP_OLDEST,
|
||||
)
|
||||
val toolCallEvents: SharedFlow<ToolCallEvent> get() = _toolCallEvents.asSharedFlow()
|
||||
|
||||
private val _compactionEvents = MutableSharedFlow<CompactionEvent>(
|
||||
replay = 0,
|
||||
extraBufferCapacity = 16,
|
||||
onBufferOverflow = BufferOverflow.DROP_OLDEST,
|
||||
)
|
||||
val compactionEvents: SharedFlow<CompactionEvent> get() = _compactionEvents.asSharedFlow()
|
||||
|
||||
private val _lifecycleEvents = MutableSharedFlow<ConversationLifecycleEvent>(
|
||||
replay = 0,
|
||||
extraBufferCapacity = 16,
|
||||
onBufferOverflow = BufferOverflow.DROP_OLDEST,
|
||||
)
|
||||
val lifecycleEvents: SharedFlow<ConversationLifecycleEvent> get() = _lifecycleEvents.asSharedFlow()
|
||||
|
||||
fun tryEmit(event: ToolCallEvent): Boolean = _toolCallEvents.tryEmit(event)
|
||||
fun tryEmit(event: CompactionEvent): Boolean = _compactionEvents.tryEmit(event)
|
||||
fun tryEmit(event: ConversationLifecycleEvent): Boolean = _lifecycleEvents.tryEmit(event)
|
||||
}
|
||||
@@ -6,6 +6,7 @@ import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.SupervisorJob
|
||||
import kotlinx.coroutines.flow.filterIsInstance
|
||||
import kotlinx.coroutines.flow.map
|
||||
import kotlinx.coroutines.flow.mapNotNull
|
||||
import kotlinx.coroutines.flow.merge
|
||||
import kotlinx.coroutines.runBlocking
|
||||
import kotlinx.coroutines.sync.Mutex
|
||||
@@ -19,12 +20,12 @@ import pw.binom.agentik.agent.SystemPromptProvider
|
||||
import pw.binom.agentik.agent.ToolProvider
|
||||
import pw.binom.agentik.llm.tools.ContextCompactor
|
||||
import pw.binom.agentik.llm.tools.LlmReflector
|
||||
import pw.binom.agentik.skill.mining.SkillMiner
|
||||
import pw.binom.agentik.memory.MemoryPrefetcher
|
||||
import pw.binom.agentik.memory.MemoryReviewer
|
||||
import pw.binom.agentik.memory.MemorySystemGuidance
|
||||
import pw.binom.agentik.proto.Agent as ProtoAgent
|
||||
import pw.binom.agentik.outbox.AgentEvent
|
||||
import pw.binom.agentik.outbox.Event as OutboxEvent
|
||||
import pw.binom.agentik.outbox.CommonEvent
|
||||
import pw.binom.agentik.outbox.MutableOutboxStore
|
||||
import pw.binom.agentik.journal.JournalStore
|
||||
@@ -98,7 +99,6 @@ internal class ChatAgent(
|
||||
private val skillStore: pw.binom.agentik.skills.SkillStore? = null,
|
||||
private val memoryStore: pw.binom.agentik.memory.MemoryStore? = null,
|
||||
private val memoryPrefetcher: MemoryPrefetcher? = null,
|
||||
internal val _backgroundEvents: BackgroundEventBus = BackgroundEventBus(),
|
||||
internal val agentScope: CoroutineScope = CoroutineScope(SupervisorJob() + Dispatchers.Default),
|
||||
private val memoryReviewer: MemoryReviewer? = null,
|
||||
/**
|
||||
@@ -129,12 +129,6 @@ internal class ChatAgent(
|
||||
* Исполнитель рефлексий (one-shot LiteLlm вызов). `null` = self-reflection выключен.
|
||||
*/
|
||||
private val reflector: LlmReflector? = null,
|
||||
/**
|
||||
* Фоновый минер скилов: LLM-вызов, который запускается на compaction
|
||||
* (`turnsToDelete > 10`) или при closing conversation. `null` = mining выключен.
|
||||
* Сетка безопасности, если модель забыла вызвать `skill_save` сама.
|
||||
*/
|
||||
private val skillMiner: SkillMiner? = null,
|
||||
/**
|
||||
* Тулсеты, доступные агенту. Пустой список (по умолчанию) — модель не знает
|
||||
* о механике toolsets: enable_toolset/disable_toolset НЕ регистрируются,
|
||||
@@ -184,12 +178,30 @@ internal class ChatAgent(
|
||||
reflections = recentReflections,
|
||||
)
|
||||
|
||||
/**
|
||||
/**
|
||||
* Тестовый инжектор: позволяет `ChatAgentTest` подсунуть тул, который
|
||||
* скриптованная LLM будет вызывать. В production-конфигурации всегда пуст —
|
||||
* инструменты приходят через [toolProviders] (компоненты, [install]).
|
||||
*/
|
||||
private val testTools: MutableList<LiteTool> = mutableListOf()
|
||||
private val testTools: MutableList<LiteTool> = mutableListOf()
|
||||
|
||||
/**
|
||||
* Единый канал всех событий агента — bounded tail с auto-TTL.
|
||||
*
|
||||
* Заменил ранее существовавшие два канала:
|
||||
* - `agentEvents: MutableSharedFlow<AgentEvent>` (agent lifecycle)
|
||||
* - per-conv `ConversationEvents._flow: MutableSharedFlow<ProtoEvent>`
|
||||
*
|
||||
* Теперь оба пишут сюда через [MutableOutboxStore.append], а consumer'ы
|
||||
* читают через [EventStore.events]/[conversationEvents]/[agentEvents].
|
||||
*
|
||||
* **Live tail + auto-TTL** — клиенты больше не должны заботиться о persistence
|
||||
* или подписке на два отдельных канала.
|
||||
*/
|
||||
private val eventStore: MutableOutboxStore = pw.binom.agentik.outbox.inmemory.InMemoryOutboxStore(
|
||||
maxMessages = null,
|
||||
ttl = null,
|
||||
)
|
||||
|
||||
/**
|
||||
* Собирает **актуальный** список тулов для диспетчеризации:
|
||||
@@ -265,18 +277,22 @@ internal class ChatAgent(
|
||||
private val systemPrompt: String = baseSystemPrompt
|
||||
|
||||
/**
|
||||
* Runtime system prompt = база + динамические секции от `systemProviders`.
|
||||
* Вызывается на каждом turn'е (через [systemPromptResolver]). Пустые секции
|
||||
* пропускаются; компонент, у которого `getSection` вернул пустую строку,
|
||||
* не появляется в промте.
|
||||
* Подписка на skill-mining через outbox (единый канал всех событий):
|
||||
* - [OutboxEvent.ConversationClosing] → [SkillMiningEvent.ConversationClosing]
|
||||
* - [OutboxEvent.CompactionTriggered] → [SkillMiningEvent.ConversationCompacted]
|
||||
*
|
||||
* Оба маппатся в один [Flow], `SkillMiningComponent` сам решает что делать.
|
||||
*/
|
||||
private fun skillMiningEvents(): kotlinx.coroutines.flow.Flow<SkillMiningEvent> =
|
||||
merge(
|
||||
_backgroundEvents.lifecycleEvents.filterIsInstance<ConversationLifecycleEvent.Closing>()
|
||||
.map { SkillMiningEvent.ConversationClosing(it.conversationId) },
|
||||
_backgroundEvents.compactionEvents.filterIsInstance<CompactionEvent.Triggered>()
|
||||
.map { SkillMiningEvent.ConversationCompacted(it.conversationId) },
|
||||
)
|
||||
eventStore.events(after = null)
|
||||
.filterIsInstance<CommonEvent.Conversation>()
|
||||
.mapNotNull { ce ->
|
||||
when (val e = ce.event) {
|
||||
is OutboxEvent.ConversationClosing -> SkillMiningEvent.ConversationClosing(ce.conversationId)
|
||||
is OutboxEvent.CompactionTriggered -> SkillMiningEvent.ConversationCompacted(ce.conversationId)
|
||||
else -> null
|
||||
}
|
||||
}
|
||||
|
||||
internal fun buildRuntimeSystemPrompt(): String {
|
||||
// Секции от systemProviders идут ПЕРЕД базой — это позволяет компонентам
|
||||
@@ -326,7 +342,8 @@ internal class ChatAgent(
|
||||
* чтобы они не лезли в чужие обязанности.
|
||||
*
|
||||
* [recentTurns] берёт последние [limit] реплик из working memory —
|
||||
* это всё, что нужно [SkillMiner] / [LlmReflector] для добычи знаний.
|
||||
* это всё, что нужно [pw.binom.agentik.skill.mining.SkillMiningComponent] /
|
||||
* [LlmReflector] для добычи знаний.
|
||||
*/
|
||||
private fun ChatConversation.asHandle(): ConversationHandle = object : ConversationHandle {
|
||||
override val id: String get() = this@asHandle.id
|
||||
@@ -365,24 +382,6 @@ internal class ChatAgent(
|
||||
override fun close() = Unit
|
||||
}
|
||||
|
||||
/**
|
||||
* Единый канал всех событий агента — bounded tail с auto-TTL.
|
||||
*
|
||||
* Заменил ранее существовавшие два канала:
|
||||
* - `agentEvents: MutableSharedFlow<AgentEvent>` (agent lifecycle)
|
||||
* - per-conv `ConversationEvents._flow: MutableSharedFlow<ProtoEvent>`
|
||||
*
|
||||
* Теперь оба пишут сюда через [MutableOutboxStore.append], а consumer'ы
|
||||
* читают через [EventStore.events]/[conversationEvents]/[agentEvents].
|
||||
*
|
||||
* **Live tail + auto-TTL** — клиенты больше не должны заботиться о persistence
|
||||
* или подписке на два отдельных канала.
|
||||
*/
|
||||
private val eventStore: MutableOutboxStore = pw.binom.agentik.outbox.inmemory.InMemoryOutboxStore(
|
||||
maxMessages = null,
|
||||
ttl = null,
|
||||
)
|
||||
|
||||
/**
|
||||
* Read-only view of [messageStore] для HTTP-фасада в `:server`
|
||||
* (`Route.agentikAgent` → `/journal/...` endpoint'ы).
|
||||
@@ -460,8 +459,6 @@ internal class ChatAgent(
|
||||
compressionThreshold = compressionThreshold,
|
||||
contextCompactor = contextCompactor,
|
||||
reflector = reflector,
|
||||
skillMiner = skillMiner,
|
||||
skillMiningStore = skillStore,
|
||||
)
|
||||
runBlocking {
|
||||
liveLock.withLock { live[conv.id] = conv }
|
||||
@@ -534,8 +531,6 @@ internal class ChatAgent(
|
||||
compressionThreshold = compressionThreshold,
|
||||
contextCompactor = contextCompactor,
|
||||
reflector = reflector,
|
||||
skillMiner = skillMiner,
|
||||
skillMiningStore = skillStore,
|
||||
)
|
||||
|
||||
override fun close() {
|
||||
|
||||
+1
-1
@@ -3,7 +3,7 @@ package pw.binom.agentik.standalone.agent
|
||||
/**
|
||||
* Историческое имя класса. До 2025-Q4 разбиения god-class а на компоненты
|
||||
* (`ConversationState`, `ConversationEvents`, `ContextBuilder`,
|
||||
* `CompactionCoordinator`, `ToolDispatcher`, `BackgroundScheduler`) вся
|
||||
* `CompactionCoordinator`, `ToolDispatcher`, `ReflectionScheduler`) вся
|
||||
* логика жила в `class ChatConversation` здесь, ~1400 строк.
|
||||
*
|
||||
* После рефактора — реализация переехала в [ConversationLoop] (этот же пакет).
|
||||
|
||||
+5
-7
@@ -18,6 +18,7 @@ import pw.binom.litert.LiteMessage
|
||||
import pw.binom.litert.LiteRole
|
||||
import pw.binom.agentik.llm.tools.ContextCompactor
|
||||
import pw.binom.agentik.llm.tools.SummaryTurn
|
||||
import kotlin.time.Instant
|
||||
|
||||
internal class CompactionCoordinator(
|
||||
private val state: ConversationState,
|
||||
@@ -28,8 +29,9 @@ internal class CompactionCoordinator(
|
||||
private val memoryStoreForReview: MemoryStore?,
|
||||
private val workingMemory: ContextStore,
|
||||
private val liteLlm: LiteLlm,
|
||||
private val events: ConversationEvents,
|
||||
private val now: () -> Instant,
|
||||
private val systemPrompt: String,
|
||||
private val backgroundEvents: BackgroundEventBus,
|
||||
) {
|
||||
private val log = KotlinLogging.logger {}
|
||||
|
||||
@@ -66,11 +68,7 @@ internal class CompactionCoordinator(
|
||||
}
|
||||
if (toCompact.isEmpty()) return false
|
||||
|
||||
// Emit CompactionEvent BEFORE deletion. BackgroundScheduler может trigger'ить
|
||||
// skill mining на основе turnsToDelete (есть контент — есть что майнить).
|
||||
backgroundEvents.tryEmit(
|
||||
CompactionEvent.Triggered(turnsToDelete = toCompact.size, conversationId = state.id),
|
||||
)
|
||||
events.tryEmit(pw.binom.agentik.outbox.Event.CompactionTriggered(date = now(), conversationId = state.id, turnsCompacted = toCompact.size))
|
||||
|
||||
val turns = toCompact.mapNotNull { row ->
|
||||
when (val e = row.entry) {
|
||||
@@ -252,7 +250,7 @@ internal class CompactionCoordinator(
|
||||
companion object {
|
||||
private const val KEEP_RECENT_TURNS = 4
|
||||
private const val MAX_SEEDED_MESSAGES = 50
|
||||
/** Минимум удаляемых ходов чтобы BackgroundScheduler trigger'ил skill mining. */
|
||||
/** Минимум удаляемых ходов чтобы SkillMiningComponent trigger'ил skill mining. */
|
||||
private const val MIN_COMPACTION_FOR_MINING = 10
|
||||
}
|
||||
}
|
||||
|
||||
+15
-23
@@ -25,7 +25,6 @@ import pw.binom.agentik.outbox.Event as ProtoEvent
|
||||
import pw.binom.agentik.reflection.ReflectionStore
|
||||
import pw.binom.agentik.proto.Message as ProtoMessage
|
||||
import pw.binom.agentik.proto.MessageContext as ProtoMessageContext
|
||||
import pw.binom.agentik.skills.SkillStore
|
||||
import pw.binom.agentik.journal.Content
|
||||
import pw.binom.agentik.journal.ConversationRecord
|
||||
import pw.binom.agentik.journal.MutableConversationStore
|
||||
@@ -45,7 +44,6 @@ import pw.binom.litert.LiteToolCall
|
||||
import java.util.concurrent.atomic.AtomicBoolean
|
||||
import kotlin.time.Instant
|
||||
import pw.binom.agentik.llm.tools.LlmReflector
|
||||
import pw.binom.agentik.skill.mining.SkillMiner
|
||||
import pw.binom.agentik.llm.tools.ContextCompactor
|
||||
|
||||
class ConversationLoop(
|
||||
@@ -88,8 +86,6 @@ class ConversationLoop(
|
||||
private val compressionThreshold: Double = 0.8,
|
||||
private val contextCompactor: ContextCompactor? = null,
|
||||
private val reflector: LlmReflector? = null,
|
||||
private val skillMiner: SkillMiner? = null,
|
||||
private val skillMiningStore: SkillStore? = null,
|
||||
) : ProtoConversation, AutoCloseable {
|
||||
|
||||
private val log = KotlinLogging.logger {}
|
||||
@@ -109,9 +105,6 @@ class ConversationLoop(
|
||||
conversationId = state.id,
|
||||
)
|
||||
|
||||
/** Per-conversation background event bus. Lifecycle scoped к этому ConversationLoop. */
|
||||
private val backgroundEvents = BackgroundEventBus()
|
||||
|
||||
private val toolsByName: MutableMap<String, LiteTool> = tools.associateBy { it.name }.toMutableMap()
|
||||
|
||||
private val contextBuilder = ContextBuilder(memoryPrefetcher = memoryPrefetcher)
|
||||
@@ -126,32 +119,29 @@ class ConversationLoop(
|
||||
workingMemory = workingMemoryStore,
|
||||
liteLlm = llm,
|
||||
systemPrompt = systemPrompt,
|
||||
backgroundEvents = backgroundEvents,
|
||||
events = events,
|
||||
now = ::now,
|
||||
)
|
||||
|
||||
private val toolDispatcher = ToolDispatcher(
|
||||
state = state,
|
||||
messageStore = messageStore,
|
||||
events = events,
|
||||
backgroundEvents = backgroundEvents,
|
||||
toolsByName = toolsByName,
|
||||
toolsetDispatch = toolsetDispatch,
|
||||
newId = ::newId,
|
||||
now = ::now,
|
||||
)
|
||||
|
||||
private val backgroundScheduler = BackgroundScheduler(
|
||||
private val reflectionScheduler = ReflectionScheduler(
|
||||
state = state,
|
||||
workingMemory = workingMemoryStore,
|
||||
config = BackgroundConfig(
|
||||
memoryReviewer = memoryReviewer,
|
||||
memoryStore = memoryStoreForReview,
|
||||
reflectionStore = reflectionStore,
|
||||
config = ReflectionConfig(
|
||||
reflector = reflector,
|
||||
skillMiner = skillMiner,
|
||||
skillMiningStore = skillMiningStore,
|
||||
reflectionStore = reflectionStore,
|
||||
),
|
||||
backgroundEvents = backgroundEvents,
|
||||
eventStore = eventStore,
|
||||
conversationIdProvider = { id },
|
||||
).also { it.start(agentScope) }
|
||||
|
||||
override val id: String get() = state.id
|
||||
@@ -244,10 +234,12 @@ class ConversationLoop(
|
||||
override fun close() {
|
||||
if (state.isClosed) return
|
||||
state.markClosed()
|
||||
// Emit Closing event BEFORE agentScope.cancel() — BackgroundScheduler's подписка
|
||||
// ловит это и делает final reflection + skill mining (last chance вытащить insights).
|
||||
// Emit Closing event BEFORE agentScope.cancel() — ReflectionScheduler's подписка
|
||||
// ловит это и делает final reflection (last chance вытащить insights).
|
||||
// После cancel() подписка умерла бы.
|
||||
backgroundEvents.tryEmit(ConversationLifecycleEvent.Closing(conversationId = id))
|
||||
runCatching {
|
||||
events.tryEmit(pw.binom.agentik.outbox.Event.ConversationClosing(date = now(), conversationId = id))
|
||||
}
|
||||
state.liteConvRef.getAndSet(null)?.let { runCatching { it.close() } }
|
||||
runCatching { runBlocking { activeTurn?.cancelAndJoin() } }
|
||||
agentScope.cancel()
|
||||
@@ -474,9 +466,9 @@ class ConversationLoop(
|
||||
)
|
||||
}
|
||||
|
||||
// BackgroundScheduler is event-driven — подписан на BackgroundEventBus
|
||||
// (compaction/lifecycle/tool-failure events). Никаких interval-based
|
||||
// вызовов сюда больше не идёт. См. BackgroundScheduler.kt.
|
||||
// ReflectionScheduler is event-driven — подписан на outbox
|
||||
// (Event.ConversationClosing / Event.ToolFailed). Никаких
|
||||
// interval-based вызовов сюда больше не идёт. См. ReflectionScheduler.kt.
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
+31
-88
@@ -2,57 +2,57 @@ package pw.binom.agentik.standalone.agent
|
||||
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.Job
|
||||
import kotlinx.coroutines.flow.filter
|
||||
import kotlinx.coroutines.flow.filterIsInstance
|
||||
import kotlinx.coroutines.flow.launchIn
|
||||
import kotlinx.coroutines.flow.merge
|
||||
import kotlinx.coroutines.flow.onEach
|
||||
import kotlinx.coroutines.launch
|
||||
import mu.KotlinLogging
|
||||
import pw.binom.agentik.memory.ConversationTurn
|
||||
import pw.binom.agentik.memory.MemoryReviewer
|
||||
import pw.binom.agentik.memory.MemoryStore
|
||||
import pw.binom.agentik.skills.SkillStore
|
||||
import pw.binom.agentik.journal.Content
|
||||
import pw.binom.agentik.reflection.ReflectionStore
|
||||
import pw.binom.agentik.context.WorkingMemoryEntry
|
||||
import pw.binom.agentik.context.ContextStore
|
||||
import pw.binom.agentik.llm.tools.LlmReflector
|
||||
import pw.binom.agentik.skill.mining.SkillMiner
|
||||
import pw.binom.agentik.outbox.CommonEvent
|
||||
import pw.binom.agentik.outbox.Event
|
||||
import pw.binom.agentik.outbox.OutboxStore
|
||||
import java.util.concurrent.atomic.AtomicLong
|
||||
|
||||
/**
|
||||
* Background work подписчик на [BackgroundEventBus]. Заменяет старую interval-based
|
||||
* логику (`maybeScheduleReview/Reflection/SkillMining` с `userTurnCount % N == 0`).
|
||||
* Event-driven reflection trigger. Подписан на [OutboxStore] и пишет
|
||||
* в [ReflectionStore] когда:
|
||||
*
|
||||
* Подписки:
|
||||
* - [ConversationLifecycleEvent.Closing] → финальный reflection + skill mining
|
||||
* перед закрытием conversation (last chance вытащить insights).
|
||||
* - [CompactionEvent.Triggered] → skill mining если `turnsToDelete > MIN_COMPACTION_FOR_MINING`.
|
||||
* Review уже сделан внутри CompactionCoordinator (`reviewPreCompaction`) — не дублируем.
|
||||
* - [ToolCallEvent.Failed] (накопительно) → reflection если 2+ фейлов в окне 60 сек
|
||||
* (что-то пошло не так — самоанализ полезен).
|
||||
* - `Event.ConversationClosing` для нашего `conversationId` — финальная
|
||||
* рефлексия перед закрытием. Skill mining на Closing уже делает
|
||||
* [pw.binom.agentik.skill.mining.SkillMiningComponent].
|
||||
* - `Event.ToolFailed` для нашего `conversationId` — накопительно:
|
||||
* 2+ фейлов за 60 сек запускают рефлексию (что-то идёт не так —
|
||||
* самоанализ полезен).
|
||||
*
|
||||
* НЕ текстовая инспекция, НЕ regex, НЕ interval-polling. Только структурные
|
||||
* события с явным семантическим смыслом.
|
||||
*
|
||||
* Подписка идёт через outbox (а не через per-conversation BackgroundEventBus
|
||||
* который был раньше): outbox — единый канал для всех событий (как клиентских,
|
||||
* так и внутренних), persistent tail с TTL работает из коробки, а клиенты
|
||||
* по тому же потоку могут самостоятельно видеть/логировать [Event.ToolFailed]
|
||||
* без скрытой телеметрии.
|
||||
*/
|
||||
internal data class BackgroundConfig(
|
||||
val memoryReviewer: MemoryReviewer?,
|
||||
val memoryStore: MemoryStore?,
|
||||
val reflectionStore: ReflectionStore?,
|
||||
internal data class ReflectionConfig(
|
||||
val reflector: LlmReflector?,
|
||||
val skillMiner: SkillMiner?,
|
||||
val skillMiningStore: SkillStore?,
|
||||
val reflectionStore: ReflectionStore?,
|
||||
)
|
||||
|
||||
internal class BackgroundScheduler(
|
||||
internal class ReflectionScheduler(
|
||||
private val state: ConversationState,
|
||||
private val workingMemory: ContextStore,
|
||||
private val config: BackgroundConfig,
|
||||
private val backgroundEvents: BackgroundEventBus,
|
||||
private val config: ReflectionConfig,
|
||||
private val eventStore: OutboxStore,
|
||||
private val conversationIdProvider: () -> String,
|
||||
) {
|
||||
private val log = KotlinLogging.logger {}
|
||||
|
||||
private val lastSkillMiningAt = AtomicLong(0)
|
||||
private val lastReflectionAt = AtomicLong(0)
|
||||
|
||||
/** Recent tool failure timestamps (ms). Trimmed to [FAILURE_WINDOW_MS]. */
|
||||
@@ -65,17 +65,14 @@ internal class BackgroundScheduler(
|
||||
fun start(scope: CoroutineScope) {
|
||||
if (subscriptionJob?.isActive == true) return
|
||||
subscriptionJob = scope.launch {
|
||||
// Merge all three event flows into one subscription scope. Each onEach
|
||||
// returns Unit, so launchIn merges them as cold flows.
|
||||
merge(
|
||||
backgroundEvents.lifecycleEvents
|
||||
.filterIsInstance<ConversationLifecycleEvent.Closing>()
|
||||
eventStore.events(after = null)
|
||||
.filterIsInstance<CommonEvent.Conversation>()
|
||||
.filter { it.event is Event.ConversationClosing && it.conversationId == conversationIdProvider() }
|
||||
.onEach { onClosing() },
|
||||
backgroundEvents.compactionEvents
|
||||
.filterIsInstance<CompactionEvent.Triggered>()
|
||||
.onEach { onCompaction(it) },
|
||||
backgroundEvents.toolCallEvents
|
||||
.filterIsInstance<ToolCallEvent.Failed>()
|
||||
eventStore.events(after = null)
|
||||
.filterIsInstance<CommonEvent.Conversation>()
|
||||
.filter { it.event is Event.ToolFailed && it.conversationId == conversationIdProvider() }
|
||||
.onEach { onToolFailure() },
|
||||
).collect {}
|
||||
}
|
||||
@@ -83,9 +80,7 @@ internal class BackgroundScheduler(
|
||||
|
||||
private suspend fun onClosing() {
|
||||
if (state.isTemporal) return
|
||||
val convId = state.id
|
||||
runFinalReflection(convId)
|
||||
runFinalSkillMining(convId)
|
||||
runFinalReflection(state.id)
|
||||
}
|
||||
|
||||
private fun runFinalReflection(convId: String) {
|
||||
@@ -105,51 +100,6 @@ internal class BackgroundScheduler(
|
||||
}
|
||||
}
|
||||
|
||||
private fun runFinalSkillMining(convId: String) {
|
||||
val miner = config.skillMiner ?: return
|
||||
val store = config.skillMiningStore ?: return
|
||||
state.agentScope.launch {
|
||||
try {
|
||||
val turns = recentTurns(miner.maxTurns)
|
||||
if (turns.isEmpty()) return@launch
|
||||
val mined = miner.mine(turns, store.catalog.skills)
|
||||
for (s in mined) {
|
||||
runCatching { store.upsert(s) }
|
||||
.onFailure { log.warn(it) { "final skill mining upsert '${s.name}' failed: ${it.message}" } }
|
||||
}
|
||||
log.info { "final skill mining on close: conv=$convId turns=${turns.size} existing=${store.catalog.skills.size} mined=${mined.size}" }
|
||||
} catch (e: Throwable) {
|
||||
log.warn(e) { "final skill mining failed for $convId: ${e.message}" }
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private fun onCompaction(event: CompactionEvent.Triggered) {
|
||||
if (state.isTemporal) return
|
||||
if (event.turnsToDelete < MIN_COMPACTION_FOR_MINING) return
|
||||
val miner = config.skillMiner ?: return
|
||||
val store = config.skillMiningStore ?: return
|
||||
val now = System.currentTimeMillis()
|
||||
// Debounce: не чаще раза в минуту
|
||||
if (now - lastSkillMiningAt.get() < MINING_DEBOUNCE_MS) return
|
||||
lastSkillMiningAt.set(now)
|
||||
val convId = event.conversationId
|
||||
state.agentScope.launch {
|
||||
try {
|
||||
val turns = recentTurns(miner.maxTurns)
|
||||
if (turns.isEmpty()) return@launch
|
||||
val mined = miner.mine(turns, store.catalog.skills)
|
||||
for (s in mined) {
|
||||
runCatching { store.upsert(s) }
|
||||
.onFailure { log.warn(it) { "compaction skill mining upsert '${s.name}' failed: ${it.message}" } }
|
||||
}
|
||||
log.info { "compaction skill mining: conv=$convId turnsToDelete=${event.turnsToDelete} mined=${mined.size}" }
|
||||
} catch (e: Throwable) {
|
||||
log.warn(e) { "compaction skill mining failed for $convId: ${e.message}" }
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private fun onToolFailure() {
|
||||
val reflector = config.reflector ?: return
|
||||
if (state.isTemporal) return
|
||||
@@ -158,21 +108,18 @@ internal class BackgroundScheduler(
|
||||
val now = System.currentTimeMillis()
|
||||
val shouldReflect = synchronized(toolFailuresLock) {
|
||||
toolFailures.add(now)
|
||||
// Trim old failures outside the window
|
||||
val cutoff = now - FAILURE_WINDOW_MS
|
||||
toolFailures.removeAll { it < cutoff }
|
||||
toolFailures.size >= FAILURE_THRESHOLD
|
||||
}
|
||||
if (!shouldReflect) return
|
||||
|
||||
// Debounce reflection globally (не чаще раза в 5 мин)
|
||||
if (now - lastReflectionAt.get() < REFLECTION_DEBOUNCE_MS) {
|
||||
log.debug { "reflection debounced: ${toolFailures.size} failures accumulated but reflection fired recently" }
|
||||
return
|
||||
}
|
||||
lastReflectionAt.set(now)
|
||||
|
||||
// Clear failure window — fresh accounting period
|
||||
synchronized(toolFailuresLock) { toolFailures.clear() }
|
||||
|
||||
val convId = state.id
|
||||
@@ -212,10 +159,6 @@ internal class BackgroundScheduler(
|
||||
filterIsInstance<Content.Text>().joinToString("\n") { it.body }
|
||||
|
||||
companion object {
|
||||
/** Минимум ходов, удаляемых compaction'ом, чтобы trigger'ить skill mining. */
|
||||
private const val MIN_COMPACTION_FOR_MINING = 10
|
||||
/** Дебаунс skill mining между запусками. */
|
||||
private const val MINING_DEBOUNCE_MS = 60_000L
|
||||
/** Debounce reflection между запусками (накопительный, не per-failure). */
|
||||
private const val REFLECTION_DEBOUNCE_MS = 300_000L
|
||||
/** Сколько tool-failures в окне должно накопиться чтобы trigger reflection. */
|
||||
+12
-6
@@ -17,7 +17,6 @@ internal class ToolDispatcher(
|
||||
private val state: ConversationState,
|
||||
private val messageStore: MutableJournalStore,
|
||||
private val events: ConversationEvents,
|
||||
private val backgroundEvents: BackgroundEventBus,
|
||||
private val toolsByName: MutableMap<String, LiteTool>,
|
||||
private val toolsetDispatch: ToolsetDispatchPolicy?,
|
||||
private val newId: (String) -> String,
|
||||
@@ -90,13 +89,20 @@ internal class ToolDispatcher(
|
||||
val resultAt = now()
|
||||
events.tryEmit(ProtoEvent.ToolResult(date = resultAt, toolCallId = callId, toolName = call.name, result = resultText))
|
||||
|
||||
// Эмитим background event — другие компоненты (BackgroundScheduler)
|
||||
// решают, делать ли что-то. Cancellation = not a failure (не эмитим Failed).
|
||||
// Только реальные падения тула попадают в outbox как background-event
|
||||
// (reflection-подобные потребители). Cancellation — not a failure,
|
||||
// успех — тоже без фонового события (ToolResult уже всё сказал).
|
||||
val durationMs = System.currentTimeMillis() - startMs
|
||||
if (failureError != null) {
|
||||
backgroundEvents.tryEmit(ToolCallEvent.Failed(toolName = call.name, error = failureError))
|
||||
} else if (resultText != "[cancelled by user]") {
|
||||
backgroundEvents.tryEmit(ToolCallEvent.Succeeded(toolName = call.name, durationMs = durationMs))
|
||||
events.tryEmit(
|
||||
ProtoEvent.ToolFailed(
|
||||
date = resultAt,
|
||||
toolCallId = callId,
|
||||
toolName = call.name,
|
||||
message = failureError,
|
||||
durationMs = durationMs,
|
||||
),
|
||||
)
|
||||
}
|
||||
|
||||
if (!state.isTemporal) {
|
||||
|
||||
@@ -1,46 +0,0 @@
|
||||
package pw.binom.agentik.standalone.llm
|
||||
|
||||
import pw.binom.litert.LiteConfig
|
||||
import pw.binom.litert.LiteLlm
|
||||
|
||||
/**
|
||||
* Factory для движка Google LiteRT-LM на JVM.
|
||||
*
|
||||
* `litert-google` опубликован только как Android AAR (с .so внутри), поэтому на JVM
|
||||
* приходится вручную подгружать JNI-библиотеки LiteRT-LM, прежде чем инстанциировать
|
||||
* движок. Эта функция:
|
||||
*
|
||||
* 1. Резолвит `pw.binom.litert.google.GoogleLiteLlm` через reflection.
|
||||
* 2. Если процесс уже загрузил нативные библиотеки (`-Djava.library.path`) — успешно
|
||||
* создаёт движок.
|
||||
* 3. Если нет — кидает `IllegalStateException` с инструкцией по настройке.
|
||||
*/
|
||||
fun googleLiteLlmJvm(config: LiteConfig): LiteLlm {
|
||||
val cls = try {
|
||||
Class.forName("pw.binom.litert.google.GoogleLiteLlm")
|
||||
} catch (e: ClassNotFoundException) {
|
||||
error(
|
||||
"litert-google classes are not on the classpath. " +
|
||||
"Add 'pw.binom.litert:litert-google-android:6' as a runtime dependency " +
|
||||
"to use the GOOGLE backend on JVM."
|
||||
)
|
||||
}
|
||||
val ctor = cls.constructors.firstOrNull { it.parameterCount == 1 }
|
||||
?: error("pw.binom.litert.google.GoogleLiteLlm constructor not found")
|
||||
val engine = try {
|
||||
ctor.newInstance(config)
|
||||
} catch (e: UnsatisfiedLinkError) {
|
||||
throw IllegalStateException(
|
||||
"LiteRT-LM native libraries are not loaded. " +
|
||||
"Extract .so/.dylib/.dll from litertlm-android-0.16.1.aar and pass " +
|
||||
"-Djava.library.path=<dir>, or build :standalone for the androidJvm target.",
|
||||
e,
|
||||
)
|
||||
} catch (e: java.lang.reflect.InvocationTargetException) {
|
||||
throw e.targetException ?: e
|
||||
}
|
||||
@Suppress("UNCHECKED_CAST")
|
||||
return engine as LiteLlm
|
||||
}
|
||||
|
||||
private fun error(message: String): Nothing = throw IllegalStateException(message)
|
||||
@@ -5,6 +5,7 @@ import pw.binom.litert.LiteBackend
|
||||
import pw.binom.litert.LiteConfig
|
||||
import pw.binom.litert.LiteExperimental
|
||||
import pw.binom.litert.LiteLlm
|
||||
import pw.binom.litert.google.googleLiteLlm
|
||||
import pw.binom.litert.openai.OpenAiConfig as LitertOpenAiConfig
|
||||
import pw.binom.litert.openai.openAiLiteLlm
|
||||
|
||||
@@ -39,7 +40,7 @@ data class LlmConfig(
|
||||
|
||||
fun createLlm(): LiteLlm = when (backend) {
|
||||
LlmBackend.OPENAI -> openAiLiteLlm(checkNotNull(openai).toLitertConfig())
|
||||
LlmBackend.GOOGLE -> googleLiteLlmJvm(google!!.toLiteConfig())
|
||||
LlmBackend.GOOGLE -> googleLiteLlm(google!!.toLiteConfig())
|
||||
}
|
||||
|
||||
fun modelInfo(): String = when (backend) {
|
||||
|
||||
Reference in New Issue
Block a user