From 9276de3c94a8f31a680237b86bd647b0b8c4cd7a 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 08:54:01 +0300 Subject: [PATCH] fix: collect response from events, remove outcome polling MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit outcome: succeeded = сессия запущена, НЕ ответ готов. Ответ приходит ТОЛЬКО через EventV2 subscription (live-only text). Изменения: - Убран polling ctx.session.get() + outcome - Event subscription: парсит 7 структур событий для извлечения текста - Логгирует ВСЕ session-события с полной структурой (для отладки) - Promise разрешается сразу при получении первого текста из событий - Safety timeout 60s (не основной механизм, только защита) --- index.ts | 165 ++++++++++++++++++++++++++----------------------------- 1 file changed, 78 insertions(+), 87 deletions(-) diff --git a/index.ts b/index.ts index da0c237..3d161d2 100644 --- a/index.ts +++ b/index.ts @@ -93,45 +93,88 @@ 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 } + // Pending prompt responses — maps sessionID -> { resolve, reject, text } const pendingResponses = new Map void reject: (err: any) => void + text: string }>() - // V2 Event subscription — logs all events for debugging + // V2 Event subscription — collects assistant response text from events void (async () => { try { for await (const event of ctx.event.subscribe({ signal: eventController.signal })) { - // Log event type for debugging - let eventType = "" + // Parse event — V2EventEncoded is a JSON string let eventData: any = event if (typeof event === "string") { try { eventData = JSON.parse(event) - eventType = eventData?.type || "" } catch { - eventType = "(raw string)" + // not JSON — log raw for debugging + log(`Raw event (not JSON): ${event.slice(0, 200)}`) + continue } - } else { - eventType = eventData?.type || "" } - // Check for session completion events + const eventType = eventData?.type || "" const sessionId = eventData?.sessionID || eventData?.session?.id || "" - if (sessionId && pendingResponses.has(sessionId)) { - // Check if event indicates session completion - const outcome = eventData?.outcome || eventData?.data?.outcome || eventData?.update?.outcome - if (outcome && ["succeeded", "failed", "interrupted"].includes(outcome)) { - log(`Session ${sessionId} completed with outcome: ${outcome} (from event)`) - // Will be picked up by polling loop - } + + // Log ALL session events for debugging + if (eventType.includes("session") || sessionId) { + const preview = JSON.stringify(eventData).slice(0, 500) + log(`Event: type=${eventType}, session=${sessionId}, data=${preview}`) } - // Log first 200 chars of interesting events - if (eventType.includes("session") && sessionId) { - const preview = typeof event === "string" ? event.slice(0, 200) : JSON.stringify(eventData).slice(0, 200) - log(`Event: type=${eventType}, session=${sessionId}, preview=${preview}`) + // Check if this event has response text for a pending prompt + if (sessionId && pendingResponses.has(sessionId)) { + const pending = pendingResponses.get(sessionId)! + + // Try multiple structures to extract assistant text + let newText = "" + + // 1. event.data.content[].text (assistant text parts) + const parts = eventData?.data?.content || eventData?.content || eventData?.update?.content + if (Array.isArray(parts)) { + for (const part of parts) { + if (part?.type === "text" && part?.text) { + newText += part.text + } + } + } + // 2. event.data.text (direct text) + if (!newText && typeof eventData?.data?.text === "string") { + newText = eventData.data.text + } + // 3. event.data.message.text (message text) + if (!newText && eventData?.data?.message?.text) { + newText = eventData.data.message.text + } + // 4. event.payload.text (synthetic message) + if (!newText && typeof eventData?.payload?.text === "string") { + newText = eventData.payload.text + } + // 5. event.update.content (simple content update) + if (!newText && typeof eventData?.update?.content === "string") { + newText = eventData.update.content + } + // 6. event.delta (streaming delta) + if (!newText && typeof eventData?.delta === "string") { + newText = eventData.delta + } + // 7. event.content (direct content) + if (!newText && typeof eventData?.content === "string") { + newText = eventData.content + } + + if (newText) { + pending.text += newText + log(`Collected ${newText.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) + } } } } catch (err) { @@ -211,74 +254,22 @@ async function runPlugin(ctx: any, options: any, dataDir: string) { }) log(`Prompt sent to session ${opencodeSessionId}`) - // Poll session.get() to wait for completion (outcome field) - let responseText = "" - let lastAssistantText = "" - - while (true) { - await new Promise((r) => setTimeout(r, 500)) - - try { - const sessionInfo = await ctx.session.get({ sessionID: opencodeSessionId }) - const outcome = sessionInfo?.outcome || sessionInfo?.data?.outcome - - // Session completed — collect response text from messages - if (outcome && ["succeeded", "failed", "interrupted"].includes(outcome)) { - log(`Session ${opencodeSessionId} completed with outcome: ${outcome}`) - - // Collect assistant text from session messages - const messages = sessionInfo?.messages || sessionInfo?.data?.messages - if (Array.isArray(messages)) { - for (const msg of messages) { - if (msg?.type === "assistant" || msg?.role === "assistant") { - const content = msg?.content || msg?.text - if (typeof content === "string") { - lastAssistantText += content - } else if (Array.isArray(content)) { - for (const part of content) { - if (part?.type === "text" && part?.text) { - lastAssistantText += part.text - } - } - } - } - } - } - - // Also check session metadata for response text - if (!lastAssistantText) { - const metaText = sessionInfo?.metadata?.responseText || sessionInfo?.data?.metadata?.responseText - if (metaText) lastAssistantText = String(metaText) - } - - responseText = lastAssistantText || "No response text collected" - break - } - - // Track assistant text as it streams (for partial responses) - const streamMessages = sessionInfo?.streamMessages || sessionInfo?.data?.streamMessages - if (Array.isArray(streamMessages)) { - for (const msg of streamMessages) { - if (msg?.type === "assistant" || msg?.role === "assistant") { - const content = msg?.content || msg?.text - if (typeof content === "string" && content !== lastAssistantText) { - lastAssistantText = content - } else if (Array.isArray(content)) { - let text = "" - for (const part of content) { - if (part?.type === "text" && part?.text) text += part.text - } - if (text && text !== lastAssistantText) lastAssistantText = text - } - } - } - } - } catch (getErr: any) { - log(`Session get error (will retry): ${String(getErr)}`) - // Session might be temporarily unavailable, keep polling - } - } + // Wait for response text from event subscription + 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 + }) + const responseText = await responsePromise session.outputChars += responseText.length await matrix.sendReply(context, responseText) } catch (err: any) {