diff --git a/plugins/matrix-plugin/.gitignore b/.gitignore similarity index 100% rename from plugins/matrix-plugin/.gitignore rename to .gitignore diff --git a/plugins/matrix-plugin/PLUGIN-DOCS-v2.md b/PLUGIN-DOCS-v2.md similarity index 100% rename from plugins/matrix-plugin/PLUGIN-DOCS-v2.md rename to PLUGIN-DOCS-v2.md diff --git a/plugins/matrix-plugin/README.md b/README.md similarity index 100% rename from plugins/matrix-plugin/README.md rename to README.md diff --git a/plugins/matrix-plugin/bun.lock b/bun.lock similarity index 100% rename from plugins/matrix-plugin/bun.lock rename to bun.lock diff --git a/plugins/matrix-plugin/config-loader.ts b/config-loader.ts similarity index 100% rename from plugins/matrix-plugin/config-loader.ts rename to config-loader.ts diff --git a/plugins/matrix-plugin/index.ts b/index.ts similarity index 71% rename from plugins/matrix-plugin/index.ts rename to index.ts index b10487e..9b0ad51 100644 --- a/plugins/matrix-plugin/index.ts +++ b/index.ts @@ -21,7 +21,6 @@ function isV2Context(ctx: any): boolean { async function runPlugin(ctx: any, options: any, sdkClient?: any) { const isV2 = isV2Context(ctx) const eventController = new AbortController() - let expiryInterval: ReturnType | undefined let cleanup: (() => void) | undefined try { @@ -39,6 +38,58 @@ async function runPlugin(ctx: any, options: any, sdkClient?: any) { const matrix = new MatrixBotClient(options) const sessionManager = new SessionManager(options.rateLimitSeconds || 5) + // Track pending tool calls per session: sessionID -> Map + const pendingToolCalls = new Map>() + + // Track sent messages per session to avoid duplicates + const sentMessageIds = new Map>() + + // Track Matrix context per OpenCode session + const sessionContextMap = new Map() + + // Helper: send a message to Matrix + async function sendMatrixMessage(roomId: string, threadRootId: string, context: any, text: string): Promise { + if (!text.trim()) return + await matrix.sendReply(context, text) + log(`Matrix: Sent to ${roomId}:${threadRootId} (${text.length} chars)`) + } + + // Helper: extract content from parts + function extractContent(parts: any[]): { reasoningText: string; mainText: string; toolCallText: string } { + let reasoningText = "" + let mainText = "" + let toolCallText = "" + const reasoningTypes = new Set(["reasoning", "thinking", "redwood", "reason"]) + + for (const p of parts) { + const content = p?.content || p?.text || "" + const pType = p?.type || "" + if (reasoningTypes.has(pType)) { + reasoningText += content + } else if (pType === "tool_call" || pType === "tool") { + toolCallText += content + } else { + mainText += content + } + } + + return { reasoningText, mainText, toolCallText } + } + + // Helper: format message text + function formatMessage(toolCallText: string, reasoningText: string, mainText: string): string { + if (toolCallText.trim()) { + return `Tool call: ${toolCallText.trim()}` + } + + let responseText = mainText + if (reasoningText.trim()) { + responseText = `
Thoughts\n\n${reasoningText.trim()}\n\n
\n\n---\n\n${mainText}` + } + + return responseText + } + // V2 Event subscription for response streaming if (isV2) { void (async () => { @@ -50,6 +101,7 @@ async function runPlugin(ctx: any, options: any, sdkClient?: any) { if (!sessionId) continue const update = (event as { update?: { type: string; content?: string; toolName?: string } }).update if (!update) continue + log(`V2 event: session=${sessionId}, update.type=${update.type}`) } } } catch (err) { @@ -62,30 +114,6 @@ async function runPlugin(ctx: any, options: any, sdkClient?: any) { })() } - // V1 Event subscription for response streaming - if (!isV2 && sdkClient) { - void (async () => { - try { - const events = await sdkClient.event.subscribe() - for await (const event of events.stream) { - const eventType = event.type || "" - if (eventType === "session.updated") { - const sessionId = event.properties?.sessionID - if (!sessionId) continue - const update = event.properties?.update - if (!update) continue - } - } - } catch (err) { - if (String(err).includes("AbortError") || (err as { name?: string }).name === "AbortError") { - // expected during cleanup - } else { - log(`V1 event subscription error: ${String(err)}`) - } - } - })() - } - // Message handler matrix.on(async (data: any) => { const { context, query, sender, roomId, eventId, timestamp } = data @@ -133,6 +161,7 @@ async function runPlugin(ctx: any, options: any, sdkClient?: any) { } session.opencodeSessionId = opencodeSessionId sessionManager.saveSession(roomId, threadRootId, opencodeSessionId) + sessionContextMap.set(opencodeSessionId, { roomId, threadRootId, context }) } catch (createErr: any) { log(`Failed to create session: ${String(createErr)}`) session.isActive = false @@ -157,66 +186,87 @@ async function runPlugin(ctx: any, options: any, sdkClient?: any) { // V1: use SDK client.session.prompt() log(`V1: Sending prompt to session ${opencodeSessionId}: ${query}`) + // Start polling for new messages + const pollInterval = setInterval(async () => { + if (!session.isActive) { + clearInterval(pollInterval) + return + } + + try { + const existingIds = sentMessageIds.get(opencodeSessionId) || new Set() + const messagesResult = await sdkClient.session.messages({ + path: { id: opencodeSessionId }, + }) + + const messages = messagesResult?.data || [] + if (Array.isArray(messages)) { + for (const msg of messages) { + const msgId = msg?.info?.id || msg?.id + if (!msgId || existingIds.has(msgId)) continue + + const role = msg?.info?.role || "" + if (role !== "assistant" && role !== "tool_call") continue + + const parts = msg?.parts || msg?.info?.parts || [] + const { reasoningText, mainText, toolCallText } = extractContent(parts) + const responseText = formatMessage(toolCallText, reasoningText, mainText) + + if (responseText.trim()) { + await matrix.sendReply(context, responseText) + log(`V1: Sent message ${msgId} (role=${role}, ${responseText.length} chars)`) + existingIds.add(msgId) + } + } + } + + sentMessageIds.set(opencodeSessionId, existingIds) + } catch (err) { + log(`V1: Error during polling: ${String(err)}`) + } + }, 1000) + try { const result = await sdkClient.session.prompt({ path: { id: opencodeSessionId }, body: { parts: [{ type: "text", text: query }] }, }) - log(`V1: prompt result keys=${Object.keys(result || {})}`) - log(`V1: prompt result.data.keys=${result?.data ? Object.keys(result.data).join(",") : "null"}`) - if (result?.data?.parts) { - log(`V1: data.parts count=${result.data.parts.length}`) - log(`V1: data.parts.types=${result.data.parts.map((p: any) => p.type).join(",")}`) - log(`V1: data.parts[0]=${JSON.stringify(result.data.parts[0]).slice(0, 200)}`) - } - // Try to extract content from various locations, separating reasoning from main text - let reasoningText = "" - let mainText = "" - const reasoningTypes = new Set(["reasoning", "thinking", "redwood", "reason"]) - if (result?.data?.parts) { - for (const p of result.data.parts) { - const content = p.content || p.text || "" - if (reasoningTypes.has(p.type)) { - reasoningText += content - } else { - mainText += content + // Stop polling + clearInterval(pollInterval) + + // Send any remaining messages + try { + const existingIds = sentMessageIds.get(opencodeSessionId) || new Set() + const messagesResult = await sdkClient.session.messages({ + path: { id: opencodeSessionId }, + }) + + const messages = messagesResult?.data || [] + if (Array.isArray(messages)) { + for (const msg of messages) { + const msgId = msg?.info?.id || msg?.id + if (!msgId || existingIds.has(msgId)) continue + + const role = msg?.info?.role || "" + if (role !== "assistant" && role !== "tool_call") continue + + const parts = msg?.parts || msg?.info?.parts || [] + const { reasoningText, mainText, toolCallText } = extractContent(parts) + const responseText = formatMessage(toolCallText, reasoningText, mainText) + + if (responseText.trim()) { + await matrix.sendReply(context, responseText) + log(`V1: Sent message ${msgId} (role=${role}, ${responseText.length} chars)`) + existingIds.add(msgId) + } } } - } else if (result?.data?.info?.parts) { - for (const p of result.data.info.parts) { - const content = p.content || p.text || "" - if (reasoningTypes.has(p.type)) { - reasoningText += content - } else { - mainText += content - } - } - } else if (result?.data?.message?.parts) { - for (const p of result.data.message.parts) { - const content = p.content || p.text || "" - if (reasoningTypes.has(p.type)) { - reasoningText += content - } else { - mainText += content - } - } - } else if (result?.data?.content) { - mainText = result.data.content - } else if (result?.response?.content) { - mainText = result.response.content - } else if (result?.content) { - mainText = result.content + + sentMessageIds.set(opencodeSessionId, existingIds) + } catch (err) { + log(`V1: Error sending final messages: ${String(err)}`) } - - let responseText = mainText - if (reasoningText.trim()) { - responseText = `
Thoughts\n\n${reasoningText.trim()}\n\n
\n\n---\n\n${mainText}` - } - - session.outputChars += responseText.length - await matrix.sendReply(context, responseText) - log(`V1: Response sent (${responseText.length} chars)`) } catch (promptErr: any) { log(`V1 prompt error: ${String(promptErr)}`) const errMsg = promptErr?.response?.error?.message || promptErr?.error?.message || String(promptErr) @@ -236,6 +286,7 @@ async function runPlugin(ctx: any, options: any, sdkClient?: any) { } finally { session.isActive = false session.lastActivity = Date.now() + sessionContextMap.delete(opencodeSessionId) releaseQuery() } }) @@ -317,7 +368,7 @@ async function runPlugin(ctx: any, options: any, sdkClient?: any) { } }) - // Start Matrix client + // V2 Storage - defensive access if (isV2 && ctx?.storage) { await ctx.storage.set("matrix/plugin_version", "1.0.0") await ctx.storage.set("matrix/started_at", new Date().toISOString()) diff --git a/plugins/matrix-plugin/logger.ts b/logger.ts similarity index 100% rename from plugins/matrix-plugin/logger.ts rename to logger.ts diff --git a/plugins/matrix-plugin/matrix-client.ts b/matrix-client.ts similarity index 100% rename from plugins/matrix-plugin/matrix-client.ts rename to matrix-client.ts diff --git a/plugins/matrix-plugin/package-lock.json b/package-lock.json similarity index 100% rename from plugins/matrix-plugin/package-lock.json rename to package-lock.json diff --git a/plugins/matrix-plugin/package.json b/package.json similarity index 100% rename from plugins/matrix-plugin/package.json rename to package.json diff --git a/plugins/matrix-plugin/session-manager.ts b/session-manager.ts similarity index 100% rename from plugins/matrix-plugin/session-manager.ts rename to session-manager.ts diff --git a/plugins/matrix-plugin/tsconfig.json b/tsconfig.json similarity index 100% rename from plugins/matrix-plugin/tsconfig.json rename to tsconfig.json diff --git a/plugins/matrix-plugin/types.ts b/types.ts similarity index 100% rename from plugins/matrix-plugin/types.ts rename to types.ts