matrix-plugin/index.ts

399 lines
14 KiB
TypeScript

import fs from "fs"
import os from "os"
import path from "path"
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: string): void {
if (!fs.existsSync(dir)) {
fs.mkdirSync(dir, { recursive: true })
}
}
// =============================================================================
// Global mutex — prevents double startup (server + client both load plugins)
// =============================================================================
let LOCK_FILE = ""
function tryAcquireLock(): boolean {
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(): void {
try {
fs.unlinkSync(LOCK_FILE)
} catch {
// ignore
}
}
// =============================================================================
// Plugin Setup
// =============================================================================
async function runPlugin(ctx: any, options: any, dataDir: string) {
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: (() => void) | undefined
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<string, { roomId: string; threadRootId: string; context: any }>()
// Pending prompt responses — maps sessionID -> { resolve, reject, reasoning }
const pendingResponses = new Map<string, {
resolve: (text: string) => void
reject: (err: any) => void
reasoning: string
}>()
// V2 Event subscription — waits for response completion signals
// Event structure: { type: "...", data: { sessionID: "...", ... } }
void (async () => {
try {
for await (const event of ctx.event.subscribe({ signal: eventController.signal })) {
let eventData: any = event
if (typeof event === "string") {
try { eventData = JSON.parse(event) } catch { continue }
}
const eventType = eventData?.type || ""
const sessionId = eventData?.data?.sessionID || ""
// Session reasoning delta — streaming reasoning chunks
if (eventType === "session.reasoning.delta" && sessionId && pendingResponses.has(sessionId)) {
const delta = eventData?.data?.delta || ""
pendingResponses.get(sessionId)!.reasoning += delta
continue
}
// Session reasoning ended — full reasoning text
if (eventType === "session.reasoning.ended" && sessionId && pendingResponses.has(sessionId)) {
const text = eventData?.data?.text || ""
pendingResponses.get(sessionId)!.reasoning = text
continue
}
// Response text complete — full assembled text available
// Structure: { type: "session.text.ended", data: { sessionID, text: "full text" } }
if (eventType === "session.text.ended" && sessionId && pendingResponses.has(sessionId)) {
const fullText = eventData?.data?.text || ""
const pending = pendingResponses.get(sessionId)!
pendingResponses.delete(sessionId)
// Combine reasoning (in spoiler) + main text
let responseText = fullText
if (pending.reasoning) {
responseText = `<details>\n<summary>Thinking</summary>\n\n${pending.reasoning}\n\n</details>\n\n---\n\n${fullText}`
}
log(`Response ready: ${responseText.length} chars (${pending.reasoning ? 'with reasoning' : 'no reasoning'}) for session ${sessionId}`)
pending.resolve(responseText)
continue
}
// Execution failed — session errored
if (eventType === "session.execution.failed" && sessionId && pendingResponses.has(sessionId)) {
const pending = pendingResponses.get(sessionId)!
pendingResponses.delete(sessionId)
pending.reject(new Error("Session execution failed"))
continue
}
// Execution interrupted — session was interrupted
if (eventType === "session.execution.interrupted" && sessionId && pendingResponses.has(sessionId)) {
const pending = pendingResponses.get(sessionId)!
pendingResponses.delete(sessionId)
pending.reject(new Error("Session execution interrupted"))
continue
}
}
} catch (err) {
if (String(err).includes("AbortError") || (err as { name?: string }).name === "AbortError") {
// expected during cleanup
} else {
log(`Event subscription error: ${String(err)}`)
}
}
})()
// Message handler
matrix.on(async (data: any) => {
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 = () => { /* placeholder for cancel */ }
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: any) {
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: any) {
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}`)
// Wait for response text from event subscription
// Resolved by: session.text.ended (success)
// Rejected by: session.execution.failed / session.execution.interrupted
const responsePromise = new Promise<string>((resolve, reject) => {
pendingResponses.set(opencodeSessionId, { resolve, reject, reasoning: "" })
})
const responseText = await responsePromise
session.outputChars += responseText.length
await matrix.sendReply(context, responseText)
} catch (err: any) {
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: { roomId: string; replyThreadRootId?: string; threadRootId?: string; eventId: string },
cmdName: string,
sendFn: (text: string) => Promise<void>,
): Promise<void> {
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} <your question>`,
`Or @${botName} <your question>`,
`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 as any)["listeners"][0]
;(matrix as any)["listeners"] = []
matrix.on(async (data: any) => {
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: string) => 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: any) {
const rawStorage = (ctx.options as any)?.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: any) {
log("Using V2 path (full context)")
return await v2Setup(ctx)
},
}