From b97c82910c817e91a2ed60e0e94875d1ac85f658 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: Sat, 26 Sep 2026 09:55:12 +0300 Subject: [PATCH] feat: per-step messages, full tool output, no truncation MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Каждый step (LLM turn) — отдельное сообщение в Matrix: Step с finish:tool-calls:
Thinking reasoning...
🔧 webfetch Args: {...} Result: {full output}
Step с finish:stop: (то же самое) --- финальный текст ответа Архитектура: - SessionState per OpenCode session: steps[], currentStep, context - step.started → новый StepData - reasoning.* → в step.reasoning - tool.* → в step.tools (имя, args, result — БЕЗ ОБРЕЗОК) - text.ended → в step.text - step.ended → форматирование + отправка сообщения - step.ended(finish:stop) → resolve pending promise Спойлеры (
): - Thinking — для reasoning - 🔧 toolName — для каждого инструмента - Полный вывод инструментов внутри спойлера --- index.ts | 268 +++++++++++++++++++++++++++++++++++++++++++++++-------- 1 file changed, 231 insertions(+), 37 deletions(-) diff --git a/index.ts b/index.ts index dd2ffc5..8d5016e 100644 --- a/index.ts +++ b/index.ts @@ -1,7 +1,7 @@ import fs from "fs" import os from "os" import path from "path" -import { MatrixBotClient } from "./matrix-client.js" +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" @@ -93,15 +93,47 @@ async function runPlugin(ctx: any, options: any, dataDir: string) { // Track Matrix context per OpenCode session const sessionContextMap = new Map() - // Pending prompt responses — maps sessionID -> { resolve, reject, reasoning } + // ============================================================================= + // Pending prompt responses — maps sessionID -> { resolve, reject } + // ============================================================================= + const pendingResponses = new Map void reject: (err: any) => void - reasoning: string }>() - // V2 Event subscription — waits for response completion signals - // Event structure: { type: "...", data: { sessionID: "...", ... } } + // ============================================================================= + // 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 // assistantMessageID → step + currentStepID: string | null + context: MatrixEventContext + finalResponse?: string + } + + const sessionsBySessionID = new Map() + + // ============================================================================= + // Event subscription — collects step data and sends messages per step + // ============================================================================= + void (async () => { try { for await (const event of ctx.event.subscribe({ signal: eventController.signal })) { @@ -112,52 +144,155 @@ async function runPlugin(ctx: any, options: any, dataDir: string) { const eventType = eventData?.type || "" const sessionId = eventData?.data?.sessionID || "" + const assistantMsgID = eventData?.data?.assistantMessageID || "" - // 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 + // ── 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 } - // 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 + // ── 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 } - // 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) + 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 + } - // Combine reasoning (in spoiler) + main text - let responseText = fullText - if (pending.reasoning) { - responseText = `
\nThinking\n\n${pending.reasoning}\n\n
\n\n---\n\n${fullText}` + // ── 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 + } + 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) + 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) + 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, resolve the pending promise + if (finish === "stop" && state) { + const pending = pendingResponses.get(sessionId) + if (pending && state.finalResponse) { + pendingResponses.delete(sessionId) + pending.resolve(state.finalResponse) + } } - 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")) + // ── 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 } - // 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")) + 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 } } @@ -170,6 +305,57 @@ async function runPlugin(ctx: any, options: any, dataDir: string) { } })() + // ============================================================================= + // Format step into Matrix message + // ============================================================================= + + function formatStepMessage(step: StepData): string { + const parts: string[] = [] + + // Reasoning in spoiler + if (step.reasoning) { + parts.push(`
\nThinking\n\n${step.reasoning}\n\n
`) + } + + // Tool calls + if (step.tools.length > 0) { + 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: MatrixEventContext, + text: string, + ): Promise { + 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 @@ -238,11 +424,19 @@ async function runPlugin(ctx: any, options: any, dataDir: string) { }) log(`Prompt sent to session ${opencodeSessionId}`) + // Create session state for step tracking + sessionsBySessionID.set(opencodeSessionId, { + steps: new Map(), + currentStepID: null, + context, + finalResponse: undefined, + }) + // Wait for response text from event subscription - // Resolved by: session.text.ended (success) + // 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, reasoning: "" }) + pendingResponses.set(opencodeSessionId, { resolve, reject }) }) const responseText = await responsePromise