fix: wait for session.text.ended, no timeouts

Полный flow сессии из логов:
  prompt → inbox.enqueued → execution.started → step.started
    → reasoning → tool-calls → step.ended (finish: tool-calls)
    → step.started → reasoning → text.started → text.delta*42
    → text.ended ← ОТВЕТ ГОТОВ
    → step.ended (finish: stop) → execution.succeeded

Ключевые сигналы:
  session.text.ended  → ответ готов (полный текст в data.text)
  session.execution.failed → ошибка
  session.execution.interrupted → прервано
  session.execution.succeeded → 'запрос принят' (НЕ ответ готов!)

Изменения:
  - Убраны timeout 60s (бесконечное ожидание)
  - Resolved на session.text.ended (полный assembled text)
  - Rejected на session.execution.failed/interrupted
  - Убраны reasoning.* сборы (не для пользователя)
  - Убраны text.delta ранние resolve (race condition)
This commit is contained in:
Бородин Роман 2026-09-26 09:22:32 +03:00
parent 76db8b5d36
commit 280eaa1e9c
1 changed files with 29 additions and 65 deletions

View File

@ -93,80 +93,52 @@ async function runPlugin(ctx: any, options: any, dataDir: string) {
// Track Matrix context per OpenCode session
const sessionContextMap = new Map<string, { roomId: string; threadRootId: string; context: any }>()
// Pending prompt responses — maps sessionID -> { resolve, reject, text }
// Pending prompt responses — maps sessionID -> { resolve, reject }
const pendingResponses = new Map<string, {
resolve: (text: string) => void
reject: (err: any) => void
text: string
}>()
// V2 Event subscription — collects assistant response text from events
// 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 })) {
// Parse event — V2EventEncoded is a JSON string
let eventData: any = event
if (typeof event === "string") {
try {
eventData = JSON.parse(event)
} catch {
continue
}
try { eventData = JSON.parse(event) } catch { continue }
}
const eventType = eventData?.type || ""
// Session text delta — streaming text chunks
// Structure: { type: "session.text.delta", data: { sessionID, delta: "text" } }
if (eventType === "session.text.delta") {
const sessionId = eventData?.data?.sessionID || ""
const delta = eventData?.data?.delta || ""
if (sessionId && delta && pendingResponses.has(sessionId)) {
const pending = pendingResponses.get(sessionId)!
pending.text += delta
log(`Collected delta: ${delta.length} chars, total: ${pending.text.length} chars for session ${sessionId}`)
// Resolve immediately when we get any text
const timer = (pending as any)._safetyTimer
if (timer) clearTimeout(timer)
pendingResponses.delete(sessionId)
pending.resolve(pending.text)
}
}
// Session text ended — complete text available
// Response text complete — full assembled text available
// Structure: { type: "session.text.ended", data: { sessionID, text: "full text" } }
if (eventType === "session.text.ended") {
const sessionId = eventData?.data?.sessionID || ""
if (eventType === "session.text.ended" && sessionId && pendingResponses.has(sessionId)) {
const fullText = eventData?.data?.text || ""
if (sessionId && fullText && pendingResponses.has(sessionId)) {
const pending = pendingResponses.get(sessionId)!
pending.text = fullText
log(`Collected ended text: ${fullText.length} chars for session ${sessionId}`)
const timer = (pending as any)._safetyTimer
if (timer) clearTimeout(timer)
log(`Response ready: ${fullText.length} chars for session ${sessionId}`)
pendingResponses.delete(sessionId)
pending.resolve(pending.text)
}
pending.resolve(fullText)
continue
}
// Session reasoning delta — also collect (optional, for full response)
if (eventType === "session.reasoning.delta") {
const sessionId = eventData?.data?.sessionID || ""
const delta = eventData?.data?.delta || ""
if (sessionId && delta && pendingResponses.has(sessionId)) {
// Execution failed — session errored
if (eventType === "session.execution.failed" && sessionId && pendingResponses.has(sessionId)) {
const pending = pendingResponses.get(sessionId)!
pending.text += delta
}
log(`Session ${sessionId} execution failed`)
pendingResponses.delete(sessionId)
pending.reject(new Error("Session execution failed"))
continue
}
// Session reasoning ended — full reasoning text
if (eventType === "session.reasoning.ended") {
const sessionId = eventData?.data?.sessionID || ""
const text = eventData?.data?.text || ""
if (sessionId && text && pendingResponses.has(sessionId)) {
// Execution interrupted — session was interrupted
if (eventType === "session.execution.interrupted" && sessionId && pendingResponses.has(sessionId)) {
const pending = pendingResponses.get(sessionId)!
pending.text += text
}
log(`Session ${sessionId} execution interrupted`)
pendingResponses.delete(sessionId)
pending.reject(new Error("Session execution interrupted"))
continue
}
}
} catch (err) {
@ -247,18 +219,10 @@ async function runPlugin(ctx: any, options: any, dataDir: string) {
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) => {
const safetyTimer = setTimeout(() => {
pendingResponses.delete(opencodeSessionId)
reject(new Error("Timeout waiting for response event (60s)"))
}, 60_000)
pendingResponses.set(opencodeSessionId, {
resolve,
reject,
text: "",
})
// Store timer reference for cleanup in event subscription
;(pendingResponses.get(opencodeSessionId) as any)._safetyTimer = safetyTimer
pendingResponses.set(opencodeSessionId, { resolve, reject })
})
const responseText = await responsePromise