Files
openchamber/packages/web/server/lib/event-stream/upstream-reader.js
T
Shyamalan KannanandBohdan Triapitsyn a9ee499ae4 fix: add concurrency controls for multiple sessions using the same provider (#1069)
* fix: add concurrency controls for multiple sessions using the same provider

Adds OS-inspired scheduling primitives (from HiveMind/AIMD research) to prevent
concurrent sessions from the same provider from experiencing slowdowns, random
stops, and cascading failures.

Server-side:
- Health check skips OpenCode restart when sessions are actively busy — a busy
  server under concurrent load can fail the health check timeout without being
  dead. Staleness guard forces restart if unhealthy+busy persists >2 minutes.
- Upstream SSE stall timeout scaled from 20s to 60s to avoid unnecessary
  reconnections when multiple sessions are waiting for LLM responses.

Client-side (HiveMind primitives, arXiv:2604.17111):
- Transparent retry with exponential backoff (1s→2s→4s, max 32s) for
  429/502/503/504 errors — the #1 most effective primitive from the paper.
- Circuit breaker: opens after 3 consecutive retryable errors, cooldown
  doubles each trip (30s→60s→120s, capped 128s), matching TCP AIMD.
- Per-provider session tracking with TTL eviction (1h idle sweep).
- Fetch-level retry gated on AbortError/TypeError only (not DNS failures).

Refs github-code-review skill findings (all 8 issues resolved).

* fix: use definite assignment assertion for response variable

Fixes TS2454: Variable 'response' is used before being assigned
in strict mode. The for-loop body always assigns it on every path
that reaches the post-loop code, but TS can't prove that.

* fix: add cleanupSession to error paths and remove unreachable code

P1 fixes (Greptile review):
- cleanupSession called on fetch error throw path
- cleanupSession called on non-retryable HTTP error throw path
- Removed unreachable post-loop code (loop always terminates via return or throw)

Adds explicit post-loop throw to satisfy TypeScript strict return check.

* fix: address Greptile review feedback on concurrent session controls

Removes client-side session tracking that leaked on normal completion paths.

The session tracking was redundant — the server-side health check already reads from

sessionRuntime.getSessionActivitySnapshot() for busy-session detection.

Changes:

- Remove activeSessions Set and all session-tracking functions from provider-tracker

- Remove trackSessionStarted/cleanupSession calls from client.ts

- Remove unreachable (response as Response) block after retry loop

- Make upstreamStallTimeoutMs conditional: 60s when >1 sessions, 20s otherwise

Refs #1069

* fix: enforce dynamic concurrency safeguards

---------

Co-authored-by: Bohdan Triapitsyn <artmore@protonmail.com>
2026-05-01 14:57:31 +03:00

212 lines
6.1 KiB
JavaScript

import { parseSseEventEnvelope } from './protocol.js';
export const DEFAULT_UPSTREAM_STALL_TIMEOUT_MS = 20_000;
export const UPSTREAM_STALL_TIMEOUT_CONCURRENT_MS = DEFAULT_UPSTREAM_STALL_TIMEOUT_MS * 3;
export const DEFAULT_UPSTREAM_RECONNECT_DELAY_MS = 250;
function resolveTimeoutMs(value, fallback) {
const resolved = typeof value === 'function' ? value() : value;
return Number.isFinite(resolved) ? resolved : fallback;
}
function waitForReconnectDelay(ms, signal) {
if (signal?.aborted) {
return Promise.resolve();
}
return new Promise((resolve) => {
const timeout = setTimeout(resolve, Math.max(0, ms));
signal?.addEventListener('abort', () => {
clearTimeout(timeout);
resolve();
}, { once: true });
});
}
function normalizeHeaders(headers) {
if (!headers || typeof headers !== 'object') {
return {};
}
return { ...headers };
}
export function createUpstreamSseReader({
buildUrl,
getHeaders = () => ({}),
fetchImpl = fetch,
parseBlock = parseSseEventEnvelope,
initialLastEventId = '',
signal,
stallTimeoutMs = DEFAULT_UPSTREAM_STALL_TIMEOUT_MS,
reconnectDelayMs = DEFAULT_UPSTREAM_RECONNECT_DELAY_MS,
onEvent,
onConnect,
onDisconnect,
onError,
}) {
let running = null;
let stopped = false;
let activeController = null;
let lastEventId = typeof initialLastEventId === 'string' ? initialLastEventId : '';
const stop = () => {
stopped = true;
if (activeController && !activeController.signal.aborted) {
activeController.abort();
}
};
signal?.addEventListener('abort', stop, { once: true });
const start = () => {
if (running) {
return running;
}
stopped = false;
running = (async () => {
while (!stopped && !signal?.aborted) {
const controller = new AbortController();
activeController = controller;
const abortActive = () => controller.abort();
signal?.addEventListener('abort', abortActive, { once: true });
let abortReason = null;
let stallTimer = null;
const clearStallTimer = () => {
if (stallTimer) {
clearTimeout(stallTimer);
stallTimer = null;
}
};
const resetStallTimer = () => {
clearStallTimer();
const currentStallTimeoutMs = resolveTimeoutMs(stallTimeoutMs, DEFAULT_UPSTREAM_STALL_TIMEOUT_MS);
if (currentStallTimeoutMs <= 0) {
return;
}
stallTimer = setTimeout(() => {
abortReason = 'upstream_stalled';
controller.abort();
}, currentStallTimeoutMs);
};
try {
const url = buildUrl();
const headers = {
Accept: 'text/event-stream',
'Cache-Control': 'no-cache',
Connection: 'keep-alive',
...normalizeHeaders(getHeaders()),
};
if (lastEventId) {
headers['Last-Event-ID'] = lastEventId;
}
const response = await fetchImpl(url.toString(), {
headers,
signal: controller.signal,
});
if (!response?.ok || !response.body) {
onError?.({
type: 'upstream_unavailable',
status: response?.status ?? 0,
response,
});
await waitForReconnectDelay(reconnectDelayMs, signal);
continue;
}
onConnect?.({ response, lastEventId });
const decoder = new TextDecoder();
const reader = response.body.getReader();
let buffer = '';
resetStallTimer();
while (!stopped && !signal?.aborted) {
const { value, done } = await reader.read();
if (done) {
break;
}
resetStallTimer();
buffer += decoder.decode(value, { stream: true }).replace(/\r\n/g, '\n');
let separatorIndex = buffer.indexOf('\n\n');
while (separatorIndex !== -1 && !stopped && !signal?.aborted) {
const block = buffer.slice(0, separatorIndex);
buffer = buffer.slice(separatorIndex + 2);
const envelope = parseBlock(block);
if (envelope?.payload) {
if (typeof envelope.eventId === 'string' && envelope.eventId.length > 0) {
lastEventId = envelope.eventId;
}
onEvent?.({
block,
envelope,
payload: envelope.payload,
eventId: envelope.eventId,
directory: envelope.directory,
});
}
separatorIndex = buffer.indexOf('\n\n');
}
}
if (!stopped && !signal?.aborted && buffer.trim().length > 0) {
const block = buffer.trim();
const envelope = parseBlock(block);
if (envelope?.payload) {
if (typeof envelope.eventId === 'string' && envelope.eventId.length > 0) {
lastEventId = envelope.eventId;
}
onEvent?.({
block,
envelope,
payload: envelope.payload,
eventId: envelope.eventId,
directory: envelope.directory,
});
}
}
} catch (error) {
if (!stopped && !signal?.aborted && abortReason !== 'upstream_stalled') {
onError?.({
type: 'stream_error',
error,
});
}
} finally {
clearStallTimer();
signal?.removeEventListener('abort', abortActive);
if (activeController === controller) {
activeController = null;
}
onDisconnect?.({ reason: abortReason ?? (stopped || signal?.aborted ? 'stopped' : 'closed') });
}
if (!stopped && !signal?.aborted) {
await waitForReconnectDelay(reconnectDelayMs, signal);
}
}
})().finally(() => {
running = null;
});
return running;
};
return {
start,
stop,
getLastEventId() {
return lastEventId;
},
};
}