fix: wait for session outcome via polling, no timeouts
Убраны таймауты. Вместо этого: - Polling ctx.session.get() каждые 500ms - Ожидаем outcome: succeeded/failed/interrupted - Собираем текст assistant сообщений из session.get() - Event subscription логирует события для отладки - Session.get() error — retry, без abort
This commit is contained in:
parent
930cb16495
commit
8e182bdc17
144
index.ts
144
index.ts
|
|
@ -93,72 +93,45 @@ 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, timer }
|
// 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
|
||||||
timer: ReturnType<typeof setTimeout>
|
|
||||||
}>()
|
}>()
|
||||||
|
|
||||||
// V2 Event subscription — collects assistant response text from events
|
// V2 Event subscription — logs all events for debugging
|
||||||
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 })) {
|
||||||
// Parse event — V2EventEncoded is a JSON string
|
// Log event type for debugging
|
||||||
|
let eventType = ""
|
||||||
let eventData: any = event
|
let eventData: any = event
|
||||||
if (typeof event === "string") {
|
if (typeof event === "string") {
|
||||||
try {
|
try {
|
||||||
eventData = JSON.parse(event)
|
eventData = JSON.parse(event)
|
||||||
} catch { /* not JSON, use raw */ }
|
eventType = eventData?.type || ""
|
||||||
|
} catch {
|
||||||
|
eventType = "(raw string)"
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
eventType = eventData?.type || ""
|
||||||
}
|
}
|
||||||
|
|
||||||
const eventType = eventData?.type || ""
|
// Check for session completion events
|
||||||
const sessionId = eventData?.sessionID || eventData?.session?.id || ""
|
const sessionId = eventData?.sessionID || eventData?.session?.id || ""
|
||||||
|
|
||||||
// Check if this event contains assistant response text for a pending prompt
|
|
||||||
if (sessionId && pendingResponses.has(sessionId)) {
|
if (sessionId && pendingResponses.has(sessionId)) {
|
||||||
const pending = pendingResponses.get(sessionId)!
|
// Check if event indicates session completion
|
||||||
|
const outcome = eventData?.outcome || eventData?.data?.outcome || eventData?.update?.outcome
|
||||||
// Extract text from various event structures
|
if (outcome && ["succeeded", "failed", "interrupted"].includes(outcome)) {
|
||||||
let responseText = ""
|
log(`Session ${sessionId} completed with outcome: ${outcome} (from event)`)
|
||||||
|
// Will be picked up by polling loop
|
||||||
// Try: event.messages[].content[].text
|
|
||||||
const messages = eventData?.messages || eventData?.data?.messages || eventData?.update?.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") responseText += content
|
|
||||||
else if (Array.isArray(content)) {
|
|
||||||
for (const part of content) {
|
|
||||||
if (part?.type === "text" && part?.text) responseText += part.text
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// Try: event.update.content (simple text update)
|
// Log first 200 chars of interesting events
|
||||||
if (!responseText && eventData?.update?.content) {
|
if (eventType.includes("session") && sessionId) {
|
||||||
responseText = String(eventData.update.content)
|
const preview = typeof event === "string" ? event.slice(0, 200) : JSON.stringify(eventData).slice(0, 200)
|
||||||
}
|
log(`Event: type=${eventType}, session=${sessionId}, preview=${preview}`)
|
||||||
|
|
||||||
// Try: event.payload.text (synthetic message)
|
|
||||||
if (!responseText && eventData?.payload?.text) {
|
|
||||||
responseText = eventData.payload.text
|
|
||||||
}
|
|
||||||
|
|
||||||
// Try: event.data.text (generate response)
|
|
||||||
if (!responseText && eventData?.data?.text) {
|
|
||||||
responseText = eventData.data.text
|
|
||||||
}
|
|
||||||
|
|
||||||
if (responseText) {
|
|
||||||
log(`Collected response text (${responseText.length} chars) for session ${sessionId}`)
|
|
||||||
clearTimeout(pending.timer)
|
|
||||||
pendingResponses.delete(sessionId)
|
|
||||||
pending.resolve(responseText)
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
} catch (err) {
|
} catch (err) {
|
||||||
|
|
@ -236,17 +209,76 @@ async function runPlugin(ctx: any, options: any, dataDir: string) {
|
||||||
sessionID: opencodeSessionId,
|
sessionID: opencodeSessionId,
|
||||||
text: query,
|
text: query,
|
||||||
})
|
})
|
||||||
|
log(`Prompt sent to session ${opencodeSessionId}`)
|
||||||
|
|
||||||
// Wait for response text from event subscription
|
// Poll session.get() to wait for completion (outcome field)
|
||||||
const responsePromise = new Promise<string>((resolve, reject) => {
|
let responseText = ""
|
||||||
const timeout = setTimeout(() => {
|
let lastAssistantText = ""
|
||||||
pendingResponses.delete(opencodeSessionId)
|
|
||||||
reject(new Error("Timeout waiting for response (30s)"))
|
while (true) {
|
||||||
}, 30_000)
|
await new Promise((r) => setTimeout(r, 500))
|
||||||
pendingResponses.set(opencodeSessionId, { resolve, reject, timer: timeout })
|
|
||||||
})
|
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
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
const responseText = await responsePromise
|
|
||||||
session.outputChars += responseText.length
|
session.outputChars += responseText.length
|
||||||
await matrix.sendReply(context, responseText)
|
await matrix.sendReply(context, responseText)
|
||||||
} catch (err: any) {
|
} catch (err: any) {
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue