Files
memind/session-reply-wait.mjs
T
john 3155b5805f feat(agent-run): optional session reply idle timeout and finish reason
MEMIND_SESSION_REPLY_IDLE_TIMEOUT_MS (off by default, min 30s) fails a reply
that stops streaming while no tool is running, and is recoverable like other
finish errors. The goose Finish reason is surfaced in session_finished.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-09-23 14:25:24 +08:00

299 lines
9.9 KiB
JavaScript

import { Readable } from 'node:stream';
function parseSseDataLines(frame) {
let data = '';
for (const line of String(frame ?? '').split('\n')) {
if (line.startsWith('data:')) data += line.slice(5).trim();
}
return data;
}
export function parseSessionStreamEvent(frame) {
const data = parseSseDataLines(frame);
if (!data) return null;
try {
return JSON.parse(data);
} catch {
return null;
}
}
export function eventMatchesRequest(event, requestId) {
if (!requestId) return true;
const routingId = event?.chat_request_id ?? event?.request_id ?? null;
return !routingId || routingId === requestId;
}
function messageContentItems(event) {
const candidates = [
event?.message?.content,
event?.data?.message?.content,
event?.content,
];
return candidates.find((content) => Array.isArray(content)) ?? [];
}
export function classifySessionProviderErrorMessage(message) {
const normalized = String(message ?? '').trim();
if (!normalized) return 'SESSION_REPLY_ERROR';
if (
/unknown variant [`']?image_url[`']?.*expected [`']?text/i.test(normalized)
|| /failed to deserialize.*image_url/i.test(normalized)
) {
return 'SESSION_VISUAL_CONTEXT_UNSUPPORTED';
}
if (/tool_calls|tool_call_id|insufficient tool messages/i.test(normalized)) {
return 'SESSION_TOOL_HISTORY_POISONED';
}
// DeepSeek/Kimi thinking mode: assistant tool-call turns must replay reasoning_content.
if (/reasoning_content/i.test(normalized)) {
return 'SESSION_REASONING_CONTENT_POISONED';
}
return 'SESSION_REPLY_ERROR';
}
function providerReplyError(event) {
if (event?.type !== 'Message') return null;
const text = messageContentItems(event)
.filter((item) => item?.type === 'text')
.map((item) => String(item?.text ?? ''))
.join('')
.trim();
if (!/^Ran into this error:\s*/i.test(text)) return null;
const normalized = text
.replace(/^Ran into this error:\s*/i, '')
.replace(/\n\nPlease retry if you think this is a transient or recoverable error\.?\s*$/i, '')
.trim();
const error = new Error(normalized || 'session reply failed');
error.code = classifySessionProviderErrorMessage(normalized);
if (error.code === 'SESSION_VISUAL_CONTEXT_UNSUPPORTED') {
error.retryable = false;
}
return error;
}
function parseGenerateImageResult(toolResult) {
const content = toolResult?.value?.content;
if (!Array.isArray(content)) return null;
for (const item of content) {
if (item?.type !== 'text' || !String(item.text ?? '').trim()) continue;
try {
const result = JSON.parse(String(item.text));
const mimeType = String(result?.source?.mimeType ?? result?.asset?.mimeType ?? '').toLowerCase();
if (
result?.ok === true
&& String(result?.jobId ?? '').trim()
&& /^image\/(?:png|jpeg|webp)$/.test(mimeType)
) {
return {
jobId: String(result.jobId),
mimeType,
assetId: String(result?.asset?.id ?? '').trim() || null,
publicUrl: String(result?.asset?.publicUrl ?? '').trim() || null,
workspaceRelativePath: String(result?.asset?.workspaceRelativePath ?? '').trim() || null,
};
}
} catch {
// Ignore non-JSON tool text; it cannot prove a successful image_make delivery.
}
}
return null;
}
function createToolEvidenceCollector() {
const requestNames = new Map();
const pendingToolIds = new Set();
const calls = new Set();
const successfulCalls = new Set();
let generateImage = { called: false, succeeded: false };
return {
observe(event) {
for (const item of messageContentItems(event)) {
if (item?.type === 'toolRequest') {
if (item.id) pendingToolIds.add(String(item.id));
const name = String(item?.toolCall?.value?.name ?? '').trim();
if (!name) continue;
calls.add(name);
if (item.id) requestNames.set(String(item.id), name);
if (name.endsWith('generate_image')) generateImage = { called: true, succeeded: false };
continue;
}
if (item?.type !== 'toolResponse') continue;
pendingToolIds.delete(String(item.id ?? ''));
const name = requestNames.get(String(item.id ?? '')) ?? '';
const result = item.toolResult ?? {};
const succeeded = result.status === 'success' && result?.value?.isError !== true;
if (name && succeeded) successfulCalls.add(name);
if (!name.endsWith('generate_image')) continue;
const generated = succeeded ? parseGenerateImageResult(result) : null;
generateImage = generated
? { called: true, succeeded: true, ...generated }
: { called: true, succeeded: false };
}
},
hasPendingTools() {
return pendingToolIds.size > 0;
},
snapshot() {
return {
calls: [...calls],
successfulCalls: [...successfulCalls],
generateImage,
};
},
};
}
export function resolveSessionReplyIdleTimeoutMs(env = process.env) {
const value = Number(env.MEMIND_SESSION_REPLY_IDLE_TIMEOUT_MS ?? 0);
return Number.isFinite(value) && value > 0 ? Math.max(30_000, Math.floor(value)) : 0;
}
const IDLE_RESETTING_EVENT_TYPES = new Set(['Message', 'Notification', 'UpdateConversation']);
export async function consumeSessionEventsUntilFinish(
body,
{
requestId = null,
timeoutMs = 15 * 60 * 1000,
onEvent = null,
requireRequestActivityBeforeUnscopedFinish = false,
emptyFinishGraceMs = 20_000,
idleTimeoutMs = 0,
} = {},
) {
if (!body) {
const err = new Error('session event stream unavailable');
err.code = 'SESSION_EVENT_STREAM_UNAVAILABLE';
throw err;
}
const reader = Readable.fromWeb(body);
const decoder = new TextDecoder();
const toolEvidence = createToolEvidenceCollector();
let pendingProviderError = null;
let sawRequestScopedActivity = false;
let buffer = '';
const deadline = Date.now() + Math.max(1, Number(timeoutMs) || 1);
let timeoutHandle = null;
const armTimeout = (delayMs, { code, message }) => {
if (timeoutHandle) clearTimeout(timeoutHandle);
timeoutHandle = setTimeout(() => {
const error = new Error(message);
error.code = code;
error.retryable = false;
reader.destroy(error);
}, Math.max(1, delayMs));
};
const armOriginalDeadline = () => {
armTimeout(Math.max(1, deadline - Date.now()), {
code: 'SESSION_REPLY_TIMEOUT',
message: `session reply timed out after ${timeoutMs}ms`,
});
};
// Idle only counts while waiting on the model; tool calls (aider/openhands/image)
// can legitimately run silent for minutes and stay bounded by the overall deadline.
const idleMs = Math.max(0, Number(idleTimeoutMs) || 0);
const armActivityTimeout = () => {
const remaining = Math.max(1, deadline - Date.now());
if (idleMs > 0 && idleMs < remaining && !toolEvidence.hasPendingTools()) {
armTimeout(idleMs, {
code: 'SESSION_REPLY_IDLE_TIMEOUT',
message: `session reply produced no events for ${idleMs}ms while waiting on the model`,
});
return;
}
armOriginalDeadline();
};
armActivityTimeout();
try {
for await (const chunk of reader) {
if (Date.now() > deadline) {
const err = new Error(`session reply timed out after ${timeoutMs}ms`);
err.code = 'SESSION_REPLY_TIMEOUT';
err.retryable = false;
throw err;
}
buffer += decoder.decode(chunk, { stream: true });
const frames = buffer.split('\n\n');
buffer = frames.pop() ?? '';
for (const frame of frames) {
const trimmed = frame.trim();
if (!trimmed || trimmed.startsWith(':')) continue;
const event = parseSessionStreamEvent(trimmed);
if (!event) continue;
if (!eventMatchesRequest(event, requestId)) continue;
const routingId = event?.chat_request_id ?? event?.request_id ?? null;
if (requestId && routingId === requestId) {
sawRequestScopedActivity = true;
}
pendingProviderError = providerReplyError(event) ?? pendingProviderError;
toolEvidence.observe(event);
onEvent?.(event);
if (event.type === 'Error') {
const err = new Error(String(event.error ?? 'session reply failed'));
err.code = classifySessionProviderErrorMessage(err.message);
err.retryable = false;
throw err;
}
if (IDLE_RESETTING_EVENT_TYPES.has(event.type)) {
armActivityTimeout();
}
if (event.type !== 'Finish') continue;
const tokenState = event.token_state ?? null;
const tokenActivity = [
tokenState?.inputTokens,
tokenState?.outputTokens,
tokenState?.totalTokens,
tokenState?.accumulatedInputTokens,
tokenState?.accumulatedOutputTokens,
tokenState?.accumulatedTotalTokens,
].some((value) => Number(value) > 0);
const toolActivity = toolEvidence.snapshot().calls.length > 0;
if (
requireRequestActivityBeforeUnscopedFinish
&& requestId
&& (
(!routingId && !sawRequestScopedActivity)
|| (!tokenActivity && !toolActivity)
)
) {
armTimeout(
Math.min(
Math.max(1, deadline - Date.now()),
Math.max(1, Number(emptyFinishGraceMs) || 1),
),
{
code: 'SESSION_EMPTY_FINISH',
message: 'session emitted an empty Finish without model or tool activity',
},
);
continue;
}
if (pendingProviderError) throw pendingProviderError;
return {
finishEvent: event,
finishReason: String(event.reason ?? '').trim() || 'stop',
tokenState,
toolEvidence: toolEvidence.snapshot(),
};
}
}
} finally {
if (timeoutHandle) clearTimeout(timeoutHandle);
}
const err = new Error('session event stream ended before Finish');
err.code = 'SESSION_REPLY_INCOMPLETE';
err.retryable = false;
throw err;
}