diff --git a/index.ts b/index.ts index 2c9c659..da0c237 100644 --- a/index.ts +++ b/index.ts @@ -93,72 +93,45 @@ 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, timer } + // Pending prompt responses — maps sessionID -> { resolve, reject } const pendingResponses = new Map void reject: (err: any) => void - timer: ReturnType }>() - // V2 Event subscription — collects assistant response text from events + // V2 Event subscription — logs all events for debugging void (async () => { try { for await (const event of ctx.event.subscribe({ signal: eventController.signal })) { - // Parse event — V2EventEncoded is a JSON string + // Log event type for debugging + let eventType = "" let eventData: any = event if (typeof event === "string") { try { eventData = JSON.parse(event) - } catch { /* not JSON, use raw */ } + eventType = eventData?.type || "" + } catch { + eventType = "(raw string)" + } + } else { + eventType = eventData?.type || "" } - const eventType = eventData?.type || "" + // Check for session completion events const sessionId = eventData?.sessionID || eventData?.session?.id || "" - - // Check if this event contains assistant response text for a pending prompt if (sessionId && pendingResponses.has(sessionId)) { - const pending = pendingResponses.get(sessionId)! - - // Extract text from various event structures - let responseText = "" - - // Try: event.messages[].content[].text - const messages = eventData?.messages || eventData?.data?.messages || eventData?.update?.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") responseText += content - else if (Array.isArray(content)) { - for (const part of content) { - if (part?.type === "text" && part?.text) responseText += part.text - } - } - } - } + // 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 } + } - // Try: event.update.content (simple text update) - if (!responseText && eventData?.update?.content) { - responseText = String(eventData.update.content) - } - - // Try: event.payload.text (synthetic message) - if (!responseText && eventData?.payload?.text) { - responseText = eventData.payload.text - } - - // Try: event.data.text (generate response) - if (!responseText && eventData?.data?.text) { - responseText = eventData.data.text - } - - if (responseText) { - log(`Collected response text (${responseText.length} chars) for session ${sessionId}`) - clearTimeout(pending.timer) - pendingResponses.delete(sessionId) - pending.resolve(responseText) - } + // 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}`) } } } catch (err) { @@ -236,17 +209,76 @@ async function runPlugin(ctx: any, options: any, dataDir: string) { sessionID: opencodeSessionId, text: query, }) + log(`Prompt sent to session ${opencodeSessionId}`) - // Wait for response text from event subscription - const responsePromise = new Promise((resolve, reject) => { - const timeout = setTimeout(() => { - pendingResponses.delete(opencodeSessionId) - reject(new Error("Timeout waiting for response (30s)")) - }, 30_000) - pendingResponses.set(opencodeSessionId, { resolve, reject, timer: timeout }) - }) + // 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 + } + } - const responseText = await responsePromise session.outputChars += responseText.length await matrix.sendReply(context, responseText) } catch (err: any) {