import fs from "fs"; import path from "path"; import os from "os"; // matrix-bot-sdk with native Rust crypto for E2EE import { AutojoinRoomsMixin, LogLevel, LogService, MatrixAuth, MatrixClient, MessageEvent, RichConsoleLogger, RustSdkCryptoStorageProvider, SimpleFsStorageProvider, } from "matrix-bot-sdk"; import { log } from "./logger.js"; // Storage paths function getStoragePaths(options) { const STORAGE_PATH = options.storagePath || path.join(os.homedir(), ".local", "share", "opencode-matrix-bot"); const STATE_STORAGE_PATH = path.join(STORAGE_PATH, "bot-state.json"); const CRYPTO_STORAGE_PATH = path.join(STORAGE_PATH, "crypto"); const TOKEN_FILE_PATH = path.join(STORAGE_PATH, "access_token"); return { STORAGE_PATH, STATE_STORAGE_PATH, CRYPTO_STORAGE_PATH, TOKEN_FILE_PATH }; } function writePrivateFileAtomically(filePath, content) { const pid = process.pid; const ts = Date.now(); const rand = Math.random().toString(36).substring(2, 10); const temporaryPath = `${filePath}.${pid}.${ts}.${rand}.tmp`; try { fs.writeFileSync(temporaryPath, content, { mode: 0o600, flag: "wx" }); fs.renameSync(temporaryPath, filePath); } catch (error) { try { fs.unlinkSync(temporaryPath); } catch { } throw error; } } function matrixErrorCode(error) { if (!error || typeof error !== "object") return ""; const candidate = error; if (typeof candidate.errcode === "string") return candidate.errcode; return typeof candidate.body?.errcode === "string" ? candidate.body.errcode : ""; } function isMatrixAuthenticationError(error) { const code = matrixErrorCode(error); return code === "M_UNKNOWN_TOKEN" || code === "M_MISSING_TOKEN"; } async function resolveMatrixAccessToken(options, dependencies) { if (options.explicitToken) { dependencies.log("Using access token from config/env"); return options.explicitToken; } if (options.savedToken) { try { const validation = await dependencies.validateToken(options.savedToken); if (!options.expectedUserId || validation.userId === options.expectedUserId) { dependencies.log("Using validated saved access token"); return options.savedToken; } dependencies.log(`Saved access token belongs to ${validation.userId}, not ${options.expectedUserId}; logging in again`); } catch (error) { if (!isMatrixAuthenticationError(error)) throw error; dependencies.log("Saved access token is no longer valid; logging in again"); } } if (!options.passwordConfigured) return null; const token = await dependencies.loginWithPassword(); if (!token) return null; dependencies.saveToken(token); return token; } // Thread helpers from matrix-thread-helpers.ts function extractThreadRootId(event) { const relatesTo = event?.content?.["m.relates_to"]; if (relatesTo?.rel_type === "m.thread" && relatesTo?.event_id) { return relatesTo.event_id; } return ""; } function extractBotNameQuery(text, botName) { const name = botName.trim(); if (!name) return null; const prefixes = [`@${name}`, name]; for (const prefix of prefixes) { if (text.slice(0, prefix.length).toLowerCase() !== prefix.toLowerCase()) continue; const separator = text.charAt(prefix.length); if (separator !== ":" && !/\s/.test(separator)) continue; return text.slice(prefix.length).replace(/^[:\s]+/, "").trim(); } return null; } function resolveThreadRoot(threadRootEventId, eventId) { return threadRootEventId || eventId; } function buildMatrixSessionId(roomId, replyThreadRootId, threadIsolation) { if (threadIsolation) { return `${roomId}:${replyThreadRootId}`; } return roomId; } function normalizeMatrixEventContext(input, threadIsolation) { const roomId = input.roomId; const eventId = input.eventId; const threadRootEventId = input.threadRootEventId || ""; const replyThreadRootId = resolveThreadRoot(threadRootEventId, eventId); return { roomId, sender: input.sender || "unknown", query: input.text || "", eventId, replyThreadRootId, sessionId: buildMatrixSessionId(roomId, replyThreadRootId, threadIsolation), }; } function buildThreadRelation(threadRootEventId, lastEventId) { return { rel_type: "m.thread", event_id: threadRootEventId, is_falling_back: true, "m.in_reply_to": { event_id: lastEventId }, }; } function shouldHandleThreadReply(input) { if (input.enabled === false) return false; const text = input.text.trim(); if (!text) return false; if (!input.threadRootEventId) return false; if (text.toLowerCase().startsWith(`${input.trigger.toLowerCase()} `)) return false; if (text.toLowerCase().startsWith(`${input.trigger.toLowerCase()}`)) return false; if (text.includes(input.botUserId)) return false; return true; } export class MatrixBotClient { options; matrix = null; userId = null; listeners = []; trigger; triggerPatterns; botName; threadIsolation; respondToThreadReplies; allowedUsers; ignoreRooms; ignoreUsers; formatHtml; constructor(options) { this.options = options; this.trigger = options.triggerPatterns?.[0] || "!oc "; this.triggerPatterns = options.triggerPatterns || []; this.botName = options.botName || "opencode"; this.threadIsolation = options.threadIsolation !== false; this.respondToThreadReplies = options.respondToThreadReplies !== false; this.allowedUsers = options.allowedUsers || []; this.ignoreRooms = options.ignoreRooms || []; this.ignoreUsers = options.ignoreUsers || []; this.formatHtml = options.formatHtml || false; } async start() { const { STORAGE_PATH, STATE_STORAGE_PATH, CRYPTO_STORAGE_PATH, TOKEN_FILE_PATH } = getStoragePaths(this.options); if (!this.options.accessToken && !this.options.password) { log("Error: Either MATRIX_ACCESS_TOKEN or MATRIX_PASSWORD must be set"); throw new Error("No credentials provided"); } log("Starting Matrix bot..."); log(` Homeserver: ${this.options.homeserver}`); log(` User: ${this.options.userId}`); log(` Storage: ${STORAGE_PATH}`); log(` E2EE: enabled (Rust crypto with SQLite)`); log(` Thread isolation: ${this.threadIsolation ? "on" : "off"}`); fs.mkdirSync(STORAGE_PATH, { recursive: true }); fs.mkdirSync(CRYPTO_STORAGE_PATH, { recursive: true }); // Get access token using the same pattern as chat-bridge let accessToken = await resolveMatrixAccessToken({ explicitToken: this.options.accessToken || "", savedToken: fs.existsSync(TOKEN_FILE_PATH) ? fs.readFileSync(TOKEN_FILE_PATH, "utf-8").trim() : "", passwordConfigured: Boolean(this.options.password), expectedUserId: this.options.userId || "", }, { validateToken: async (token) => { const client = new MatrixClient(this.options.homeserver || "https://matrix.org", token); const whoami = await client.getWhoAmI(); return { userId: whoami.user_id }; }, loginWithPassword: async () => { log("Logging in with password..."); try { const auth = new MatrixAuth(this.options.homeserver || "https://matrix.org"); const username = this.options.userId.split(":")[0].replace("@", ""); const client = await auth.passwordLogin(username, this.options.password, "OpenCode Matrix Plugin"); log("Password login successful"); return client.accessToken; } catch (err) { log(`Password login failed: ${err.message || err}`); return null; } }, saveToken: (token) => writePrivateFileAtomically(TOKEN_FILE_PATH, token), log: (message) => log(message), }); if (!accessToken) { log("Error: Could not obtain access token"); throw new Error("Could not obtain access token"); } // Configure logging LogService.setLogger(new RichConsoleLogger()); LogService.setLevel(LogLevel.INFO); LogService.muteModule("Metrics"); // Create client with storage providers (same as chat-bridge) const stateStorage = new SimpleFsStorageProvider(STATE_STORAGE_PATH); const cryptoStorage = new RustSdkCryptoStorageProvider(CRYPTO_STORAGE_PATH, 0 /* RustSdkCryptoStoreType.Sqlite */); this.matrix = new MatrixClient(this.options.homeserver || "https://matrix.org", accessToken, stateStorage, cryptoStorage); // Setup auto-join mixin (same as chat-bridge) if (this.options.autoJoin !== false) { AutojoinRoomsMixin.setupOnClient(this.matrix); } // Handle decryption failures this.matrix.on("room.failed_decryption", async (roomId, event, error) => { log(`[CRYPTO] Failed to decrypt in ${roomId}: ${error.message}`); }); // Get user ID this.userId = await this.matrix.getUserId(); log(`Matrix bot started as ${this.userId}`); // Handle messages this.matrix.on("room.message", this.handleRoomMessage.bind(this)); // Start syncing (same as chat-bridge: await this.matrix.start()) await this.matrix.start(); log("Matrix bot listening for messages"); } async handleRoomMessage(roomId, event) { if (!this.matrix || !this.userId) return; const message = new MessageEvent(event); if (message.messageType !== "m.text") return; if (message.sender === this.userId) return; // Check ignored rooms if (this.ignoreRooms.length > 0 && this.ignoreRooms.includes(roomId)) { return; } // Check ignored users if (this.ignoreUsers.length > 0 && this.ignoreUsers.includes(message.sender)) { return; } // Check allowed users if (this.allowedUsers.length > 0 && !this.allowedUsers.includes(message.sender)) { return; } const body = message.textBody?.trim(); if (!body) return; // Deduplicate events if (this.isDuplicateEvent(event.event_id || `${roomId}:${Date.now()}`)) return; const threadRootEventId = extractThreadRootId(event); const context = normalizeMatrixEventContext({ roomId, sender: message.sender, text: body, eventId: event.event_id, threadRootEventId, }, this.threadIsolation); // Check if this is a DM const members = await this.matrix.getJoinedRoomMembers(roomId); const isDM = members.length === 2; // Extract query let query = ""; const botNameQuery = extractBotNameQuery(body, this.botName); if (this.triggerPatterns.length === 0) { // No triggers configured — intercept all messages query = body; } else { // Check all trigger patterns let matchedTrigger = ""; for (const pattern of this.triggerPatterns) { if (body.startsWith(pattern + " ")) { matchedTrigger = pattern; query = body.slice(pattern.length + 1).trim(); break; } else if (body.startsWith(pattern)) { matchedTrigger = pattern; query = body.slice(pattern.length).trim(); break; } } // If no trigger matched, check bot mention if (!matchedTrigger) { if (body.includes(this.userId)) { query = body.replace(this.userId, "").trim(); } else if (botNameQuery !== null) { query = botNameQuery; } else if (this.threadIsolation && shouldHandleThreadReply({ enabled: this.respondToThreadReplies, text: body, threadRootEventId, trigger: this.triggerPatterns[0], botUserId: this.userId, })) { query = body; log(`[THREAD] ${message.sender} in ${context.sessionId}: ${body}`); } else { return; } } } query = query.replace(/^[:\s]+/, "").trim(); if (!query) return; log(`[MSG] ${message.sender} in ${context.sessionId}: ${body}`); const timestamp = event.ts; for (const listener of this.listeners) { try { await listener({ context, query, sender: message.sender, roomId, eventId: event.event_id, timestamp }); } catch (e) { log(`Error in message listener: ${e}`); } } } processedEvents = new Set(); isDuplicateEvent(eventId) { if (!eventId) return false; if (this.processedEvents.has(eventId)) return true; if (this.processedEvents.size > 10000) { this.processedEvents.clear(); } this.processedEvents.add(eventId); return false; } async stop() { log("Stopping..."); if (this.matrix) { this.matrix.stop(); this.matrix = null; } log("Stopped."); } async sendMessage(roomId, text) { if (!this.matrix) throw new Error("Matrix client not started"); if (this.formatHtml) { const { marked } = await import("marked"); const html = await marked.parse(text); await this.matrix.sendMessage(roomId, { msgtype: "m.text", body: text, format: "org.matrix.custom.html", formatted_body: html, }); } else { await this.matrix.sendText(roomId, text); } } async sendReply(context, text) { if (!this.matrix) return null; try { if (this.threadIsolation) { const threadRoot = context.replyThreadRootId || context.eventId; const relation = buildThreadRelation(threadRoot, context.eventId); let content; if (this.formatHtml) { const { marked } = await import("marked"); const html = await marked.parse(text); content = { msgtype: "m.text", body: text, format: "org.matrix.custom.html", formatted_body: html, "m.relates_to": relation, }; } else { content = { msgtype: "m.text", body: text, "m.relates_to": relation, }; } const eventId = await this.matrix.sendMessage(context.roomId, content); return eventId; } else { await this.sendMessage(context.roomId, text); return null; } } catch (err) { log(`Failed to send reply to ${context.roomId}: ${err}`); return null; } } async sendNotice(context, text) { if (!this.matrix) return; try { if (this.threadIsolation) { const threadRoot = context.replyThreadRootId || context.eventId; const relation = buildThreadRelation(threadRoot, context.eventId); await this.matrix.sendMessage(context.roomId, { msgtype: "m.notice", body: text, "m.relates_to": relation, }); } else { await this.matrix.sendNotice(context.roomId, text); } } catch (err) { log(`Failed to send notice to ${context.roomId}: ${err}`); } } on(listener) { this.listeners.push(listener); } }