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 { StoreType as RustSdkCryptoStoreType } from "@matrix-org/matrix-sdk-crypto-nodejs" import { log } from "./logger.js" export interface MatrixOptions { homeserver?: string userId?: string accessToken?: string password?: string deviceId?: string autoJoin?: boolean triggerPatterns?: string[] ignoreRooms?: string[] ignoreUsers?: string[] allowedUsers?: string[] formatHtml?: boolean threadIsolation?: boolean respondToThreadReplies?: boolean rateLimitSeconds?: number botName?: string storagePath?: string } export interface MatrixEventContext { sessionId: string roomId: string sender: string query: string replyThreadRootId?: string eventId: string } export interface MessageListener { (data: { context: MatrixEventContext query: string sender: string roomId: string eventId: string timestamp: number }): Promise } // Storage paths function getStoragePaths(options: MatrixOptions) { 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 } } // Token helpers from matrix-auth.ts pattern interface MatrixTokenValidation { userId: string } interface MatrixAccessTokenOptions { explicitToken: string savedToken: string passwordConfigured: boolean expectedUserId: string } interface MatrixAccessTokenDependencies { validateToken: (token: string) => Promise loginWithPassword: () => Promise saveToken: (token: string) => void log: (message: string) => void } function writePrivateFileAtomically(filePath: string, content: string): void { 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: unknown): string { if (!error || typeof error !== "object") return "" const candidate = error as { errcode?: unknown; body?: { errcode?: unknown } } if (typeof candidate.errcode === "string") return candidate.errcode return typeof candidate.body?.errcode === "string" ? candidate.body.errcode : "" } function isMatrixAuthenticationError(error: unknown): boolean { const code = matrixErrorCode(error) return code === "M_UNKNOWN_TOKEN" || code === "M_MISSING_TOKEN" } async function resolveMatrixAccessToken( options: MatrixAccessTokenOptions, dependencies: MatrixAccessTokenDependencies, ): Promise { 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: any): string { const relatesTo = event?.content?.["m.relates_to"] if (relatesTo?.rel_type === "m.thread" && relatesTo?.event_id) { return relatesTo.event_id } return "" } function extractBotNameQuery(text: string, botName: string): string | null { 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: string, eventId: string): string { return threadRootEventId || eventId } function buildMatrixSessionId(roomId: string, replyThreadRootId: string, threadIsolation: boolean): string { if (threadIsolation) { return `${roomId}:${replyThreadRootId}` } return roomId } function normalizeMatrixEventContext( input: { roomId: string; sender?: string; text?: string; eventId: string; threadRootEventId?: string }, threadIsolation: boolean, ): MatrixEventContext { 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: string, lastEventId: string): object { return { rel_type: "m.thread", event_id: threadRootEventId, is_falling_back: true, "m.in_reply_to": { event_id: lastEventId }, } } function shouldHandleThreadReply(input: { enabled?: boolean text: string threadRootEventId: string trigger: string botUserId: string }): boolean { 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 { private options: MatrixOptions private matrix: MatrixClient | null = null public userId: string | null = null private listeners: MessageListener[] = [] private readonly trigger: string private readonly triggerPatterns: string[] private readonly botName: string private readonly threadIsolation: boolean private readonly respondToThreadReplies: boolean private readonly allowedUsers: string[] private readonly ignoreRooms: string[] private readonly ignoreUsers: string[] private readonly formatHtml: boolean constructor(options: MatrixOptions) { 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(): Promise { 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: any) { 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, 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: string, event: any, error: 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") } private async handleRoomMessage(roomId: string, event: any): Promise { 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}`) } } } private processedEvents = new Set() isDuplicateEvent(eventId: string | undefined): boolean { 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(): Promise { log("Stopping...") if (this.matrix) { this.matrix.stop() this.matrix = null } log("Stopped.") } private async sendMessage(roomId: string, text: string): Promise { 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: MatrixEventContext, text: string): Promise { if (!this.matrix) return null try { if (this.threadIsolation) { const threadRoot = context.replyThreadRootId || context.eventId const relation = buildThreadRelation(threadRoot, context.eventId) let content: any 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: MatrixEventContext, text: string): Promise { 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: MessageListener): void { this.listeners.push(listener) } }