import fs from "fs" import { MatrixBotClient } from "./matrix-client.js" import { SessionManager, encodeSessionId } from "./session-manager.js" import { loadConfig } from "./config-loader.js" import { log } from "./logger.js" log("Plugin module loaded") // ============================================================================= // V2 Support Detection // ============================================================================= function isV2Context(ctx: any): boolean { return !!(ctx?.session?.prompt && ctx?.event?.subscribe && ctx?.storage) } // ============================================================================= // Shared initialization logic // ============================================================================= async function runPlugin(ctx: any, options: any, sdkClient?: any) { const isV2 = isV2Context(ctx) const eventController = new AbortController() let cleanup: (() => void) | undefined try { const directory = ctx?.location?.directory || process.cwd() fs.appendFileSync("/tmp/opencode-matrix-plugin/setup-started.log", `[${new Date().toISOString()}] setup() called, dir=${directory}, v2=${isV2}, sdk=${!!sdkClient}\n`) log(`Plugin setup started, directory=${directory}, v2=${isV2}, sdk=${!!sdkClient}`) if (!options.homeserver) { log("Plugin disabled (no homeserver configured)") return () => {} } log(`Configuration loaded`) 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 () => { try { for await (const event of ctx.event.subscribe({ signal: eventController.signal })) { const eventType = typeof event === "object" && event !== null && "type" in event ? (event as { type: string }).type : "" if (eventType === "session.updated") { const sessionId = (event as { sessionID?: string }).sessionID 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) { if (String(err).includes("AbortError") || (err as { name?: string }).name === "AbortError") { // expected during cleanup } else { log(`Event subscription error: ${String(err)}`) } } })() } // Message handler matrix.on(async (data: any) => { const { context, query, sender, roomId, eventId, timestamp } = data // Check rate limit if (sessionManager.isRateLimited(sender)) { await matrix.sendNotice(context, "Rate limited. Please wait a moment.") return } // Check for active query if (sessionManager.hasActiveQuery(roomId, context.replyThreadRootId || context.eventId)) { await matrix.sendNotice(context, "A request is already running in this thread. Please wait.") return } const abortFn = () => { /* placeholder for cancel */ } const releaseQuery = sessionManager.markQueryActive(roomId, context.replyThreadRootId || context.eventId, abortFn) const threadRootId = context.replyThreadRootId || context.eventId const session = sessionManager.getOrCreateSession(roomId, threadRootId, eventId) session.isActive = true session.messageCount++ session.inputChars += query.length session.lastEventIds.set(context.replyThreadRootId || context.eventId, eventId) session.lastActivity = Date.now() let opencodeSessionId = session.opencodeSessionId // If no session ID yet, create one if (!opencodeSessionId) { log(`Creating new OpenCode session for ${roomId}:${threadRootId}`) try { if (isV2 && ctx?.session?.create) { // V2: can pass custom ID const encodedId = encodeSessionId(roomId, threadRootId) const createResult = await ctx.session.create({ id: encodedId, title: `Matrix: ${roomId.slice(0, 30)}...` }) opencodeSessionId = createResult?.id || createResult?.data?.id || encodedId log(`V2: Created OpenCode session ${opencodeSessionId} (custom ID: ${encodedId})`) } else if (sdkClient) { // V1: cannot pass custom ID, need to save mapping const createResult = await sdkClient.session.create({ body: { title: `Matrix: ${roomId.slice(0, 30)}...` } }) opencodeSessionId = createResult?.data?.id || createResult?.id || "" log(`V1: Created OpenCode session ${opencodeSessionId}`) } 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 session.lastActivity = Date.now() releaseQuery() await matrix.sendNotice(context, "Failed to create session. Check logs.") return } } try { if (isV2) { // V2: use ctx.session.prompt() const result = await ctx.session.prompt({ sessionID: opencodeSessionId, text: query, }) const responseText = result ? String(result) : "No response" session.outputChars += responseText.length await matrix.sendReply(context, responseText) } else if (sdkClient) { // 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 }] }, }) // 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) } } } sentMessageIds.set(opencodeSessionId, existingIds) } catch (err) { log(`V1: Error sending final messages: ${String(err)}`) } } catch (promptErr: any) { log(`V1 prompt error: ${String(promptErr)}`) const errMsg = promptErr?.response?.error?.message || promptErr?.error?.message || String(promptErr) await matrix.sendNotice(context, `Error: ${errMsg.slice(0, 200)}`) } } else { // No SDK client available log(`V1 mode without SDK client: Would send prompt to session ${opencodeSessionId}: ${query}`) const responseText = "Matrix plugin running in V1 mode — prompt forwarding requires OpenCode SDK client" session.outputChars += responseText.length await matrix.sendReply(context, responseText) } } catch (err: any) { log(`Error processing message: ${String(err)}`) const errMsg = err?.response?.error?.message || String(err) await matrix.sendNotice(context, `Error: ${errMsg.slice(0, 200)}`) } finally { session.isActive = false session.lastActivity = Date.now() sessionContextMap.delete(opencodeSessionId) releaseQuery() } }) // Bridge commands async function handleBridgeCommand( context: { roomId: string; replyThreadRootId?: string; threadRootId?: string; eventId: string }, cmdName: string, sendFn: (text: string) => Promise, ): Promise { const threadRootId = context.threadRootId || context.replyThreadRootId || context.eventId const key = sessionManager.getSessionKey(context.roomId, threadRootId) const sess = sessionManager.sessions.get(key) switch (cmdName) { case "status": { if (!sess) { await sendFn("No active session.") return } const age = Math.round((Date.now() - sess.lastActivity) / 60000) const msg = `Session status:\n- Messages: ${sess.messageCount}\n- Age: ~${age} min ago\n- Input: ~${sess.inputChars} chars\n- Output: ~${sess.outputChars} chars` await sendFn(msg) break } case "clear": case "reset": { sessionManager.removeSession(context.roomId, threadRootId) await sendFn("Session cleared. Next message will start a fresh session.") break } case "help": case "h": { const trigger = options.triggerPatterns?.[0] || "!oc " const botName = options.botName || "opencode" const mode = isV2 ? "OpenCode V2 (full)" : sdkClient ? "OpenCode V1 (full via SDK)" : "OpenCode V1 (limited — requires SDK client)" const msg = [ `${botName} - OpenCode Matrix Plugin`, "", "Bridge commands:", `- /h or /help - Show this help`, `- /status - Show current chat session info`, `- /clear or /reset - Delete current session history`, "", `Usage: ${trigger} `, `Or @${botName} `, `Or reply in a thread`, "", `Mode: ${mode}`, ].join("\n") await sendFn(msg) break } default: await sendFn(`Unknown command: ${cmdName}. Try /help`) } } // Override the matrix.on handler to also process commands const originalListener = (matrix as any)["listeners"][0] ;(matrix as any)["listeners"] = [] matrix.on(async (data: any) => { const { context, query, sender, eventId } = data // Handle bridge commands if (query.startsWith("/")) { const cmdName = query.slice(1).split(" ")[0].toLowerCase() const bridgeCommands = ["status", "clear", "reset", "help", "h"] if (bridgeCommands.includes(cmdName)) { const sendFn = (text: string) => matrix.sendNotice(context, text) await handleBridgeCommand(context, cmdName, sendFn) return } } // Forward to original handler if (originalListener) { await originalListener(data) } }) // 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()) } // Start Matrix client await matrix.start() log(`Plugin setup completed successfully (V2=${isV2}, sdk=${!!sdkClient})`) cleanup = () => { log("Plugin cleaning up...") eventController.abort() matrix.stop() } return cleanup } catch (err) { fs.appendFileSync("/tmp/opencode-matrix-plugin/setup-error.log", `[${new Date().toISOString()}] ${String(err)}\n`) log(`Setup error: ${String(err)}`) return () => { cleanup?.() } } } // ============================================================================= // V1 Server Entry Point — receives PluginInput with SDK client // ============================================================================= async function v1Server(input: any, options?: any) { try { log("V1 server entry point called") const { client, directory, worktree, project } = input const sdkClient = client if (!sdkClient) { log("V1: No SDK client in PluginInput") return {} } log(`V1: SDK client available from PluginInput`) // Load config using loadConfig (reads from matrix.json or env vars) const pluginOptions = options || {} const { options: configOptions } = await loadConfig(pluginOptions) // Create a context that runPlugin can use const v1Ctx = { location: { directory: directory || process.cwd(), project: project || { id: "v1" } }, storage: { set: async (_key: string, _value: any) => { /* V1 storage not available */ }, get: async (_key: string) => null, }, session: {}, // Not V2 context event: {}, // Not V2 context } await runPlugin(v1Ctx as any, configOptions, sdkClient) return {} } catch (err) { log(`V1 server error: ${String(err)}`) return {} } } // ============================================================================= // V2 Setup Function // ============================================================================= async function v2Setup(ctx: any) { try { const { options } = await loadConfig(ctx.options) return await runPlugin(ctx, options) } catch (err) { fs.appendFileSync("/tmp/opencode-matrix-plugin/setup-error.log", `[${new Date().toISOString()}] ${String(err)}\n`) log(`V2 setup error: ${String(err)}`) return () => {} } } // ============================================================================= // Dual V1/V2 Export // ============================================================================= // OpenCode 1.18.x calls server(PluginInput, options) where PluginInput contains SDK client. // V2 reads id + setup() and ignores server(). export default { id: "matrix-plugin", async setup(ctx: any) { const isV2 = isV2Context(ctx) if (isV2) { log("Using V2 path (full context)") return await v2Setup(ctx) } log("Using V1 path (server() will be called by OpenCode)") return await v1Server({}, {}) }, async server(input: any, options?: any) { log("V1 server() called by OpenCode") return await v1Server(input, options) }, }