549 lines
17 KiB
TypeScript
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)
|
|
}
|
|
}
|