feat: per-step messages, full tool output, no truncation
Каждый step (LLM turn) — отдельное сообщение в Matrix: Step с finish:tool-calls: <details> <summary>Thinking</summary> reasoning... </details> <details> <summary>🔧 webfetch</summary> Args: {...} Result: {full output} </details> 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 Спойлеры (<details>): - Thinking — для reasoning - 🔧 toolName — для каждого инструмента - Полный вывод инструментов внутри спойлера
This commit is contained in:
parent
ca4863feb8
commit
b97c82910c
268
index.ts
268
index.ts
|
|
@ -1,7 +1,7 @@
|
||||||
import fs from "fs"
|
import fs from "fs"
|
||||||
import os from "os"
|
import os from "os"
|
||||||
import path from "path"
|
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 { SessionManager } from "./session-manager.js"
|
||||||
import { loadConfig } from "./config-loader.js"
|
import { loadConfig } from "./config-loader.js"
|
||||||
import { setAppLog, log, logError } from "./logger.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
|
// Track Matrix context per OpenCode session
|
||||||
const sessionContextMap = new Map<string, { roomId: string; threadRootId: string; context: any }>()
|
const sessionContextMap = new Map<string, { roomId: string; threadRootId: string; context: any }>()
|
||||||
|
|
||||||
// Pending prompt responses — maps sessionID -> { resolve, reject, reasoning }
|
// =============================================================================
|
||||||
|
// Pending prompt responses — maps sessionID -> { resolve, reject }
|
||||||
|
// =============================================================================
|
||||||
|
|
||||||
const pendingResponses = new Map<string, {
|
const pendingResponses = new Map<string, {
|
||||||
resolve: (text: string) => void
|
resolve: (text: string) => void
|
||||||
reject: (err: any) => 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<string, StepData> // assistantMessageID → step
|
||||||
|
currentStepID: string | null
|
||||||
|
context: MatrixEventContext
|
||||||
|
finalResponse?: string
|
||||||
|
}
|
||||||
|
|
||||||
|
const sessionsBySessionID = new Map<string, SessionState>()
|
||||||
|
|
||||||
|
// =============================================================================
|
||||||
|
// Event subscription — collects step data and sends messages per step
|
||||||
|
// =============================================================================
|
||||||
|
|
||||||
void (async () => {
|
void (async () => {
|
||||||
try {
|
try {
|
||||||
for await (const event of ctx.event.subscribe({ signal: eventController.signal })) {
|
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 eventType = eventData?.type || ""
|
||||||
const sessionId = eventData?.data?.sessionID || ""
|
const sessionId = eventData?.data?.sessionID || ""
|
||||||
|
const assistantMsgID = eventData?.data?.assistantMessageID || ""
|
||||||
|
|
||||||
// Session reasoning delta — streaming reasoning chunks
|
// ── Step lifecycle ──────────────────────────────────────────────
|
||||||
if (eventType === "session.reasoning.delta" && sessionId && pendingResponses.has(sessionId)) {
|
|
||||||
const delta = eventData?.data?.delta || ""
|
if (eventType === "session.step.started") {
|
||||||
pendingResponses.get(sessionId)!.reasoning += delta
|
// 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
|
continue
|
||||||
}
|
}
|
||||||
|
|
||||||
// Session reasoning ended — full reasoning text
|
// ── Reasoning ───────────────────────────────────────────────────
|
||||||
if (eventType === "session.reasoning.ended" && sessionId && pendingResponses.has(sessionId)) {
|
|
||||||
const text = eventData?.data?.text || ""
|
if (eventType === "session.reasoning.delta" && assistantMsgID) {
|
||||||
pendingResponses.get(sessionId)!.reasoning = text
|
const state = sessionsBySessionID.get(sessionId)
|
||||||
|
const step = state?.steps.get(assistantMsgID)
|
||||||
|
if (step) step.reasoning += eventData?.data?.delta || ""
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
|
|
||||||
// Response text complete — full assembled text available
|
if (eventType === "session.reasoning.ended" && assistantMsgID) {
|
||||||
// Structure: { type: "session.text.ended", data: { sessionID, text: "full text" } }
|
const state = sessionsBySessionID.get(sessionId)
|
||||||
if (eventType === "session.text.ended" && sessionId && pendingResponses.has(sessionId)) {
|
const step = state?.steps.get(assistantMsgID)
|
||||||
const fullText = eventData?.data?.text || ""
|
if (step) step.reasoning = eventData?.data?.text || ""
|
||||||
const pending = pendingResponses.get(sessionId)!
|
continue
|
||||||
pendingResponses.delete(sessionId)
|
}
|
||||||
|
|
||||||
// Combine reasoning (in spoiler) + main text
|
// ── Tool calls ──────────────────────────────────────────────────
|
||||||
let responseText = fullText
|
|
||||||
if (pending.reasoning) {
|
if (eventType === "session.tool.input.started" && assistantMsgID) {
|
||||||
responseText = `<details>\n<summary>Thinking</summary>\n\n${pending.reasoning}\n\n</details>\n\n---\n\n${fullText}`
|
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
|
continue
|
||||||
}
|
}
|
||||||
|
|
||||||
// Execution failed — session errored
|
// ── Execution errors ───────────────────────────────────────────
|
||||||
if (eventType === "session.execution.failed" && sessionId && pendingResponses.has(sessionId)) {
|
|
||||||
const pending = pendingResponses.get(sessionId)!
|
if (eventType === "session.execution.failed" && sessionId) {
|
||||||
pendingResponses.delete(sessionId)
|
const pending = pendingResponses.get(sessionId)
|
||||||
pending.reject(new Error("Session execution failed"))
|
if (pending) {
|
||||||
|
pendingResponses.delete(sessionId)
|
||||||
|
pending.reject(new Error("Session execution failed"))
|
||||||
|
}
|
||||||
|
sessionsBySessionID.delete(sessionId)
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
|
|
||||||
// Execution interrupted — session was interrupted
|
if (eventType === "session.execution.interrupted" && sessionId) {
|
||||||
if (eventType === "session.execution.interrupted" && sessionId && pendingResponses.has(sessionId)) {
|
const pending = pendingResponses.get(sessionId)
|
||||||
const pending = pendingResponses.get(sessionId)!
|
if (pending) {
|
||||||
pendingResponses.delete(sessionId)
|
pendingResponses.delete(sessionId)
|
||||||
pending.reject(new Error("Session execution interrupted"))
|
pending.reject(new Error("Session execution interrupted"))
|
||||||
|
}
|
||||||
|
sessionsBySessionID.delete(sessionId)
|
||||||
continue
|
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(`<details>\n<summary>Thinking</summary>\n\n${step.reasoning}\n\n</details>`)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Tool calls
|
||||||
|
if (step.tools.length > 0) {
|
||||||
|
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
|
// Message handler
|
||||||
matrix.on(async (data: any) => {
|
matrix.on(async (data: any) => {
|
||||||
const { context, query, sender, roomId, eventId, timestamp } = data
|
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}`)
|
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
|
// 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
|
// Rejected by: session.execution.failed / session.execution.interrupted
|
||||||
const responsePromise = new Promise<string>((resolve, reject) => {
|
const responsePromise = new Promise<string>((resolve, reject) => {
|
||||||
pendingResponses.set(opencodeSessionId, { resolve, reject, reasoning: "" })
|
pendingResponses.set(opencodeSessionId, { resolve, reject })
|
||||||
})
|
})
|
||||||
|
|
||||||
const responseText = await responsePromise
|
const responseText = await responsePromise
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue