From 280eaa1e9cb76b2fac0011e79ba6387e99efadbd Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=D0=91=D0=BE=D1=80=D0=BE=D0=B4=D0=B8=D0=BD=20=D0=A0=D0=BE?= =?UTF-8?q?=D0=BC=D0=B0=D0=BD?= Date: Sat, 26 Sep 2026 09:22:32 +0300 Subject: [PATCH] fix: wait for session.text.ended, no timeouts MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Полный flow сессии из логов: prompt → inbox.enqueued → execution.started → step.started → reasoning → tool-calls → step.ended (finish: tool-calls) → step.started → reasoning → text.started → text.delta*42 → text.ended ← ОТВЕТ ГОТОВ → step.ended (finish: stop) → execution.succeeded Ключевые сигналы: session.text.ended → ответ готов (полный текст в data.text) session.execution.failed → ошибка session.execution.interrupted → прервано session.execution.succeeded → 'запрос принят' (НЕ ответ готов!) Изменения: - Убраны timeout 60s (бесконечное ожидание) - Resolved на session.text.ended (полный assembled text) - Rejected на session.execution.failed/interrupted - Убраны reasoning.* сборы (не для пользователя) - Убраны text.delta ранние resolve (race condition) --- index.ts | 94 +++++++++++++++++--------------------------------------- 1 file changed, 29 insertions(+), 65 deletions(-) diff --git a/index.ts b/index.ts index cbe1e37..09912b0 100644 --- a/index.ts +++ b/index.ts @@ -93,80 +93,52 @@ async function runPlugin(ctx: any, options: any, dataDir: string) { // Track Matrix context per OpenCode session const sessionContextMap = new Map() - // Pending prompt responses — maps sessionID -> { resolve, reject, text } + // Pending prompt responses — maps sessionID -> { resolve, reject } const pendingResponses = new Map void reject: (err: any) => void - text: string }>() - // V2 Event subscription — collects assistant response text from events + // V2 Event subscription — waits for response completion signals + // Event structure: { type: "...", data: { sessionID: "...", ... } } void (async () => { try { for await (const event of ctx.event.subscribe({ signal: eventController.signal })) { - // Parse event — V2EventEncoded is a JSON string let eventData: any = event if (typeof event === "string") { - try { - eventData = JSON.parse(event) - } catch { - continue - } + try { eventData = JSON.parse(event) } catch { continue } } const eventType = eventData?.type || "" + const sessionId = eventData?.data?.sessionID || "" - // Session text delta — streaming text chunks - // Structure: { type: "session.text.delta", data: { sessionID, delta: "text" } } - if (eventType === "session.text.delta") { - const sessionId = eventData?.data?.sessionID || "" - const delta = eventData?.data?.delta || "" - if (sessionId && delta && pendingResponses.has(sessionId)) { - const pending = pendingResponses.get(sessionId)! - pending.text += delta - log(`Collected delta: ${delta.length} chars, total: ${pending.text.length} chars for session ${sessionId}`) - // Resolve immediately when we get any text - const timer = (pending as any)._safetyTimer - if (timer) clearTimeout(timer) - pendingResponses.delete(sessionId) - pending.resolve(pending.text) - } - } - - // Session text ended — complete text available + // Response text complete — full assembled text available // Structure: { type: "session.text.ended", data: { sessionID, text: "full text" } } - if (eventType === "session.text.ended") { - const sessionId = eventData?.data?.sessionID || "" + if (eventType === "session.text.ended" && sessionId && pendingResponses.has(sessionId)) { const fullText = eventData?.data?.text || "" - if (sessionId && fullText && pendingResponses.has(sessionId)) { - const pending = pendingResponses.get(sessionId)! - pending.text = fullText - log(`Collected ended text: ${fullText.length} chars for session ${sessionId}`) - const timer = (pending as any)._safetyTimer - if (timer) clearTimeout(timer) - pendingResponses.delete(sessionId) - pending.resolve(pending.text) - } + const pending = pendingResponses.get(sessionId)! + log(`Response ready: ${fullText.length} chars for session ${sessionId}`) + pendingResponses.delete(sessionId) + pending.resolve(fullText) + continue } - // Session reasoning delta — also collect (optional, for full response) - if (eventType === "session.reasoning.delta") { - const sessionId = eventData?.data?.sessionID || "" - const delta = eventData?.data?.delta || "" - if (sessionId && delta && pendingResponses.has(sessionId)) { - const pending = pendingResponses.get(sessionId)! - pending.text += delta - } + // Execution failed — session errored + if (eventType === "session.execution.failed" && sessionId && pendingResponses.has(sessionId)) { + const pending = pendingResponses.get(sessionId)! + log(`Session ${sessionId} execution failed`) + pendingResponses.delete(sessionId) + pending.reject(new Error("Session execution failed")) + continue } - // Session reasoning ended — full reasoning text - if (eventType === "session.reasoning.ended") { - const sessionId = eventData?.data?.sessionID || "" - const text = eventData?.data?.text || "" - if (sessionId && text && pendingResponses.has(sessionId)) { - const pending = pendingResponses.get(sessionId)! - pending.text += text - } + // Execution interrupted — session was interrupted + if (eventType === "session.execution.interrupted" && sessionId && pendingResponses.has(sessionId)) { + const pending = pendingResponses.get(sessionId)! + log(`Session ${sessionId} execution interrupted`) + pendingResponses.delete(sessionId) + pending.reject(new Error("Session execution interrupted")) + continue } } } catch (err) { @@ -247,18 +219,10 @@ async function runPlugin(ctx: any, options: any, dataDir: string) { log(`Prompt sent to session ${opencodeSessionId}`) // Wait for response text from event subscription + // Resolved by: session.text.ended (success) + // Rejected by: session.execution.failed / session.execution.interrupted const responsePromise = new Promise((resolve, reject) => { - const safetyTimer = setTimeout(() => { - pendingResponses.delete(opencodeSessionId) - reject(new Error("Timeout waiting for response event (60s)")) - }, 60_000) - pendingResponses.set(opencodeSessionId, { - resolve, - reject, - text: "", - }) - // Store timer reference for cleanup in event subscription - ;(pendingResponses.get(opencodeSessionId) as any)._safetyTimer = safetyTimer + pendingResponses.set(opencodeSessionId, { resolve, reject }) }) const responseText = await responsePromise