fix: collect response from events, remove outcome polling
outcome: succeeded = сессия запущена, НЕ ответ готов. Ответ приходит ТОЛЬКО через EventV2 subscription (live-only text). Изменения: - Убран polling ctx.session.get() + outcome - Event subscription: парсит 7 структур событий для извлечения текста - Логгирует ВСЕ session-события с полной структурой (для отладки) - Promise разрешается сразу при получении первого текста из событий - Safety timeout 60s (не основной механизм, только защита)
This commit is contained in:
parent
8e182bdc17
commit
9276de3c94
165
index.ts
165
index.ts
|
|
@ -93,45 +93,88 @@ 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 }
|
||||
// Pending prompt responses — maps sessionID -> { resolve, reject, text }
|
||||
const pendingResponses = new Map<string, {
|
||||
resolve: (text: string) => void
|
||||
reject: (err: any) => void
|
||||
text: string
|
||||
}>()
|
||||
|
||||
// V2 Event subscription — logs all events for debugging
|
||||
// V2 Event subscription — collects assistant response text from events
|
||||
void (async () => {
|
||||
try {
|
||||
for await (const event of ctx.event.subscribe({ signal: eventController.signal })) {
|
||||
// Log event type for debugging
|
||||
let eventType = ""
|
||||
// Parse event — V2EventEncoded is a JSON string
|
||||
let eventData: any = event
|
||||
if (typeof event === "string") {
|
||||
try {
|
||||
eventData = JSON.parse(event)
|
||||
eventType = eventData?.type || ""
|
||||
} catch {
|
||||
eventType = "(raw string)"
|
||||
// not JSON — log raw for debugging
|
||||
log(`Raw event (not JSON): ${event.slice(0, 200)}`)
|
||||
continue
|
||||
}
|
||||
} else {
|
||||
eventType = eventData?.type || ""
|
||||
}
|
||||
|
||||
// Check for session completion events
|
||||
const eventType = eventData?.type || ""
|
||||
const sessionId = eventData?.sessionID || eventData?.session?.id || ""
|
||||
if (sessionId && pendingResponses.has(sessionId)) {
|
||||
// Check if event indicates session completion
|
||||
const outcome = eventData?.outcome || eventData?.data?.outcome || eventData?.update?.outcome
|
||||
if (outcome && ["succeeded", "failed", "interrupted"].includes(outcome)) {
|
||||
log(`Session ${sessionId} completed with outcome: ${outcome} (from event)`)
|
||||
// Will be picked up by polling loop
|
||||
}
|
||||
|
||||
// Log ALL session events for debugging
|
||||
if (eventType.includes("session") || sessionId) {
|
||||
const preview = JSON.stringify(eventData).slice(0, 500)
|
||||
log(`Event: type=${eventType}, session=${sessionId}, data=${preview}`)
|
||||
}
|
||||
|
||||
// Log first 200 chars of interesting events
|
||||
if (eventType.includes("session") && sessionId) {
|
||||
const preview = typeof event === "string" ? event.slice(0, 200) : JSON.stringify(eventData).slice(0, 200)
|
||||
log(`Event: type=${eventType}, session=${sessionId}, preview=${preview}`)
|
||||
// Check if this event has response text for a pending prompt
|
||||
if (sessionId && pendingResponses.has(sessionId)) {
|
||||
const pending = pendingResponses.get(sessionId)!
|
||||
|
||||
// Try multiple structures to extract assistant text
|
||||
let newText = ""
|
||||
|
||||
// 1. event.data.content[].text (assistant text parts)
|
||||
const parts = eventData?.data?.content || eventData?.content || eventData?.update?.content
|
||||
if (Array.isArray(parts)) {
|
||||
for (const part of parts) {
|
||||
if (part?.type === "text" && part?.text) {
|
||||
newText += part.text
|
||||
}
|
||||
}
|
||||
}
|
||||
// 2. event.data.text (direct text)
|
||||
if (!newText && typeof eventData?.data?.text === "string") {
|
||||
newText = eventData.data.text
|
||||
}
|
||||
// 3. event.data.message.text (message text)
|
||||
if (!newText && eventData?.data?.message?.text) {
|
||||
newText = eventData.data.message.text
|
||||
}
|
||||
// 4. event.payload.text (synthetic message)
|
||||
if (!newText && typeof eventData?.payload?.text === "string") {
|
||||
newText = eventData.payload.text
|
||||
}
|
||||
// 5. event.update.content (simple content update)
|
||||
if (!newText && typeof eventData?.update?.content === "string") {
|
||||
newText = eventData.update.content
|
||||
}
|
||||
// 6. event.delta (streaming delta)
|
||||
if (!newText && typeof eventData?.delta === "string") {
|
||||
newText = eventData.delta
|
||||
}
|
||||
// 7. event.content (direct content)
|
||||
if (!newText && typeof eventData?.content === "string") {
|
||||
newText = eventData.content
|
||||
}
|
||||
|
||||
if (newText) {
|
||||
pending.text += newText
|
||||
log(`Collected ${newText.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)
|
||||
}
|
||||
}
|
||||
}
|
||||
} catch (err) {
|
||||
|
|
@ -211,74 +254,22 @@ async function runPlugin(ctx: any, options: any, dataDir: string) {
|
|||
})
|
||||
log(`Prompt sent to session ${opencodeSessionId}`)
|
||||
|
||||
// Poll session.get() to wait for completion (outcome field)
|
||||
let responseText = ""
|
||||
let lastAssistantText = ""
|
||||
|
||||
while (true) {
|
||||
await new Promise((r) => setTimeout(r, 500))
|
||||
|
||||
try {
|
||||
const sessionInfo = await ctx.session.get({ sessionID: opencodeSessionId })
|
||||
const outcome = sessionInfo?.outcome || sessionInfo?.data?.outcome
|
||||
|
||||
// Session completed — collect response text from messages
|
||||
if (outcome && ["succeeded", "failed", "interrupted"].includes(outcome)) {
|
||||
log(`Session ${opencodeSessionId} completed with outcome: ${outcome}`)
|
||||
|
||||
// Collect assistant text from session messages
|
||||
const messages = sessionInfo?.messages || sessionInfo?.data?.messages
|
||||
if (Array.isArray(messages)) {
|
||||
for (const msg of messages) {
|
||||
if (msg?.type === "assistant" || msg?.role === "assistant") {
|
||||
const content = msg?.content || msg?.text
|
||||
if (typeof content === "string") {
|
||||
lastAssistantText += content
|
||||
} else if (Array.isArray(content)) {
|
||||
for (const part of content) {
|
||||
if (part?.type === "text" && part?.text) {
|
||||
lastAssistantText += part.text
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Also check session metadata for response text
|
||||
if (!lastAssistantText) {
|
||||
const metaText = sessionInfo?.metadata?.responseText || sessionInfo?.data?.metadata?.responseText
|
||||
if (metaText) lastAssistantText = String(metaText)
|
||||
}
|
||||
|
||||
responseText = lastAssistantText || "No response text collected"
|
||||
break
|
||||
}
|
||||
|
||||
// Track assistant text as it streams (for partial responses)
|
||||
const streamMessages = sessionInfo?.streamMessages || sessionInfo?.data?.streamMessages
|
||||
if (Array.isArray(streamMessages)) {
|
||||
for (const msg of streamMessages) {
|
||||
if (msg?.type === "assistant" || msg?.role === "assistant") {
|
||||
const content = msg?.content || msg?.text
|
||||
if (typeof content === "string" && content !== lastAssistantText) {
|
||||
lastAssistantText = content
|
||||
} else if (Array.isArray(content)) {
|
||||
let text = ""
|
||||
for (const part of content) {
|
||||
if (part?.type === "text" && part?.text) text += part.text
|
||||
}
|
||||
if (text && text !== lastAssistantText) lastAssistantText = text
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
} catch (getErr: any) {
|
||||
log(`Session get error (will retry): ${String(getErr)}`)
|
||||
// Session might be temporarily unavailable, keep polling
|
||||
}
|
||||
}
|
||||
// Wait for response text from event subscription
|
||||
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
|
||||
})
|
||||
|
||||
const responseText = await responsePromise
|
||||
session.outputChars += responseText.length
|
||||
await matrix.sendReply(context, responseText)
|
||||
} catch (err: any) {
|
||||
|
|
|
|||
Loading…
Reference in New Issue