matrix-plugin/matrix-client.ts

549 lines
17 KiB
TypeScript

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<void>
}
// 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<MatrixTokenValidation>
loginWithPassword: () => Promise<string | null>
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<string | null> {
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<void> {
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<void> {
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<string>()
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<void> {
log("Stopping...")
if (this.matrix) {
this.matrix.stop()
this.matrix = null
}
log("Stopped.")
}
private async sendMessage(roomId: string, text: string): Promise<void> {
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<string | null> {
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<void> {
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)
}
}