From 4c0ad70e47ea1b13d692cce9b8b1265e71ed7d9d Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=D0=91=D0=BE=D1=80=D0=BE=D0=B4=D0=B8=D0=BD=20=D0=A0=D0=BE?= =?UTF-8?q?=D0=BC=D0=B0=D0=BD?= Date: Wed, 30 Sep 2026 21:39:03 +0300 Subject: [PATCH] fix: compile TS to JS and set main to dist/index.js for git plugin installation - Set noEmit to false in tsconfig.json - Change main entry point from index.ts to dist/index.js - Add files field to package.json - Pre-compile TypeScript to JavaScript --- dist/config-loader.js | 197 +++++++++++ dist/index.js | 759 ++++++++++++++++++++++++++++++++++++++++ dist/logger.js | 67 ++++ dist/matrix-client.js | 436 +++++++++++++++++++++++ dist/session-manager.js | 166 +++++++++ dist/types.js | 1 + package.json | 6 +- tsconfig.json | 3 +- 8 files changed, 1632 insertions(+), 3 deletions(-) create mode 100644 dist/config-loader.js create mode 100644 dist/index.js create mode 100644 dist/logger.js create mode 100644 dist/matrix-client.js create mode 100644 dist/session-manager.js create mode 100644 dist/types.js diff --git a/dist/config-loader.js b/dist/config-loader.js new file mode 100644 index 0000000..ffc4dcb --- /dev/null +++ b/dist/config-loader.js @@ -0,0 +1,197 @@ +import fs from "fs"; +import path from "path"; +import os from "os"; +import { log } from "./logger.js"; +const CONFIG_FILENAMES = ["matrix.json", "matrix.jsonc"]; +function mergeDefaults(raw, envOverrides) { + const base = { + homeserver: envOverrides.homeserver || "https://matrix.org", + userId: envOverrides.userId, + accessToken: envOverrides.accessToken, + password: envOverrides.password, + deviceId: envOverrides.deviceId || "opencode-matrix-plugin", + autoJoin: envOverrides.autoJoin !== false, + triggerPatterns: envOverrides.triggerPatterns || ["!oc "], + ignoreRooms: envOverrides.ignoreRooms || [], + ignoreUsers: envOverrides.ignoreUsers || [], + allowedUsers: envOverrides.allowedUsers || [], + formatHtml: envOverrides.formatHtml || false, + threadIsolation: envOverrides.threadIsolation !== false, + respondToThreadReplies: envOverrides.respondToThreadReplies !== false, + rateLimitSeconds: envOverrides.rateLimitSeconds || 5, + botName: envOverrides.botName || "opencode", + storagePath: envOverrides.storagePath, + }; + if (!raw) + return base; + return { + homeserver: raw.homeserver || base.homeserver, + userId: raw.userId || base.userId, + accessToken: raw.accessToken || base.accessToken, + password: raw.password || base.password, + deviceId: raw.deviceId || base.deviceId, + autoJoin: raw.autoJoin !== undefined ? raw.autoJoin : base.autoJoin, + triggerPatterns: raw.triggerPatterns || base.triggerPatterns, + ignoreRooms: raw.ignoreRooms || base.ignoreRooms, + ignoreUsers: raw.ignoreUsers || base.ignoreUsers, + allowedUsers: raw.allowedUsers || base.allowedUsers, + formatHtml: raw.formatHtml !== undefined ? raw.formatHtml : base.formatHtml, + threadIsolation: raw.threadIsolation !== undefined ? raw.threadIsolation : base.threadIsolation, + respondToThreadReplies: raw.respondToThreadReplies !== undefined ? raw.respondToThreadReplies : base.respondToThreadReplies, + rateLimitSeconds: raw.rateLimitSeconds || base.rateLimitSeconds, + botName: raw.botName || base.botName, + storagePath: raw.storagePath || base.storagePath, + }; +} +function tryReadConfig(dir) { + for (const filename of CONFIG_FILENAMES) { + const filepath = path.join(dir, filename); + if (!fs.existsSync(filepath)) + continue; + try { + const content = fs.readFileSync(filepath, "utf-8"); + const parsed = JSON.parse(content); + if (typeof parsed === "object" && parsed !== null) { + return parsed; + } + } + catch (e) { + log(`Failed to read config: ${filepath}`); + return null; + } + } + return null; +} +function findProjectConfigDir() { + let dir = process.cwd(); + const root = path.parse(dir).root; + while (dir && dir !== root) { + const opencodeDir = path.join(dir, ".opencode"); + if (fs.existsSync(opencodeDir) && fs.statSync(opencodeDir).isDirectory()) { + return opencodeDir; + } + dir = path.dirname(dir); + } + const currentOpencode = path.join(process.cwd(), ".opencode"); + if (fs.existsSync(currentOpencode) && fs.statSync(currentOpencode).isDirectory()) { + return currentOpencode; + } + return null; +} +function findGlobalConfigDir() { + return path.join(os.homedir(), ".config", "opencode"); +} +const DEFAULT_CONFIG_CONTENT = JSON.stringify({ + homeserver: "https://matrix.org", + userId: "", + accessToken: "", + password: "", + deviceId: "opencode-matrix-plugin", + autoJoin: true, + triggerPatterns: ["!oc "], + ignoreRooms: [], + ignoreUsers: [], + allowedUsers: [], + formatHtml: false, + threadIsolation: true, + respondToThreadReplies: true, + rateLimitSeconds: 5, + botName: "opencode", + enabled: true, +}, null, 2); +async function ensureGlobalConfig(globalConfigDir) { + log(`ensureGlobalConfig called, dir=${globalConfigDir}`); + for (const filename of CONFIG_FILENAMES) { + const filepath = path.join(globalConfigDir, filename); + if (fs.existsSync(filepath)) { + log(`config already exists: ${filepath}`); + return filepath; + } + } + log("no config found, creating default"); + try { + log("mkdirSync globalConfigDir"); + fs.mkdirSync(globalConfigDir, { recursive: true }); + const filepath = path.join(globalConfigDir, "matrix.json"); + log("writing file: " + filepath); + fs.writeFileSync(filepath, DEFAULT_CONFIG_CONTENT, { mode: 0o600 }); + log("Created default config: " + filepath); + return filepath; + } + catch (err) { + log("Failed to create default config: " + err); + return null; + } +} +export async function loadConfig(pluginOptions = {}) { + pluginOptions = pluginOptions || {}; + const envOverrides = { + homeserver: process.env.MATRIX_HOMESERVER, + userId: process.env.MATRIX_USER_ID, + accessToken: process.env.MATRIX_ACCESS_TOKEN, + password: process.env.MATRIX_PASSWORD, + deviceId: process.env.MATRIX_DEVICE_ID || undefined, + triggerPatterns: process.env.MATRIX_TRIGGER ? [process.env.MATRIX_TRIGGER] : undefined, + allowedUsers: process.env.MATRIX_ALLOWED_USERS ? process.env.MATRIX_ALLOWED_USERS.split(",").map(s => s.trim()).filter(Boolean) : undefined, + storagePath: process.env.MATRIX_STORAGE_PATH, + }; + const projectConfigDir = findProjectConfigDir(); + log(`projectConfigDir=${projectConfigDir}`); + let rawConfig = null; + let configPath = null; + if (projectConfigDir) { + const projectConfig = tryReadConfig(projectConfigDir); + log(`projectConfig=${projectConfig ? "found" : "null"}`); + if (projectConfig) { + rawConfig = projectConfig; + const foundFile = CONFIG_FILENAMES.find((f) => fs.existsSync(path.join(projectConfigDir, f))); + if (foundFile) { + configPath = path.join(projectConfigDir, foundFile); + } + } + } + const globalConfigDir = findGlobalConfigDir(); + const globalConfig = tryReadConfig(globalConfigDir); + log(`globalConfigDir=${globalConfigDir}`); + log(`globalConfig=${globalConfig ? "found" : "null"}`); + log(`rawConfig=${rawConfig ? "found" : "null"}`); + if (!rawConfig && globalConfig) { + rawConfig = globalConfig; + configPath = path.join(globalConfigDir, "matrix.json"); + log("using global config"); + } + if (!rawConfig) { + log("no config found, calling ensureGlobalConfig"); + const createdPath = await ensureGlobalConfig(globalConfigDir); + log(`ensureGlobalConfig returned=${createdPath}`); + if (createdPath) { + configPath = createdPath; + } + } + if (rawConfig && rawConfig.enabled === false) { + return { + options: { ...mergeDefaults(null, envOverrides), homeserver: "" }, + configPath, + }; + } + const merged = mergeDefaults(rawConfig, envOverrides); + const finalOptions = { + homeserver: pluginOptions.homeserver || merged.homeserver, + userId: pluginOptions.userId || merged.userId, + accessToken: pluginOptions.accessToken || merged.accessToken, + password: pluginOptions.password || merged.password, + deviceId: pluginOptions.deviceId || merged.deviceId, + autoJoin: pluginOptions.autoJoin !== undefined ? pluginOptions.autoJoin : merged.autoJoin, + triggerPatterns: pluginOptions.triggerPatterns || merged.triggerPatterns, + ignoreRooms: pluginOptions.ignoreRooms || merged.ignoreRooms, + ignoreUsers: pluginOptions.ignoreUsers || merged.ignoreUsers, + allowedUsers: pluginOptions.allowedUsers || merged.allowedUsers, + formatHtml: pluginOptions.formatHtml !== undefined ? pluginOptions.formatHtml : merged.formatHtml, + threadIsolation: pluginOptions.threadIsolation !== undefined ? pluginOptions.threadIsolation : merged.threadIsolation, + respondToThreadReplies: pluginOptions.respondToThreadReplies !== undefined ? pluginOptions.respondToThreadReplies : merged.respondToThreadReplies, + rateLimitSeconds: pluginOptions.rateLimitSeconds || merged.rateLimitSeconds, + botName: pluginOptions.botName || merged.botName, + storagePath: pluginOptions.storagePath || merged.storagePath, + }; + return { options: finalOptions, configPath }; +} diff --git a/dist/index.js b/dist/index.js new file mode 100644 index 0000000..6ceb994 --- /dev/null +++ b/dist/index.js @@ -0,0 +1,759 @@ +import fs from "fs"; +import os from "os"; +import path from "path"; +import * as acorn from "acorn"; +import * as acornWalk from "acorn-walk"; +import { MatrixBotClient } from "./matrix-client.js"; +import { SessionManager } from "./session-manager.js"; +import { loadConfig } from "./config-loader.js"; +import { setAppLog, log, logError } from "./logger.js"; +// ============================================================================= +// Plugin data directory — uses the same default as matrix-client +// ============================================================================= +function ensureDataDir(dir) { + if (!fs.existsSync(dir)) { + fs.mkdirSync(dir, { recursive: true }); + } +} +// ============================================================================= +// Global mutex — prevents double startup (server + client both load plugins) +// ============================================================================= +let LOCK_FILE = ""; +function tryAcquireLock() { + try { + // Check if an existing lock is stale (process no longer running) + if (fs.existsSync(LOCK_FILE)) { + const existingPid = parseInt(fs.readFileSync(LOCK_FILE, "utf-8").trim(), 10); + if (existingPid && existingPid !== process.pid) { + try { + process.kill(existingPid, 0); // throws if process doesn't exist + } + catch { + // Process dead, remove stale lock + fs.unlinkSync(LOCK_FILE); + } + } + else if (existingPid === process.pid) { + // Same process re-entering (shouldn't happen, but be safe) + return false; + } + else { + return false; // Valid lock held by another process + } + } + fs.writeFileSync(LOCK_FILE, String(process.pid), { mode: 0o644 }); + return true; + } + catch { + return false; + } +} +function releaseLock() { + try { + fs.unlinkSync(LOCK_FILE); + } + catch { + // ignore + } +} +// ============================================================================= +// Plugin Setup +// ============================================================================= +async function runPlugin(ctx, options, dataDir) { + ensureDataDir(dataDir); + LOCK_FILE = `${dataDir}/lock`; + const directory = ctx?.location?.directory || process.cwd(); + log(`Plugin setup started, directory=${directory}`); + // Acquire global lock — only one instance (server or client) starts the bot + const locked = tryAcquireLock(); + if (!locked) { + log("Another instance already running (lock held), skipping bot startup"); + return () => { }; + } + const eventController = new AbortController(); + let cleanup; + try { + log(`Configuration loaded`); + if (!options.homeserver) { + log("Plugin disabled (no homeserver configured)"); + return () => { }; + } + log(`Configuration loaded`); + const matrix = new MatrixBotClient(options); + const sessionManager = new SessionManager(options.rateLimitSeconds || 5); + // Track Matrix context per OpenCode session + const sessionContextMap = new Map(); + // ============================================================================= + // Pending prompt responses — maps sessionID -> { resolve, reject } + // ============================================================================= + const pendingResponses = new Map(); + const sessionsBySessionID = new Map(); + /** + * Parse the CodeMode JS script and extract all tool calls. + * Uses acorn AST parser + walk — same approach as OpenCode's CodeMode interpreter. + * Does NOT execute the code, only parses and analyzes the AST. + */ + function parseCodeModeTools(code) { + const results = []; + const seen = new Set(); + // Parse the code into an AST using acorn (same parser as CodeMode) + const ast = acorn.parse(code, { + ecmaVersion: "latest", + sourceType: "script", + allowReturnOutsideFunction: true, + allowAwaitOutsideFunction: true, + locations: true, + }); + // Walk the AST and find all CallExpression nodes for tools.xxx.yyy(...) + acornWalk.simple(ast, { + CallExpression(node) { + // Check if callee is tools..(...) + const toolCall = extractToolCall(node); + if (!toolCall) + return; + const { path, args } = toolCall; + const name = path.join("_"); + if (seen.has(name)) + return; + seen.add(name); + // Evaluate the arguments to get their string representation + const argsStr = evaluateArgs(args); + results.push({ name, args: argsStr }); + }, + }); + return results; + } + /** + * Check if a CallExpression is a tool call: tools..(...) + * Returns { path, args } or null. + */ + function extractToolCall(node) { + const callee = node.callee; + // Must be a MemberExpression: tools.xxx.yyy + if (callee.type !== "MemberExpression") + return null; + if (callee.computed) + return null; + const prop = callee.property; + if (prop.type !== "Identifier") + return null; + const toolName = prop.name; + // Object must be tools.xxx (nested MemberExpression) or just "tools" + let obj = callee.object; + const path = []; + // Walk up the nested MemberExpression chain + while (obj.type === "MemberExpression" && !obj.computed) { + const innerProp = obj.property; + if (innerProp.type !== "Identifier") + return null; + path.unshift(innerProp.name); + obj = obj.object; + } + // The base must be the identifier "tools" + if (obj.type !== "Identifier" || obj.name !== "tools") + return null; + path.unshift(toolName); + return { path, args: (node.arguments ?? []) }; + } + /** + * Evaluate arguments to a string representation. + * Handles literals, identifiers, objects, arrays, and template literals. + * Matches how CodeMode's interpreter evaluates expression values. + */ + function evaluateArgs(args) { + if (!args || args.length === 0) + return ""; + const values = []; + for (const arg of args) { + values.push(evaluateExpression(arg)); + } + return values.join(", "); + } + /** + * Evaluate an AST expression node to its string value. + * Handles the expression types used in CodeMode tool calls. + */ + function evaluateExpression(node) { + switch (node.type) { + // Literals + case "Literal": + return String(node.value); + // String/number/boolean/null/regexp literals + case "Identifier": { + const name = node.name; + // Known globals with known values + if (name === "undefined") + return "undefined"; + if (name === "true") + return "true"; + if (name === "false") + return "false"; + if (name === "null") + return "null"; + // For variables, show the variable name + return name; + } + // Object literals: { key: value, ... } + case "ObjectExpression": { + const parts = []; + for (const prop of node.properties ?? []) { + if (prop.type === "Property") { + const key = prop.key.type === "Identifier" ? prop.key.name : String(prop.key.value); + const value = evaluateExpression(prop.value); + parts.push(`${key}: ${value}`); + } + } + return `{ ${parts.join(", ")} }`; + } + // Array literals: [value1, value2, ...] + case "ArrayExpression": { + const elements = (node.elements ?? []).map((el) => el ? evaluateExpression(el) : ""); + return `[${elements.join(", ")}]`; + } + // Template literals: `text ${expr} text` + case "TemplateLiteral": { + let result = "`"; + const quasis = node.quasis ?? []; + const expressions = node.expressions ?? []; + for (let i = 0; i < quasis.length; i++) { + result += quasis[i]?.value?.raw ?? ""; + if (expressions[i]) { + result += "${" + evaluateExpression(expressions[i]) + "}"; + } + } + result += "`"; + return result; + } + // Binary expressions: a + b, a === b, etc. + case "BinaryExpression": { + const left = evaluateExpression(node.left); + const right = evaluateExpression(node.right); + return `${left} ${node.operator} ${right}`; + } + // Conditional: a ? b : c + case "ConditionalExpression": { + return `${evaluateExpression(node.test)} ? ${evaluateExpression(node.consequent)} : ${evaluateExpression(node.alternate)}`; + } + // Function expressions (arrow functions, etc.) + case "ArrowFunctionExpression": + case "FunctionExpression": + return "function"; + // Call expressions: other().method() + case "CallExpression": { + if (node.callee?.type === "MemberExpression") { + const member = node.callee; + if (member.object?.type === "Identifier") { + return `${member.object.name}.${member.property?.name}(${evaluateArgs(node.arguments ?? [])})`; + } + } + return "call"; + } + // Unary expressions: !a, -b, typeof c + case "UnaryExpression": { + const arg = evaluateExpression(node.argument); + return `${node.operator}${arg}`; + } + // Logical expressions: a && b, a || b + case "LogicalExpression": { + return `${evaluateExpression(node.left)} ${node.operator} ${evaluateExpression(node.right)}`; + } + // Await expressions + case "AwaitExpression": { + return `await ${evaluateExpression(node.argument)}`; + } + // Spread: ...arr + case "SpreadElement": { + return `...${evaluateExpression(node.argument)}`; + } + // Default: show "(expr)" + default: + return `(expr:${node.type})`; + } + } + // ============================================================================= + // Event subscription — collects step data and sends messages per step + // ============================================================================= + void (async () => { + try { + for await (const event of ctx.event.subscribe({ signal: eventController.signal })) { + let eventData = event; + if (typeof event === "string") { + try { + eventData = JSON.parse(event); + } + catch { + continue; + } + } + const eventType = eventData?.type || ""; + const sessionId = eventData?.data?.sessionID || ""; + const assistantMsgID = eventData?.data?.assistantMessageID || ""; + // Skip all events after session is done + const state = sessionsBySessionID.get(sessionId); + if (state?.done) + continue; + // ── Step lifecycle ────────────────────────────────────────────── + if (eventType === "session.step.started") { + // SessionState created in message handler after prompt + if (!sessionsBySessionID.has(sessionId)) { + log(`step.started for unknown session ${sessionId}`); + continue; + } + const state = sessionsBySessionID.get(sessionId); + state.currentStepID = assistantMsgID; + state.steps.set(assistantMsgID, { + assistantMessageID: assistantMsgID, + reasoning: "", + tools: [], + text: "", + }); + continue; + } + // ── Reasoning ─────────────────────────────────────────────────── + if (eventType === "session.reasoning.delta" && assistantMsgID) { + const state = sessionsBySessionID.get(sessionId); + const step = state?.steps.get(assistantMsgID); + if (step) + step.reasoning += eventData?.data?.delta || ""; + continue; + } + if (eventType === "session.reasoning.ended" && assistantMsgID) { + const state = sessionsBySessionID.get(sessionId); + const step = state?.steps.get(assistantMsgID); + if (step) + step.reasoning = eventData?.data?.text || ""; + continue; + } + // ── Tool calls ────────────────────────────────────────────────── + if (eventType === "session.tool.input.started" && assistantMsgID) { + const toolID = eventData?.data?.id || ""; + const name = eventData?.data?.name || ""; + const state = sessionsBySessionID.get(sessionId); + const step = state?.steps.get(assistantMsgID); + if (step && toolID) { + step.tools.push({ toolID, name, args: "", result: "", failed: false }); + } + continue; + } + if (eventType === "session.tool.input.ended" && assistantMsgID) { + const toolID = eventData?.data?.id || ""; + const args = eventData?.data?.text || ""; + const state = sessionsBySessionID.get(sessionId); + const step = state?.steps.get(assistantMsgID); + if (step && toolID) { + const tool = step.tools.find(t => t.toolID === toolID); + if (tool) + tool.args = args; + } + // Parse CodeMode execute tool — extract child tool calls from JS code + if (toolID && step && args) { + let parsedArgs; + try { + parsedArgs = JSON.parse(args); + } + catch { + continue; + } + const code = parsedArgs?.code || ""; + if (code) { + const childTools = parseCodeModeTools(code); + for (const child of childTools) { + if (!step.tools.find(t => t.name === child.name)) { + step.tools.push({ + toolID: `${toolID}/${child.name}`, + name: child.name, + args: child.args, + result: "", + failed: false, + }); + } + } + } + } + continue; + } + if (eventType === "session.tool.success" && assistantMsgID) { + const toolID = eventData?.data?.id || ""; + const content = eventData?.data?.content; + let result = ""; + if (typeof content === "string") { + result = content; + } + else if (Array.isArray(content)) { + for (const part of content) { + if (part?.type === "text" && part?.text) + result += part.text; + } + } + const state = sessionsBySessionID.get(sessionId); + const step = state?.steps.get(assistantMsgID); + if (step && toolID) { + const tool = step.tools.find(t => t.toolID === toolID); + if (tool) + tool.result = result; + } + continue; + } + if (eventType === "session.tool.failed" && assistantMsgID) { + const toolID = eventData?.data?.id || ""; + const state = sessionsBySessionID.get(sessionId); + const step = state?.steps.get(assistantMsgID); + if (step && toolID) { + const tool = step.tools.find(t => t.toolID === toolID); + if (tool) { + tool.failed = true; + tool.result = String(eventData?.data?.error || "Unknown error"); + } + } + continue; + } + // ── Text output ──────────────────────────────────────────────── + if (eventType === "session.text.ended" && assistantMsgID) { + const state = sessionsBySessionID.get(sessionId); + // Skip text events after final step + if (state?.done) + continue; + const step = state?.steps.get(assistantMsgID); + if (step) { + step.text = eventData?.data?.text || ""; + if (state) + state.finalResponse = step.text; + } + continue; + } + // ── Step ended — send message for this step ──────────────────── + if (eventType === "session.step.ended" && sessionId && assistantMsgID) { + const state = sessionsBySessionID.get(sessionId); + // Skip steps after final (finish:stop) — they are re-projections/duplicates + if (state?.done) + continue; + const step = state?.steps.get(assistantMsgID); + const finish = eventData?.data?.finish || ""; + if (step && state) { + const msg = formatStepMessage(step); + await sendMessageToMatrix(state.context, msg); + log(`Step ${step.tools.length} tool(s), ${step.text.length} chars text → sent to ${state.context.roomId}`); + } + // If this is the final step, mark done and resolve the pending promise + if (finish === "stop" && state) { + state.done = true; + const pending = pendingResponses.get(sessionId); + if (pending && state.finalResponse) { + pendingResponses.delete(sessionId); + pending.resolve(state.finalResponse); + } + sessionsBySessionID.delete(sessionId); + } + continue; + } + // ── Execution errors ─────────────────────────────────────────── + if (eventType === "session.execution.failed" && sessionId) { + const pending = pendingResponses.get(sessionId); + if (pending) { + pendingResponses.delete(sessionId); + pending.reject(new Error("Session execution failed")); + } + sessionsBySessionID.delete(sessionId); + continue; + } + if (eventType === "session.execution.interrupted" && sessionId) { + const pending = pendingResponses.get(sessionId); + if (pending) { + pendingResponses.delete(sessionId); + pending.reject(new Error("Session execution interrupted")); + } + sessionsBySessionID.delete(sessionId); + continue; + } + } + } + catch (err) { + if (String(err).includes("AbortError") || err.name === "AbortError") { + // expected during cleanup + } + else { + log(`Event subscription error: ${String(err)}`); + } + } + })(); + // ============================================================================= + // Format step into Matrix message + // ============================================================================= + function formatStepMessage(step) { + const parts = []; + // Reasoning in spoiler + if (step.reasoning) { + parts.push(`
\nThinking\n\n${step.reasoning}\n\n
`); + } + // Tool calls + if (step.tools.length > 0) { + // Check if any tool is an execute with child tool calls + const executeTool = step.tools.find((t) => t.name === "execute" && t.toolID.includes("/")); + const childTools = executeTool + ? step.tools.filter((t) => t.toolID.startsWith(executeTool.toolID + "/")) + : []; + if (childTools.length > 0) { + // Execute has child tool calls — show children, then execute result + for (const child of childTools) { + parts.push(`
\n🔧 ${child.name}\n\n`); + parts.push(`Args:\n\`\`\`\n${child.args}\n\`\`\`\n\n`); + if (child.failed) { + parts.push(`Failed: ${child.result}\n\n`); + } + else if (child.result) { + parts.push(`Result:\n\`\`\`\n${child.result}\n\`\`\`\n\n`); + } + parts.push(`
`); + } + // Execute result as a separate spoiler (like a regular tool result) + if (executeTool?.result) { + parts.push(`
\n🔧 execute\n\n`); + parts.push(`Result:\n\`\`\`\n${executeTool.result}\n\`\`\`\n\n`); + parts.push(`
`); + } + } + else { + // No child tool calls — show tools normally (including execute as-is) + for (const tool of step.tools) { + parts.push(`
\n🔧 ${tool.name}\n\n`); + parts.push(`Args:\n\`\`\`\n${tool.args}\n\`\`\`\n\n`); + if (tool.failed) { + parts.push(`Failed: ${tool.result}\n\n`); + } + else { + parts.push(`Result:\n\`\`\`\n${tool.result}\n\`\`\`\n\n`); + } + parts.push(`
`); + } + } + } + // Separator before text + if (step.text) { + if (parts.length > 0) + parts.push("\n---\n"); + parts.push(step.text); + } + return parts.join("\n"); + } + // ============================================================================= + // Send message to Matrix (reuses existing sendReply/sendNotice) + // ============================================================================= + async function sendMessageToMatrix(context, text) { + try { + await matrix.sendReply(context, text); + } + catch (err) { + log(`Failed to send step message: ${String(err)}`); + await matrix.sendNotice(context, `Error sending: ${String(err).slice(0, 200)}`); + } + } + // Message handler + matrix.on(async (data) => { + const { context, query, sender, roomId, eventId, timestamp } = data; + // Check rate limit + if (sessionManager.isRateLimited(sender)) { + await matrix.sendNotice(context, "Rate limited. Please wait a moment."); + return; + } + // Check for active query + if (sessionManager.hasActiveQuery(roomId, context.replyThreadRootId || context.eventId)) { + await matrix.sendNotice(context, "A request is already running in this thread. Please wait."); + return; + } + const abortFn = () => { }; + const releaseQuery = sessionManager.markQueryActive(roomId, context.replyThreadRootId || context.eventId, abortFn); + const threadRootId = context.replyThreadRootId || context.eventId; + const session = sessionManager.getOrCreateSession(roomId, threadRootId, eventId); + session.isActive = true; + session.messageCount++; + session.inputChars += query.length; + session.lastEventIds.set(context.replyThreadRootId || context.eventId, eventId); + session.lastActivity = Date.now(); + let opencodeSessionId = session.opencodeSessionId; + // If we have a session ID from the map, verify it exists in OpenCode + if (opencodeSessionId) { + try { + await ctx.session.get({ sessionID: opencodeSessionId }); + log(`Using existing OpenCode session ${opencodeSessionId}`); + } + catch (getErr) { + log(`Session ${opencodeSessionId} not found in OpenCode, creating new one`); + opencodeSessionId = ""; + } + } + // If no session ID yet, create one + if (!opencodeSessionId) { + log(`Creating new OpenCode session for ${roomId}:${threadRootId}`); + try { + const createResult = await ctx.session.create({ title: `Matrix: ${roomId.slice(0, 30)}...` }); + opencodeSessionId = createResult?.id || createResult?.data?.id; + log(`Created OpenCode session ${opencodeSessionId} for ${roomId}:${threadRootId}`); + session.opencodeSessionId = opencodeSessionId; + sessionManager.saveSession(roomId, threadRootId, opencodeSessionId); + sessionContextMap.set(opencodeSessionId, { roomId, threadRootId, context }); + } + catch (createErr) { + log(`Failed to create session: ${String(createErr)}`); + session.isActive = false; + session.lastActivity = Date.now(); + releaseQuery(); + await matrix.sendNotice(context, "Failed to create session. Check logs."); + return; + } + } + try { + // Trigger prompt — returns admission receipt, not response text + await ctx.session.prompt({ + sessionID: opencodeSessionId, + text: query, + }); + log(`Prompt sent to session ${opencodeSessionId}`); + // Create session state for step tracking + sessionsBySessionID.set(opencodeSessionId, { + steps: new Map(), + currentStepID: null, + context, + finalResponse: undefined, + done: false, + }); + // Wait for response text from event subscription + // Resolved by: session.step.ended (finish: "stop") + // Rejected by: session.execution.failed / session.execution.interrupted + const responsePromise = new Promise((resolve, reject) => { + pendingResponses.set(opencodeSessionId, { resolve, reject }); + }); + const responseText = await responsePromise; + session.outputChars += responseText.length; + // Step messages already sent by event subscription — no double send + } + catch (err) { + log(`Error processing message: ${String(err)}`); + const errMsg = err?.response?.error?.message || String(err); + await matrix.sendNotice(context, `Error: ${errMsg.slice(0, 200)}`); + } + finally { + session.isActive = false; + session.lastActivity = Date.now(); + sessionContextMap.delete(opencodeSessionId); + releaseQuery(); + } + }); + // Bridge commands + async function handleBridgeCommand(context, cmdName, sendFn) { + const threadRootId = context.threadRootId || context.replyThreadRootId || context.eventId; + const key = sessionManager.getSessionKey(context.roomId, threadRootId); + const sess = sessionManager.sessions.get(key); + switch (cmdName) { + case "status": { + if (!sess) { + await sendFn("No active session."); + return; + } + const age = Math.round((Date.now() - sess.lastActivity) / 60000); + const msg = `Session status:\n- Messages: ${sess.messageCount}\n- Age: ~${age} min ago\n- Input: ~${sess.inputChars} chars\n- Output: ~${sess.outputChars} chars`; + await sendFn(msg); + break; + } + case "clear": + case "reset": { + sessionManager.removeSession(context.roomId, threadRootId); + await sendFn("Session cleared. Next message will start a fresh session."); + break; + } + case "help": + case "h": { + const trigger = options.triggerPatterns?.[0] || "!oc "; + const botName = options.botName || "opencode"; + const msg = [ + `${botName} - OpenCode Matrix Plugin`, + "", + "Bridge commands:", + `- /h or /help - Show this help`, + `- /status - Show current chat session info`, + `- /clear or /reset - Delete current session history`, + "", + `Usage: ${trigger} `, + `Or @${botName} `, + `Or reply in a thread`, + ].join("\n"); + await sendFn(msg); + break; + } + default: + await sendFn(`Unknown command: ${cmdName}. Try /help`); + } + } + // Override the matrix.on handler to also process commands + const originalListener = matrix["listeners"][0]; + matrix["listeners"] = []; + matrix.on(async (data) => { + const { context, query, sender, eventId } = data; + // Handle bridge commands + if (query.startsWith("/")) { + const cmdName = query.slice(1).split(" ")[0].toLowerCase(); + const bridgeCommands = ["status", "clear", "reset", "help", "h"]; + if (bridgeCommands.includes(cmdName)) { + const sendFn = (text) => matrix.sendNotice(context, text); + await handleBridgeCommand(context, cmdName, sendFn); + return; + } + } + // Forward to original handler + if (originalListener) { + await originalListener(data); + } + }); + // V2 Storage + if (ctx?.storage) { + await ctx.storage.set("matrix/plugin_version", "1.0.0"); + await ctx.storage.set("matrix/started_at", new Date().toISOString()); + } + // Start Matrix client + await matrix.start(); + log(`Plugin setup completed successfully`); + cleanup = () => { + log("Plugin cleaning up..."); + eventController.abort(); + matrix.stop(); + // Don't release lock — keep it on disk so client always sees it + }; + return cleanup; + } + catch (err) { + logError(`Setup error: ${String(err)}`); + releaseLock(); + return () => { }; + } +} +// ============================================================================= +// V2 Setup Function +// ============================================================================= +async function v2Setup(ctx) { + const rawStorage = ctx.options?.storagePath; + const dataDir = rawStorage + ? path.join(rawStorage, "plugin-data") + : path.join(os.homedir(), ".local", "share", "opencode-matrix-bot", "plugin-data"); + ensureDataDir(dataDir); + // Try to use OpenCode's logging API + if (ctx?.app?.log) { + setAppLog(ctx.app.log.bind(ctx.app)); + } + try { + const { options } = await loadConfig(ctx.options); + return await runPlugin(ctx, options, dataDir); + } + catch (err) { + logError(`V2 setup error: ${String(err)}`); + return () => { }; + } +} +// ============================================================================= +// Export +// ============================================================================= +export default { + id: "matrix-plugin", + async setup(ctx) { + log("Using V2 path (full context)"); + return await v2Setup(ctx); + }, +}; diff --git a/dist/logger.js b/dist/logger.js new file mode 100644 index 0000000..8d7d946 --- /dev/null +++ b/dist/logger.js @@ -0,0 +1,67 @@ +import fs from "fs"; +import os from "os"; +import path from "path"; +const LOG_DIR = path.join(os.homedir(), ".opencode-matrix-plugin"); +const LOG_FILE = path.join(LOG_DIR, "plugin.log"); +let appLog = null; +export function setAppLog(fn) { + appLog = fn; +} +export function log(message) { + // Try OpenCode's logging API first + if (appLog) { + appLog({ + service: "matrix-plugin", + level: "info", + message, + }).catch(() => { + // Fallback to file logging + try { + fs.mkdirSync(LOG_DIR, { recursive: true }); + fs.appendFileSync(LOG_FILE, `[${new Date().toISOString()}] [matrix-plugin] ${message}\n`); + } + catch { + // silent + } + }); + } + else { + // No app log available, use file logging + try { + fs.mkdirSync(LOG_DIR, { recursive: true }); + fs.appendFileSync(LOG_FILE, `[${new Date().toISOString()}] [matrix-plugin] ${message}\n`); + } + catch { + // silent + } + } +} +export function logError(message) { + // Try OpenCode's logging API first + if (appLog) { + appLog({ + service: "matrix-plugin", + level: "error", + message, + }).catch(() => { + // Fallback to file logging + try { + fs.mkdirSync(LOG_DIR, { recursive: true }); + fs.appendFileSync(LOG_FILE, `[${new Date().toISOString()}] [matrix-plugin] ERROR: ${message}\n`); + } + catch { + // silent + } + }); + } + else { + // No app log available, use file logging + try { + fs.mkdirSync(LOG_DIR, { recursive: true }); + fs.appendFileSync(LOG_FILE, `[${new Date().toISOString()}] [matrix-plugin] ERROR: ${message}\n`); + } + catch { + // silent + } + } +} diff --git a/dist/matrix-client.js b/dist/matrix-client.js new file mode 100644 index 0000000..ea02b2a --- /dev/null +++ b/dist/matrix-client.js @@ -0,0 +1,436 @@ +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); + } +} diff --git a/dist/session-manager.js b/dist/session-manager.js new file mode 100644 index 0000000..1f10467 --- /dev/null +++ b/dist/session-manager.js @@ -0,0 +1,166 @@ +import fs from "fs"; +import path from "path"; +const SESSION_MAP_FILE = path.join(process.env.HOME || "/tmp", ".opencode-matrix-sessions.json"); +console.log(`[SessionManager] SESSION_MAP_FILE=${SESSION_MAP_FILE}`); +function loadSessionMap() { + const result = new Map(); + try { + if (fs.existsSync(SESSION_MAP_FILE)) { + const data = JSON.parse(fs.readFileSync(SESSION_MAP_FILE, "utf-8")); + console.log(`[SessionManager] Loaded ${data.length} sessions from ${SESSION_MAP_FILE}`); + if (Array.isArray(data)) { + for (const entry of data) { + const key = `${entry.matrixRoomId}::${entry.threadRootId}`; + result.set(key, entry); + } + } + } + else { + console.log(`[SessionManager] File not found: ${SESSION_MAP_FILE}`); + } + } + catch (e) { + console.log(`[SessionManager] Error loading: ${e}`); + } + return result; +} +function saveSessionMap(map) { + try { + const entries = []; + for (const [, entry] of map) { + entries.push(entry); + } + fs.writeFileSync(SESSION_MAP_FILE, JSON.stringify(entries, null, 2)); + } + catch (e) { + // ignore + } +} +export class SessionManager { + sessions = new Map(); + persistentMap; + activeQueries = new Map(); + rateLimitMs; + lastRequestByUser = new Map(); + processedEvents = new Set(); + constructor(rateLimitSeconds) { + this.rateLimitMs = (rateLimitSeconds || 5) * 1000; + this.persistentMap = loadSessionMap(); + // Restore sessions from persistent map + for (const [key, entry] of this.persistentMap) { + const [matrixRoomId, threadRootId] = key.split("::"); + this.sessions.set(key, { + matrixRoomId, + threadRootId, + opencodeSessionId: entry.opencodeSessionId, + eventId: "", + isActive: false, + messageCount: entry.messageCount, + inputChars: entry.inputChars, + outputChars: entry.outputChars, + lastActivity: entry.lastActivity, + lastEventIds: new Map(), + }); + } + } + getSessionKey(matrixRoomId, threadRootId) { + return `${matrixRoomId}::${threadRootId}`; + } + getOrCreateSession(matrixRoomId, threadRootId, eventId) { + const key = this.getSessionKey(matrixRoomId, threadRootId); + let session = this.sessions.get(key); + if (!session) { + session = { + matrixRoomId, + threadRootId, + opencodeSessionId: "", + eventId, + isActive: false, + messageCount: 0, + inputChars: 0, + outputChars: 0, + lastActivity: Date.now(), + lastEventIds: new Map(), + }; + this.sessions.set(key, session); + } + return session; + } + getSession(matrixRoomId, threadRootId) { + const key = this.getSessionKey(matrixRoomId, threadRootId); + return this.sessions.get(key); + } + removeSession(matrixRoomId, threadRootId) { + const key = this.getSessionKey(matrixRoomId, threadRootId); + this.sessions.delete(key); + this.persistentMap.delete(key); + saveSessionMap(this.persistentMap); + this.activeQueries.delete(key); + } + has(matrixRoomId, threadRootId) { + const key = this.getSessionKey(matrixRoomId, threadRootId); + return this.sessions.has(key); + } + hasActiveQuery(matrixRoomId, threadRootId) { + const key = this.getSessionKey(matrixRoomId, threadRootId); + return this.activeQueries.has(key); + } + markQueryActive(matrixRoomId, threadRootId, abortFn) { + const key = this.getSessionKey(matrixRoomId, threadRootId); + this.activeQueries.set(key, abortFn); + return () => this.activeQueries.delete(key); + } + isRateLimited(sender) { + const now = Date.now(); + const last = this.lastRequestByUser.get(sender); + if (last && now - last < this.rateLimitMs) { + return true; + } + this.lastRequestByUser.set(sender, now); + return false; + } + 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; + } + saveSession(matrixRoomId, threadRootId, opencodeSessionId) { + const key = this.getSessionKey(matrixRoomId, threadRootId); + const session = this.sessions.get(key); + if (session) { + const entry = { + matrixRoomId, + threadRootId, + opencodeSessionId, + messageCount: session.messageCount, + inputChars: session.inputChars, + outputChars: session.outputChars, + lastActivity: session.lastActivity, + }; + this.persistentMap.set(key, entry); + saveSessionMap(this.persistentMap); + } + } + syncWithOpenCode(openCodeSessionIds) { + const toDelete = []; + for (const [, entry] of this.persistentMap) { + if (!openCodeSessionIds.has(entry.opencodeSessionId)) { + toDelete.push(`${entry.matrixRoomId}::${entry.threadRootId}`); + } + } + for (const key of toDelete) { + this.persistentMap.delete(key); + this.sessions.delete(key); + this.activeQueries.delete(key); + } + if (toDelete.length > 0) { + saveSessionMap(this.persistentMap); + } + } +} diff --git a/dist/types.js b/dist/types.js new file mode 100644 index 0000000..cb0ff5c --- /dev/null +++ b/dist/types.js @@ -0,0 +1 @@ +export {}; diff --git a/package.json b/package.json index dcafe74..70a15af 100644 --- a/package.json +++ b/package.json @@ -3,7 +3,7 @@ "version": "1.0.0", "type": "module", "description": "OpenCode V2 plugin for Matrix messaging integration", - "main": "index.ts", + "main": "dist/index.js", "scripts": { "build": "tsc", "typecheck": "tsc --noEmit" @@ -15,6 +15,10 @@ "marked": "^15.0.0", "matrix-bot-sdk": "^0.8.0" }, + "files": [ + "dist", + "package.json" + ], "devDependencies": { "@types/node": "^20.0.0", "typescript": "^5.7.0" diff --git a/tsconfig.json b/tsconfig.json index e7da084..33f4b62 100644 --- a/tsconfig.json +++ b/tsconfig.json @@ -13,8 +13,7 @@ "sourceMap": false, "outDir": "./dist", "rootDir": ".", - "noEmit": true, - "allowImportingTsExtensions": true + "noEmit": false }, "include": ["**/*.ts"], "exclude": ["node_modules", "dist"]