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