fix: collect response text from event subscription instead of prompt return
session.prompt() non-blocking — возвращает объект-приём, а не текст.
Текст ответа приходит через ctx.event.subscribe() как live-only events.
Подход:
- pendingResponses Map: sessionID -> {resolve, reject, timer}
- Event subscription: парсит события, извлекает текст из assistant сообщений
- Message handler: создаёт Promise, ждёт события (timeout 30s)
- Парсим несколько структур событий на всякий случай
This commit is contained in:
parent
e633558483
commit
930cb16495
86
index.ts
86
index.ts
|
|
@ -93,17 +93,72 @@ 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 }>()
|
||||||
|
|
||||||
// V2 Event subscription for response streaming
|
// Pending prompt responses — maps sessionID -> { resolve, reject, timer }
|
||||||
|
const pendingResponses = new Map<string, {
|
||||||
|
resolve: (text: string) => void
|
||||||
|
reject: (err: any) => void
|
||||||
|
timer: ReturnType<typeof setTimeout>
|
||||||
|
}>()
|
||||||
|
|
||||||
|
// V2 Event subscription — collects assistant response text from events
|
||||||
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 })) {
|
||||||
const eventType = typeof event === "object" && event !== null && "type" in event ? (event as { type: string }).type : ""
|
// Parse event — V2EventEncoded is a JSON string
|
||||||
if (eventType === "session.updated") {
|
let eventData: any = event
|
||||||
const sessionId = (event as { sessionID?: string }).sessionID
|
if (typeof event === "string") {
|
||||||
if (!sessionId) continue
|
try {
|
||||||
const update = (event as { update?: { type: string; content?: string; toolName?: string } }).update
|
eventData = JSON.parse(event)
|
||||||
if (!update) continue
|
} catch { /* not JSON, use raw */ }
|
||||||
log(`V2 event: session=${sessionId}, update.type=${update.type}`)
|
}
|
||||||
|
|
||||||
|
const eventType = eventData?.type || ""
|
||||||
|
const sessionId = eventData?.sessionID || eventData?.session?.id || ""
|
||||||
|
|
||||||
|
// Check if this event contains assistant response text for a pending prompt
|
||||||
|
if (sessionId && pendingResponses.has(sessionId)) {
|
||||||
|
const pending = pendingResponses.get(sessionId)!
|
||||||
|
|
||||||
|
// Extract text from various event structures
|
||||||
|
let responseText = ""
|
||||||
|
|
||||||
|
// 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)
|
||||||
|
if (!responseText && eventData?.update?.content) {
|
||||||
|
responseText = String(eventData.update.content)
|
||||||
|
}
|
||||||
|
|
||||||
|
// 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) {
|
||||||
|
|
@ -176,11 +231,22 @@ async function runPlugin(ctx: any, options: any, dataDir: string) {
|
||||||
}
|
}
|
||||||
|
|
||||||
try {
|
try {
|
||||||
const result = await ctx.session.prompt({
|
// Trigger prompt — returns admission receipt, not response text
|
||||||
|
await ctx.session.prompt({
|
||||||
sessionID: opencodeSessionId,
|
sessionID: opencodeSessionId,
|
||||||
text: query,
|
text: query,
|
||||||
})
|
})
|
||||||
const responseText = result ? String(result) : "No response"
|
|
||||||
|
// Wait for response text from event subscription
|
||||||
|
const responsePromise = new Promise<string>((resolve, reject) => {
|
||||||
|
const timeout = setTimeout(() => {
|
||||||
|
pendingResponses.delete(opencodeSessionId)
|
||||||
|
reject(new Error("Timeout waiting for response (30s)"))
|
||||||
|
}, 30_000)
|
||||||
|
pendingResponses.set(opencodeSessionId, { resolve, reject, timer: timeout })
|
||||||
|
})
|
||||||
|
|
||||||
|
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