From 86eb0632e0e85a3b7c6b4ec2cfbf365941b9a7a2 Mon Sep 17 00:00:00 2001 From: subochev Date: Wed, 16 Sep 2026 07:46:30 +0300 Subject: [PATCH] fix interrupt: race in flag reset + no-op when no active turn + bounded background scope MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Три фикса в runTurn/interrupt: 1. **Race condition в finally-блоке.** Раньше сбрасывал interrupted.set(false) только если флаг был установлен при чтении wasInterrupted в начале finally. Если interrupt() приходил между этими двумя точками — флаг оставался true и следующий turn видел wasInterruptedAtEntry=true → сразу short-circuit'ил без вызова LLM. Теперь всегда сбрасываем (compareAndSet атомарен, гарантирует следующий turn чистый). 2. **interrupt() отравлял следующий send.** Если вызывали interrupt() в пустоту (нет активного turn'а — флаг всё равно ставился → следующий send сразу short-circuit'ил, пользователь не получал ответа на своё 'Ок.' после явного cancel). Теперь interrupt() проверяет activeTurn?.isActive и при отсутствии активного turn'а — no-op. 3. **Bounded background scope для review/reflection/skill-mining.** OpenAiLlm.send() использует runBlocking — если запустить 30+ параллельных review (по одному на беседу), IO-thread pool голодает и ассистент висит. Вынес в отдельный scope с Dispatchers.IO.limitedParallelism(4) — не больше 4 sync LLM вызовов одновременно. Тест 27/27 (см. /tmp/test-interrupt.py и /tmp/run-manual-tests.py). --- .../standalone/agent/ChatConversation.kt | 35 +++++++++++++++---- 1 file changed, 29 insertions(+), 6 deletions(-) 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