matrix-plugin/dist/matrix-client.js

437 lines
17 KiB
JavaScript

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);
}
}