diff --git a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ChatConversation.kt b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ChatConversation.kt index 352caaf..8e759f4 100644 --- a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ChatConversation.kt +++ b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ChatConversation.kt @@ -179,7 +179,15 @@ class ChatConversation( ) private val turnLock = Mutex() + // Основной scope для активного turn'а — IO (много потоков, не упираемся). private val scope = CoroutineScope(SupervisorJob() + Dispatchers.IO) + // Фоновые задачи (review/reflection/skill-mining) — отдельный pool с + // bounded parallelism. Они делают sync LLM-вызовы через runBlocking внутри + // LiteLlm.send() — если запустить 30+ параллельно (по одному на беседу), + // упираемся в IO-thread starvation и все повисают в очереди. + private val backgroundScope = CoroutineScope( + SupervisorJob() + Dispatchers.IO.limitedParallelism(4), + ) @Volatile private var liteConv: LiteConversation? = null @@ -260,6 +268,15 @@ override suspend fun interrupt() { // activeTurn НЕ cancel — даём runTurn'у finally-блоку корректно // записать state (частичный assistant text + cancelled tool exchanges) // и эмитить End. close() тоже не вызываем — это сделает finally. + // + // NB: ставим флаг ТОЛЬКО если есть что прерывать. Если вызвали interrupt() + // "в пустоту" (нет активного turn'а), флаг не должен отравлять следующий + // send — иначе runTurn'у следующего хода сразу придётся short-circuit'нуть, + // и пользователь не получит ответа на свой "Ок." после явного cancel. + if (activeTurn?.isActive != true) { + log.info { "interrupt() no-op: no active turn for $id" } + return + } interrupted.set(true) runCatching { liteConv?.cancel() } currentToolJob?.cancel() @@ -582,13 +599,19 @@ override suspend fun interrupt() { // Interrupted event — клиент видит его в SSE сразу как прерывание // произошло (на самом деле он эмитится в finally, после возможного // финального ответа модели — это нормально, клиент рендерит оба). - if (wasInterrupted) { + // + // Если interrupt() пришёл ПОСЛЕ того, как мы прочитали wasInterrupted + // в начале finally — перечитываем флаг здесь, чтобы корректно эмитить + // Interrupted event и сбрасывать флаг после. + if (wasInterrupted || interrupted.get()) { emitEvent(ProtoEvent.Interrupted(date = now())) } emitEvent(ProtoEvent.End(date = now())) - // Сбрасываем флаг — следующий turn стартует чистым. - if (wasInterruptedAtEntry || wasInterrupted) interrupted.set(false) + // Сбрасываем флаг — следующий turn стартует чистым. Всегда + // (compareAndSet атомарен, защищаем от race-condition: interrupt() + // мог быть вызван между чтением wasInterrupted и этой строкой). + interrupted.set(false) } } @@ -819,7 +842,7 @@ override suspend fun interrupt() { .joinToString("\n") { it.body } if (userText.isBlank() || assistantText.isBlank()) return val convId = id - scope.launch { + backgroundScope.launch { try { val decision: MemoryReviewDecision = reviewer.review( ReviewedTurn( @@ -867,7 +890,7 @@ override suspend fun interrupt() { val userTurnCount = countUserTurnsBlocking() if (userTurnCount % reflectionInterval != 0) return val convId = id - scope.launch { + backgroundScope.launch { try { val turns = listOf( pw.binom.agentik.memory.ConversationTurn( @@ -916,7 +939,7 @@ override suspend fun interrupt() { val userTurnCount = countUserTurnsBlocking() if (userTurnCount % skillMiningInterval != 0) return val convId = id - scope.launch { + backgroundScope.launch { try { val turns = recentTurnsFromWorkingMemory(miner.maxTurns) if (turns.isEmpty()) return@launch