476 lines
18 KiB
TypeScript
476 lines
18 KiB
TypeScript
import fs from "fs"
|
|
import { MatrixBotClient } from "./matrix-client.js"
|
|
import { SessionManager, encodeSessionId } from "./session-manager.js"
|
|
import { loadConfig } from "./config-loader.js"
|
|
import { log } from "./logger.js"
|
|
|
|
log("Plugin module loaded")
|
|
|
|
// =============================================================================
|
|
// V2 Support Detection
|
|
// =============================================================================
|
|
|
|
function isV2Context(ctx: any): boolean {
|
|
return !!(ctx?.session?.prompt && ctx?.event?.subscribe && ctx?.storage)
|
|
}
|
|
|
|
// =============================================================================
|
|
// Shared initialization logic
|
|
// =============================================================================
|
|
|
|
async function runPlugin(ctx: any, options: any, sdkClient?: any) {
|
|
const isV2 = isV2Context(ctx)
|
|
const eventController = new AbortController()
|
|
let cleanup: (() => void) | undefined
|
|
|
|
try {
|
|
const directory = ctx?.location?.directory || process.cwd()
|
|
fs.appendFileSync("/tmp/opencode-matrix-plugin/setup-started.log", `[${new Date().toISOString()}] setup() called, dir=${directory}, v2=${isV2}, sdk=${!!sdkClient}\n`)
|
|
log(`Plugin setup started, directory=${directory}, v2=${isV2}, sdk=${!!sdkClient}`)
|
|
|
|
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 pending tool calls per session: sessionID -> Map<callID, {roomId, threadRootId, context}>
|
|
const pendingToolCalls = new Map<string, Map<string, { roomId: string; threadRootId: string; context: any }>>()
|
|
|
|
// Track sent messages per session to avoid duplicates
|
|
const sentMessageIds = new Map<string, Set<string>>()
|
|
|
|
// Track Matrix context per OpenCode session
|
|
const sessionContextMap = new Map<string, { roomId: string; threadRootId: string; context: any }>()
|
|
|
|
// Helper: send a message to Matrix
|
|
async function sendMatrixMessage(roomId: string, threadRootId: string, context: any, text: string): Promise<void> {
|
|
if (!text.trim()) return
|
|
await matrix.sendReply(context, text)
|
|
log(`Matrix: Sent to ${roomId}:${threadRootId} (${text.length} chars)`)
|
|
}
|
|
|
|
// Helper: extract content from parts
|
|
function extractContent(parts: any[]): { reasoningText: string; mainText: string; toolCallText: string } {
|
|
let reasoningText = ""
|
|
let mainText = ""
|
|
let toolCallText = ""
|
|
const reasoningTypes = new Set(["reasoning", "thinking", "redwood", "reason"])
|
|
|
|
for (const p of parts) {
|
|
const content = p?.content || p?.text || ""
|
|
const pType = p?.type || ""
|
|
if (reasoningTypes.has(pType)) {
|
|
reasoningText += content
|
|
} else if (pType === "tool_call" || pType === "tool") {
|
|
toolCallText += content
|
|
} else {
|
|
mainText += content
|
|
}
|
|
}
|
|
|
|
return { reasoningText, mainText, toolCallText }
|
|
}
|
|
|
|
// Helper: format message text
|
|
function formatMessage(toolCallText: string, reasoningText: string, mainText: string): string {
|
|
if (toolCallText.trim()) {
|
|
return `Tool call: ${toolCallText.trim()}`
|
|
}
|
|
|
|
let responseText = mainText
|
|
if (reasoningText.trim()) {
|
|
responseText = `<details><summary>Thoughts</summary>\n\n${reasoningText.trim()}\n\n</details>\n\n---\n\n${mainText}`
|
|
}
|
|
|
|
return responseText
|
|
}
|
|
|
|
// V2 Event subscription for response streaming
|
|
if (isV2) {
|
|
void (async () => {
|
|
try {
|
|
for await (const event of ctx.event.subscribe({ signal: eventController.signal })) {
|
|
const eventType = typeof event === "object" && event !== null && "type" in event ? (event as { type: string }).type : ""
|
|
if (eventType === "session.updated") {
|
|
const sessionId = (event as { sessionID?: string }).sessionID
|
|
if (!sessionId) continue
|
|
const update = (event as { update?: { type: string; content?: string; toolName?: string } }).update
|
|
if (!update) continue
|
|
log(`V2 event: session=${sessionId}, update.type=${update.type}`)
|
|
}
|
|
}
|
|
} 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 no session ID yet, create one
|
|
if (!opencodeSessionId) {
|
|
log(`Creating new OpenCode session for ${roomId}:${threadRootId}`)
|
|
try {
|
|
if (isV2 && ctx?.session?.create) {
|
|
// V2: can pass custom ID
|
|
const encodedId = encodeSessionId(roomId, threadRootId)
|
|
const createResult = await ctx.session.create({ id: encodedId, title: `Matrix: ${roomId.slice(0, 30)}...` })
|
|
opencodeSessionId = createResult?.id || createResult?.data?.id || encodedId
|
|
log(`V2: Created OpenCode session ${opencodeSessionId} (custom ID: ${encodedId})`)
|
|
} else if (sdkClient) {
|
|
// V1: cannot pass custom ID, need to save mapping
|
|
const createResult = await sdkClient.session.create({ body: { title: `Matrix: ${roomId.slice(0, 30)}...` } })
|
|
opencodeSessionId = createResult?.data?.id || createResult?.id || ""
|
|
log(`V1: Created OpenCode session ${opencodeSessionId}`)
|
|
}
|
|
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 {
|
|
if (isV2) {
|
|
// V2: use ctx.session.prompt()
|
|
const result = await ctx.session.prompt({
|
|
sessionID: opencodeSessionId,
|
|
text: query,
|
|
})
|
|
const responseText = result ? String(result) : "No response"
|
|
session.outputChars += responseText.length
|
|
await matrix.sendReply(context, responseText)
|
|
} else if (sdkClient) {
|
|
// V1: use SDK client.session.prompt()
|
|
log(`V1: Sending prompt to session ${opencodeSessionId}: ${query}`)
|
|
|
|
// Start polling for new messages
|
|
const pollInterval = setInterval(async () => {
|
|
if (!session.isActive) {
|
|
clearInterval(pollInterval)
|
|
return
|
|
}
|
|
|
|
try {
|
|
const existingIds = sentMessageIds.get(opencodeSessionId) || new Set()
|
|
const messagesResult = await sdkClient.session.messages({
|
|
path: { id: opencodeSessionId },
|
|
})
|
|
|
|
const messages = messagesResult?.data || []
|
|
if (Array.isArray(messages)) {
|
|
for (const msg of messages) {
|
|
const msgId = msg?.info?.id || msg?.id
|
|
if (!msgId || existingIds.has(msgId)) continue
|
|
|
|
const role = msg?.info?.role || ""
|
|
if (role !== "assistant" && role !== "tool_call") continue
|
|
|
|
const parts = msg?.parts || msg?.info?.parts || []
|
|
const { reasoningText, mainText, toolCallText } = extractContent(parts)
|
|
const responseText = formatMessage(toolCallText, reasoningText, mainText)
|
|
|
|
if (responseText.trim()) {
|
|
await matrix.sendReply(context, responseText)
|
|
log(`V1: Sent message ${msgId} (role=${role}, ${responseText.length} chars)`)
|
|
existingIds.add(msgId)
|
|
}
|
|
}
|
|
}
|
|
|
|
sentMessageIds.set(opencodeSessionId, existingIds)
|
|
} catch (err) {
|
|
log(`V1: Error during polling: ${String(err)}`)
|
|
}
|
|
}, 1000)
|
|
|
|
try {
|
|
const result = await sdkClient.session.prompt({
|
|
path: { id: opencodeSessionId },
|
|
body: { parts: [{ type: "text", text: query }] },
|
|
})
|
|
|
|
// Stop polling
|
|
clearInterval(pollInterval)
|
|
|
|
// Send any remaining messages
|
|
try {
|
|
const existingIds = sentMessageIds.get(opencodeSessionId) || new Set()
|
|
const messagesResult = await sdkClient.session.messages({
|
|
path: { id: opencodeSessionId },
|
|
})
|
|
|
|
const messages = messagesResult?.data || []
|
|
if (Array.isArray(messages)) {
|
|
for (const msg of messages) {
|
|
const msgId = msg?.info?.id || msg?.id
|
|
if (!msgId || existingIds.has(msgId)) continue
|
|
|
|
const role = msg?.info?.role || ""
|
|
if (role !== "assistant" && role !== "tool_call") continue
|
|
|
|
const parts = msg?.parts || msg?.info?.parts || []
|
|
const { reasoningText, mainText, toolCallText } = extractContent(parts)
|
|
const responseText = formatMessage(toolCallText, reasoningText, mainText)
|
|
|
|
if (responseText.trim()) {
|
|
await matrix.sendReply(context, responseText)
|
|
log(`V1: Sent message ${msgId} (role=${role}, ${responseText.length} chars)`)
|
|
existingIds.add(msgId)
|
|
}
|
|
}
|
|
}
|
|
|
|
sentMessageIds.set(opencodeSessionId, existingIds)
|
|
} catch (err) {
|
|
log(`V1: Error sending final messages: ${String(err)}`)
|
|
}
|
|
} catch (promptErr: any) {
|
|
log(`V1 prompt error: ${String(promptErr)}`)
|
|
const errMsg = promptErr?.response?.error?.message || promptErr?.error?.message || String(promptErr)
|
|
await matrix.sendNotice(context, `Error: ${errMsg.slice(0, 200)}`)
|
|
}
|
|
} else {
|
|
// No SDK client available
|
|
log(`V1 mode without SDK client: Would send prompt to session ${opencodeSessionId}: ${query}`)
|
|
const responseText = "Matrix plugin running in V1 mode — prompt forwarding requires OpenCode SDK client"
|
|
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 mode = isV2 ? "OpenCode V2 (full)" : sdkClient ? "OpenCode V1 (full via SDK)" : "OpenCode V1 (limited — requires SDK client)"
|
|
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`,
|
|
"",
|
|
`Mode: ${mode}`,
|
|
].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 - defensive access
|
|
if (isV2 && 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 (V2=${isV2}, sdk=${!!sdkClient})`)
|
|
|
|
cleanup = () => {
|
|
log("Plugin cleaning up...")
|
|
eventController.abort()
|
|
matrix.stop()
|
|
}
|
|
|
|
return cleanup
|
|
} catch (err) {
|
|
fs.appendFileSync("/tmp/opencode-matrix-plugin/setup-error.log", `[${new Date().toISOString()}] ${String(err)}\n`)
|
|
log(`Setup error: ${String(err)}`)
|
|
return () => {
|
|
cleanup?.()
|
|
}
|
|
}
|
|
}
|
|
|
|
// =============================================================================
|
|
// V1 Server Entry Point — receives PluginInput with SDK client
|
|
// =============================================================================
|
|
|
|
async function v1Server(input: any, options?: any) {
|
|
try {
|
|
log("V1 server entry point called")
|
|
|
|
const { client, directory, worktree, project } = input
|
|
const sdkClient = client
|
|
|
|
if (!sdkClient) {
|
|
log("V1: No SDK client in PluginInput")
|
|
return {}
|
|
}
|
|
|
|
log(`V1: SDK client available from PluginInput`)
|
|
|
|
// Load config using loadConfig (reads from matrix.json or env vars)
|
|
const pluginOptions = options || {}
|
|
const { options: configOptions } = await loadConfig(pluginOptions)
|
|
|
|
// Create a context that runPlugin can use
|
|
const v1Ctx = {
|
|
location: { directory: directory || process.cwd(), project: project || { id: "v1" } },
|
|
storage: {
|
|
set: async (_key: string, _value: any) => { /* V1 storage not available */ },
|
|
get: async (_key: string) => null,
|
|
},
|
|
session: {}, // Not V2 context
|
|
event: {}, // Not V2 context
|
|
}
|
|
|
|
await runPlugin(v1Ctx as any, configOptions, sdkClient)
|
|
return {}
|
|
} catch (err) {
|
|
log(`V1 server error: ${String(err)}`)
|
|
return {}
|
|
}
|
|
}
|
|
|
|
// =============================================================================
|
|
// V2 Setup Function
|
|
// =============================================================================
|
|
|
|
async function v2Setup(ctx: any) {
|
|
try {
|
|
const { options } = await loadConfig(ctx.options)
|
|
return await runPlugin(ctx, options)
|
|
} catch (err) {
|
|
fs.appendFileSync("/tmp/opencode-matrix-plugin/setup-error.log", `[${new Date().toISOString()}] ${String(err)}\n`)
|
|
log(`V2 setup error: ${String(err)}`)
|
|
return () => {}
|
|
}
|
|
}
|
|
|
|
// =============================================================================
|
|
// Dual V1/V2 Export
|
|
// =============================================================================
|
|
|
|
// OpenCode 1.18.x calls server(PluginInput, options) where PluginInput contains SDK client.
|
|
// V2 reads id + setup() and ignores server().
|
|
export default {
|
|
id: "opencode-matrix",
|
|
async setup(ctx: any) {
|
|
const isV2 = isV2Context(ctx)
|
|
if (isV2) {
|
|
log("Using V2 path (full context)")
|
|
return await v2Setup(ctx)
|
|
}
|
|
log("Using V1 path (server() will be called by OpenCode)")
|
|
return await v1Server({}, {})
|
|
},
|
|
async server(input: any, options?: any) {
|
|
log("V1 server() called by OpenCode")
|
|
return await v1Server(input, options)
|
|
},
|
|
}
|