diff --git a/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/Event.kt b/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/Event.kt index e1683de..326543e 100644 --- a/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/Event.kt +++ b/outbox-api/src/commonMain/kotlin/pw/binom/agentik/outbox/Event.kt @@ -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 } diff --git a/settings.gradle.kts b/settings.gradle.kts index fed9baf..39f8ba1 100644 --- a/settings.gradle.kts +++ b/settings.gradle.kts @@ -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") diff --git a/skill-mining/src/commonMain/kotlin/pw/binom/agentik/skill/mining/SkillMiningComponent.kt b/skill-mining/src/commonMain/kotlin/pw/binom/agentik/skill/mining/SkillMiningComponent.kt index 3389baa..1685389 100644 --- a/skill-mining/src/commonMain/kotlin/pw/binom/agentik/skill/mining/SkillMiningComponent.kt +++ b/skill-mining/src/commonMain/kotlin/pw/binom/agentik/skill/mining/SkillMiningComponent.kt @@ -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 diff --git a/standalone/build.gradle.kts b/standalone/build.gradle.kts index ebe335e..3eca6c7 100644 --- a/standalone/build.gradle.kts +++ b/standalone/build.gradle.kts @@ -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/...), нам нужен diff --git a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/Main.kt b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/Main.kt index 3545267..8addd5e 100644 --- a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/Main.kt +++ b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/Main.kt @@ -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'а, diff --git a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/BackgroundEvents.kt b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/BackgroundEvents.kt deleted file mode 100644 index e54716c..0000000 --- a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/BackgroundEvents.kt +++ /dev/null @@ -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( - replay = 0, - extraBufferCapacity = 256, - onBufferOverflow = BufferOverflow.DROP_OLDEST, - ) - val toolCallEvents: SharedFlow get() = _toolCallEvents.asSharedFlow() - - private val _compactionEvents = MutableSharedFlow( - replay = 0, - extraBufferCapacity = 16, - onBufferOverflow = BufferOverflow.DROP_OLDEST, - ) - val compactionEvents: SharedFlow get() = _compactionEvents.asSharedFlow() - - private val _lifecycleEvents = MutableSharedFlow( - replay = 0, - extraBufferCapacity = 16, - onBufferOverflow = BufferOverflow.DROP_OLDEST, - ) - val lifecycleEvents: SharedFlow 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) -} diff --git a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ChatAgent.kt b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ChatAgent.kt index c6c01e2..9b482d8 100644 --- a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ChatAgent.kt +++ b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ChatAgent.kt @@ -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 = mutableListOf() +/** + * Тестовый инжектор: позволяет `ChatAgentTest` подсунуть тул, который + * скриптованная LLM будет вызывать. В production-конфигурации всегда пуст — + * инструменты приходят через [toolProviders] (компоненты, [install]). + */ +private val testTools: MutableList = mutableListOf() + +/** + * Единый канал всех событий агента — bounded tail с auto-TTL. + * + * Заменил ранее существовавшие два канала: + * - `agentEvents: MutableSharedFlow` (agent lifecycle) + * - per-conv `ConversationEvents._flow: MutableSharedFlow` + * + * Теперь оба пишут сюда через [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 = - merge( - _backgroundEvents.lifecycleEvents.filterIsInstance() - .map { SkillMiningEvent.ConversationClosing(it.conversationId) }, - _backgroundEvents.compactionEvents.filterIsInstance() - .map { SkillMiningEvent.ConversationCompacted(it.conversationId) }, - ) + eventStore.events(after = null) + .filterIsInstance() + .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` (agent lifecycle) - * - per-conv `ConversationEvents._flow: MutableSharedFlow` - * - * Теперь оба пишут сюда через [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 } @@ -532,11 +529,9 @@ internal class ChatAgent( memoryStoreForReview = memoryStore, contextWindow = contextWindow, compressionThreshold = compressionThreshold, - contextCompactor = contextCompactor, - reflector = reflector, - skillMiner = skillMiner, - skillMiningStore = skillStore, - ) + contextCompactor = contextCompactor, + reflector = reflector, + ) override fun close() { runBlocking { diff --git a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ChatConversation.kt b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ChatConversation.kt index a07c59b..9117766 100644 --- a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ChatConversation.kt +++ b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ChatConversation.kt @@ -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] (этот же пакет). diff --git a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/CompactionCoordinator.kt b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/CompactionCoordinator.kt index af60456..c662d1a 100644 --- a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/CompactionCoordinator.kt +++ b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/CompactionCoordinator.kt @@ -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 } } diff --git a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ConversationLoop.kt b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ConversationLoop.kt index b84b5ef..13c4714 100644 --- a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ConversationLoop.kt +++ b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ConversationLoop.kt @@ -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 = 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. } } diff --git a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/BackgroundScheduler.kt b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ReflectionScheduler.kt similarity index 53% rename from standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/BackgroundScheduler.kt rename to standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ReflectionScheduler.kt index 21a3346..bf8f57e 100644 --- a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/BackgroundScheduler.kt +++ b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ReflectionScheduler.kt @@ -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() + eventStore.events(after = null) + .filterIsInstance() + .filter { it.event is Event.ConversationClosing && it.conversationId == conversationIdProvider() } .onEach { onClosing() }, - backgroundEvents.compactionEvents - .filterIsInstance() - .onEach { onCompaction(it) }, - backgroundEvents.toolCallEvents - .filterIsInstance() + eventStore.events(after = null) + .filterIsInstance() + .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().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. */ diff --git a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ToolDispatcher.kt b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ToolDispatcher.kt index c92bd72..cf7123a 100644 --- a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ToolDispatcher.kt +++ b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/agent/ToolDispatcher.kt @@ -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, 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) { diff --git a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/llm/GoogleLiteLlmJvm.kt b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/llm/GoogleLiteLlmJvm.kt deleted file mode 100644 index ff9b224..0000000 --- a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/llm/GoogleLiteLlmJvm.kt +++ /dev/null @@ -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=, 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) diff --git a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/llm/LlmConfig.kt b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/llm/LlmConfig.kt index 8a6ee4c..6bf9f47 100644 --- a/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/llm/LlmConfig.kt +++ b/standalone/src/commonMain/kotlin/pw/binom/agentik/standalone/llm/LlmConfig.kt @@ -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) {