fix: harden long conversation continuity
Memind CI / Test, build, and release guards (push) Successful in 1m28s
Memind CI / Test, build, and release guards (push) Successful in 1m28s
This commit is contained in:
+88
-24
@@ -145,6 +145,8 @@ export async function consumeSessionEventsUntilFinish(
|
||||
requestId = null,
|
||||
timeoutMs = 15 * 60 * 1000,
|
||||
onEvent = null,
|
||||
requireRequestActivityBeforeUnscopedFinish = false,
|
||||
emptyFinishGraceMs = 20_000,
|
||||
} = {},
|
||||
) {
|
||||
if (!body) {
|
||||
@@ -157,44 +159,106 @@ export async function consumeSessionEventsUntilFinish(
|
||||
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;
|
||||
|
||||
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;
|
||||
}
|
||||
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));
|
||||
};
|
||||
|
||||
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;
|
||||
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);
|
||||
const armOriginalDeadline = () => {
|
||||
armTimeout(Math.max(1, deadline - Date.now()), {
|
||||
code: 'SESSION_REPLY_TIMEOUT',
|
||||
message: `session reply timed out after ${timeoutMs}ms`,
|
||||
});
|
||||
};
|
||||
|
||||
armOriginalDeadline();
|
||||
|
||||
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;
|
||||
}
|
||||
if (event.type === 'Finish') {
|
||||
|
||||
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 (event.type === 'Message') {
|
||||
armOriginalDeadline();
|
||||
}
|
||||
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,
|
||||
tokenState: event.token_state ?? null,
|
||||
tokenState,
|
||||
toolEvidence: toolEvidence.snapshot(),
|
||||
};
|
||||
}
|
||||
}
|
||||
} finally {
|
||||
if (timeoutHandle) clearTimeout(timeoutHandle);
|
||||
}
|
||||
|
||||
const err = new Error('session event stream ended before Finish');
|
||||
|
||||
Reference in New Issue
Block a user