729 lines
26 KiB
TypeScript
729 lines
26 KiB
TypeScript
import fs from "fs"
|
|
import os from "os"
|
|
import path from "path"
|
|
import { MatrixBotClient, type MatrixEventContext } 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 }
|
|
// =============================================================================
|
|
|
|
const pendingResponses = new Map<string, {
|
|
resolve: (text: string) => void
|
|
reject: (err: any) => void
|
|
}>()
|
|
|
|
// =============================================================================
|
|
// Step tracking — each step is one LLM turn, sent as a separate Matrix message
|
|
// =============================================================================
|
|
|
|
interface ToolCall {
|
|
toolID: string
|
|
name: string
|
|
args: string
|
|
result: string
|
|
failed: boolean
|
|
}
|
|
|
|
interface StepData {
|
|
assistantMessageID: string
|
|
reasoning: string
|
|
tools: ToolCall[]
|
|
text: string
|
|
}
|
|
|
|
interface SessionState {
|
|
steps: Map<string, StepData> // assistantMessageID → step
|
|
currentStepID: string | null
|
|
context: MatrixEventContext
|
|
finalResponse?: string
|
|
done: boolean // true после finish:stop
|
|
}
|
|
|
|
const sessionsBySessionID = new Map<string, SessionState>()
|
|
|
|
// =============================================================================
|
|
// Parse tool calls from CodeMode JS script
|
|
// CodeMode creates nested tools: tools.<server>.<tool>(...)
|
|
// e.g. tools.duckduckgo.search({ query: "..." })
|
|
// =============================================================================
|
|
|
|
interface ChildToolInfo {
|
|
name: string
|
|
args: string
|
|
}
|
|
|
|
function parseCodeModeTools(code: string): ChildToolInfo[] {
|
|
const results: ChildToolInfo[] = []
|
|
const seen = new Set<string>()
|
|
const regex = /tools\.(\w+)\.(\w+)\s*\(/g
|
|
let match
|
|
|
|
while ((match = regex.exec(code)) !== null) {
|
|
const toolName = `${match[1]}_${match[2]}`
|
|
if (seen.has(toolName)) continue
|
|
seen.add(toolName)
|
|
|
|
// Extract the full argument block starting from the (
|
|
const openParenPos = match.index + match[0].indexOf("(")
|
|
const args = extractArgumentBlock(code, openParenPos)
|
|
results.push({ name: toolName, args })
|
|
}
|
|
|
|
return results
|
|
}
|
|
|
|
/**
|
|
* Extract argument block from code starting at the opening parenthesis.
|
|
* Handles nested parentheses, brackets, and braces.
|
|
* e.g. for tools.duckduckgo.fetch_content({ url: url, max_length: 5000 })
|
|
* returns "{ url: url, max_length: 5000 }"
|
|
*/
|
|
function extractArgumentBlock(code: string, openParenPos: number): string {
|
|
let depth = 0
|
|
let start = openParenPos
|
|
let end = openParenPos
|
|
|
|
// Find the opening character: (, {, [, or "
|
|
while (start < code.length) {
|
|
const ch = code[start]
|
|
if (ch === "(" || ch === "{" || ch === "[" || ch === '"') break
|
|
start++
|
|
}
|
|
|
|
if (start >= code.length) return "(variable)"
|
|
|
|
const openCh = code[start]
|
|
const closeCh = openCh === "(" ? ")" : openCh === "{" ? "}" : openCh === "[" ? "]" : '"'
|
|
|
|
for (let i = start; i < code.length; i++) {
|
|
const ch = code[i]
|
|
if (ch === openCh) depth++
|
|
else if (ch === closeCh) {
|
|
depth--
|
|
if (depth === 0) {
|
|
end = i + 1
|
|
break
|
|
}
|
|
}
|
|
// Skip string literals to avoid counting quotes inside strings
|
|
if (ch === '"' || ch === "'") {
|
|
let j = i + 1
|
|
while (j < code.length && code[j] !== ch) {
|
|
if (code[j] === "\\") j++
|
|
j++
|
|
}
|
|
i = j
|
|
}
|
|
}
|
|
|
|
const result = code.slice(start, end).trim()
|
|
return result || "(variable)"
|
|
}
|
|
|
|
// =============================================================================
|
|
// 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: any = 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: any
|
|
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 as { name?: string }).name === "AbortError") {
|
|
// expected during cleanup
|
|
} else {
|
|
log(`Event subscription error: ${String(err)}`)
|
|
}
|
|
}
|
|
})()
|
|
|
|
// =============================================================================
|
|
// Format step into Matrix message
|
|
// =============================================================================
|
|
|
|
function formatStepMessage(step: StepData): string {
|
|
const parts: string[] = []
|
|
|
|
// Reasoning in spoiler
|
|
if (step.reasoning) {
|
|
parts.push(`<details>\n<summary>Thinking</summary>\n\n${step.reasoning}\n\n</details>`)
|
|
}
|
|
|
|
// 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 only the children, skip execute
|
|
for (const child of childTools) {
|
|
parts.push(`<details>\n<summary>🔧 ${child.name}</summary>\n\n`)
|
|
parts.push(`<b>Args:</b>\n\`\`\`\n${child.args}\n\`\`\`\n\n`)
|
|
if (child.failed) {
|
|
parts.push(`<b>Failed:</b> ${child.result}\n\n`)
|
|
} else if (child.result) {
|
|
parts.push(`<b>Result:</b>\n\`\`\`\n${child.result}\n\`\`\`\n\n`)
|
|
}
|
|
parts.push(`</details>`)
|
|
}
|
|
} else {
|
|
// No child tool calls — show tools normally (including execute as-is)
|
|
for (const tool of step.tools) {
|
|
parts.push(`<details>\n<summary>🔧 ${tool.name}</summary>\n\n`)
|
|
parts.push(`<b>Args:</b>\n\`\`\`\n${tool.args}\n\`\`\`\n\n`)
|
|
if (tool.failed) {
|
|
parts.push(`<b>Failed:</b> ${tool.result}\n\n`)
|
|
} else {
|
|
parts.push(`<b>Result:</b>\n\`\`\`\n${tool.result}\n\`\`\`\n\n`)
|
|
}
|
|
parts.push(`</details>`)
|
|
}
|
|
}
|
|
}
|
|
|
|
// 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: MatrixEventContext,
|
|
text: string,
|
|
): Promise<void> {
|
|
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: 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}`)
|
|
|
|
// 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<string>((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: 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)
|
|
},
|
|
}
|