fix interrupt: race in flag reset + no-op when no active turn + bounded background scope
Три фикса в 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).
This commit is contained in:
+29
-6
@@ -179,7 +179,15 @@ class ChatConversation(
|
|||||||
)
|
)
|
||||||
|
|
||||||
private val turnLock = Mutex()
|
private val turnLock = Mutex()
|
||||||
|
// Основной scope для активного turn'а — IO (много потоков, не упираемся).
|
||||||
private val scope = CoroutineScope(SupervisorJob() + Dispatchers.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
|
@Volatile
|
||||||
private var liteConv: LiteConversation? = null
|
private var liteConv: LiteConversation? = null
|
||||||
@@ -260,6 +268,15 @@ override suspend fun interrupt() {
|
|||||||
// activeTurn НЕ cancel — даём runTurn'у finally-блоку корректно
|
// activeTurn НЕ cancel — даём runTurn'у finally-блоку корректно
|
||||||
// записать state (частичный assistant text + cancelled tool exchanges)
|
// записать state (частичный assistant text + cancelled tool exchanges)
|
||||||
// и эмитить End. close() тоже не вызываем — это сделает finally.
|
// и эмитить 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)
|
interrupted.set(true)
|
||||||
runCatching { liteConv?.cancel() }
|
runCatching { liteConv?.cancel() }
|
||||||
currentToolJob?.cancel()
|
currentToolJob?.cancel()
|
||||||
@@ -582,13 +599,19 @@ override suspend fun interrupt() {
|
|||||||
// Interrupted event — клиент видит его в SSE сразу как прерывание
|
// Interrupted event — клиент видит его в SSE сразу как прерывание
|
||||||
// произошло (на самом деле он эмитится в finally, после возможного
|
// произошло (на самом деле он эмитится в finally, после возможного
|
||||||
// финального ответа модели — это нормально, клиент рендерит оба).
|
// финального ответа модели — это нормально, клиент рендерит оба).
|
||||||
if (wasInterrupted) {
|
//
|
||||||
|
// Если interrupt() пришёл ПОСЛЕ того, как мы прочитали wasInterrupted
|
||||||
|
// в начале finally — перечитываем флаг здесь, чтобы корректно эмитить
|
||||||
|
// Interrupted event и сбрасывать флаг после.
|
||||||
|
if (wasInterrupted || interrupted.get()) {
|
||||||
emitEvent(ProtoEvent.Interrupted(date = now()))
|
emitEvent(ProtoEvent.Interrupted(date = now()))
|
||||||
}
|
}
|
||||||
emitEvent(ProtoEvent.End(date = now()))
|
emitEvent(ProtoEvent.End(date = now()))
|
||||||
|
|
||||||
// Сбрасываем флаг — следующий turn стартует чистым.
|
// Сбрасываем флаг — следующий turn стартует чистым. Всегда
|
||||||
if (wasInterruptedAtEntry || wasInterrupted) interrupted.set(false)
|
// (compareAndSet атомарен, защищаем от race-condition: interrupt()
|
||||||
|
// мог быть вызван между чтением wasInterrupted и этой строкой).
|
||||||
|
interrupted.set(false)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -819,7 +842,7 @@ override suspend fun interrupt() {
|
|||||||
.joinToString("\n") { it.body }
|
.joinToString("\n") { it.body }
|
||||||
if (userText.isBlank() || assistantText.isBlank()) return
|
if (userText.isBlank() || assistantText.isBlank()) return
|
||||||
val convId = id
|
val convId = id
|
||||||
scope.launch {
|
backgroundScope.launch {
|
||||||
try {
|
try {
|
||||||
val decision: MemoryReviewDecision = reviewer.review(
|
val decision: MemoryReviewDecision = reviewer.review(
|
||||||
ReviewedTurn(
|
ReviewedTurn(
|
||||||
@@ -867,7 +890,7 @@ override suspend fun interrupt() {
|
|||||||
val userTurnCount = countUserTurnsBlocking()
|
val userTurnCount = countUserTurnsBlocking()
|
||||||
if (userTurnCount % reflectionInterval != 0) return
|
if (userTurnCount % reflectionInterval != 0) return
|
||||||
val convId = id
|
val convId = id
|
||||||
scope.launch {
|
backgroundScope.launch {
|
||||||
try {
|
try {
|
||||||
val turns = listOf(
|
val turns = listOf(
|
||||||
pw.binom.agentik.memory.ConversationTurn(
|
pw.binom.agentik.memory.ConversationTurn(
|
||||||
@@ -916,7 +939,7 @@ override suspend fun interrupt() {
|
|||||||
val userTurnCount = countUserTurnsBlocking()
|
val userTurnCount = countUserTurnsBlocking()
|
||||||
if (userTurnCount % skillMiningInterval != 0) return
|
if (userTurnCount % skillMiningInterval != 0) return
|
||||||
val convId = id
|
val convId = id
|
||||||
scope.launch {
|
backgroundScope.launch {
|
||||||
try {
|
try {
|
||||||
val turns = recentTurnsFromWorkingMemory(miner.maxTurns)
|
val turns = recentTurnsFromWorkingMemory(miner.maxTurns)
|
||||||
if (turns.isEmpty()) return@launch
|
if (turns.isEmpty()) return@launch
|
||||||
|
|||||||
Reference in New Issue
Block a user