Share upstream event stream hub
This commit is contained in:
@@ -5,13 +5,15 @@ import {
|
||||
MESSAGE_STREAM_DIRECTORY_WS_PATH,
|
||||
MESSAGE_STREAM_GLOBAL_WS_PATH,
|
||||
MESSAGE_STREAM_WS_HEARTBEAT_INTERVAL_MS,
|
||||
parseSseEventEnvelope,
|
||||
sendMessageStreamWsEvent,
|
||||
sendMessageStreamWsFrame,
|
||||
} from './protocol.js';
|
||||
|
||||
const MESSAGE_STREAM_UPSTREAM_STALL_TIMEOUT_MS = 20_000;
|
||||
const MESSAGE_STREAM_UPSTREAM_RECONNECT_DELAY_MS = 250;
|
||||
import { createGlobalMessageStreamHub } from './global-hub.js';
|
||||
import {
|
||||
DEFAULT_UPSTREAM_RECONNECT_DELAY_MS,
|
||||
DEFAULT_UPSTREAM_STALL_TIMEOUT_MS,
|
||||
createUpstreamSseReader,
|
||||
} from './upstream-reader.js';
|
||||
|
||||
function shouldTriggerUpstreamHealthCheck(upstream) {
|
||||
if (!upstream) {
|
||||
@@ -71,14 +73,161 @@ export function createMessageStreamWsRuntime({
|
||||
wsClients,
|
||||
triggerHealthCheck,
|
||||
heartbeatIntervalMs = MESSAGE_STREAM_WS_HEARTBEAT_INTERVAL_MS,
|
||||
upstreamStallTimeoutMs = MESSAGE_STREAM_UPSTREAM_STALL_TIMEOUT_MS,
|
||||
upstreamReconnectDelayMs = MESSAGE_STREAM_UPSTREAM_RECONNECT_DELAY_MS,
|
||||
upstreamStallTimeoutMs = DEFAULT_UPSTREAM_STALL_TIMEOUT_MS,
|
||||
upstreamReconnectDelayMs = DEFAULT_UPSTREAM_RECONNECT_DELAY_MS,
|
||||
fetchImpl = fetch,
|
||||
globalEventHub = null,
|
||||
}) {
|
||||
const wsServer = new WebSocketServer({
|
||||
noServer: true,
|
||||
});
|
||||
|
||||
const ownsGlobalHub = !globalEventHub;
|
||||
const globalHub = globalEventHub ?? createGlobalMessageStreamHub({
|
||||
buildOpenCodeUrl,
|
||||
getOpenCodeAuthHeaders,
|
||||
fetchImpl,
|
||||
upstreamStallTimeoutMs,
|
||||
upstreamReconnectDelayMs,
|
||||
});
|
||||
|
||||
const globalClients = new Set();
|
||||
const globalClientLastEventIds = new Map();
|
||||
const globalReadyClients = new Set();
|
||||
|
||||
const replayGlobalEvents = (socket, requestedLastEventId) => {
|
||||
for (const entry of globalHub.replayAfter(requestedLastEventId)) {
|
||||
const sent = sendMessageStreamWsEvent(socket, entry.payload, {
|
||||
directory: entry.directory,
|
||||
eventId: entry.eventId,
|
||||
});
|
||||
if (!sent) {
|
||||
globalClients.delete(socket);
|
||||
globalClientLastEventIds.delete(socket);
|
||||
globalReadyClients.delete(socket);
|
||||
wsClients.delete(socket);
|
||||
return;
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
const markGlobalClientReady = (socket, requestedLastEventId) => {
|
||||
if (socket.readyState !== 1) {
|
||||
return;
|
||||
}
|
||||
|
||||
const sent = sendMessageStreamWsFrame(socket, {
|
||||
type: 'ready',
|
||||
scope: 'global',
|
||||
});
|
||||
if (!sent) {
|
||||
globalClients.delete(socket);
|
||||
globalClientLastEventIds.delete(socket);
|
||||
globalReadyClients.delete(socket);
|
||||
wsClients.delete(socket);
|
||||
return;
|
||||
}
|
||||
|
||||
globalReadyClients.add(socket);
|
||||
wsClients.add(socket);
|
||||
replayGlobalEvents(socket, requestedLastEventId);
|
||||
};
|
||||
|
||||
const closeGlobalClientsWithInitialError = ({ message, closeReason = message, triggerHealthCheckFor = null }) => {
|
||||
for (const socket of Array.from(globalClients)) {
|
||||
sendMessageStreamWsFrame(socket, { type: 'error', message });
|
||||
try {
|
||||
socket.close(1011, closeReason);
|
||||
} catch {
|
||||
}
|
||||
globalClients.delete(socket);
|
||||
globalClientLastEventIds.delete(socket);
|
||||
globalReadyClients.delete(socket);
|
||||
wsClients.delete(socket);
|
||||
}
|
||||
|
||||
if (triggerHealthCheckFor === true || (triggerHealthCheckFor && shouldTriggerUpstreamHealthCheck(triggerHealthCheckFor))) {
|
||||
triggerHealthCheck?.();
|
||||
}
|
||||
|
||||
if (ownsGlobalHub) {
|
||||
globalHub.stop();
|
||||
}
|
||||
};
|
||||
|
||||
const unsubscribeGlobalEvent = globalHub.subscribeEvent(({ envelope, payload, directory, eventId }) => {
|
||||
for (const socket of Array.from(globalClients)) {
|
||||
if (!globalReadyClients.has(socket)) {
|
||||
continue;
|
||||
}
|
||||
const sent = sendMessageStreamWsEvent(socket, payload, {
|
||||
directory,
|
||||
eventId,
|
||||
});
|
||||
if (!sent) {
|
||||
globalClients.delete(socket);
|
||||
globalClientLastEventIds.delete(socket);
|
||||
globalReadyClients.delete(socket);
|
||||
wsClients.delete(socket);
|
||||
}
|
||||
}
|
||||
|
||||
processForwardedEventPayload(payload, (syntheticPayload) => {
|
||||
for (const socket of Array.from(globalClients)) {
|
||||
if (!globalReadyClients.has(socket)) {
|
||||
continue;
|
||||
}
|
||||
const sent = sendMessageStreamWsEvent(socket, syntheticPayload, { directory: 'global' });
|
||||
if (!sent) {
|
||||
globalClients.delete(socket);
|
||||
globalClientLastEventIds.delete(socket);
|
||||
globalReadyClients.delete(socket);
|
||||
wsClients.delete(socket);
|
||||
}
|
||||
}
|
||||
});
|
||||
});
|
||||
|
||||
const unsubscribeGlobalStatus = globalHub.subscribeStatus((status) => {
|
||||
if (status.type === 'connect') {
|
||||
for (const socket of Array.from(globalClients)) {
|
||||
if (!globalReadyClients.has(socket)) {
|
||||
markGlobalClientReady(socket, globalClientLastEventIds.get(socket) ?? '');
|
||||
}
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
if (status.type === 'initial-error') {
|
||||
const error = status.error;
|
||||
if (error?.type === 'upstream_unavailable') {
|
||||
closeGlobalClientsWithInitialError({
|
||||
message: `OpenCode event stream unavailable (${error.status})`,
|
||||
closeReason: 'OpenCode event stream unavailable',
|
||||
triggerHealthCheckFor: error.response,
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
closeGlobalClientsWithInitialError({
|
||||
message: status.buildUrlFailed ? 'OpenCode service unavailable' : 'Failed to connect to OpenCode event stream',
|
||||
closeReason: status.buildUrlFailed ? 'OpenCode service unavailable' : 'Failed to connect to OpenCode event stream',
|
||||
triggerHealthCheckFor: !status.buildUrlFailed,
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
if (status.type === 'error' && status.error?.type === 'stream_error') {
|
||||
console.warn('Message stream WS proxy error:', status.error.error);
|
||||
}
|
||||
});
|
||||
|
||||
const stopGlobalHubIfUnused = () => {
|
||||
if (ownsGlobalHub && globalClients.size === 0) {
|
||||
globalHub.stop();
|
||||
}
|
||||
};
|
||||
|
||||
wsServer.on('connection', (socket, req) => {
|
||||
const rawUrl = typeof req?.url === 'string' ? req.url : MESSAGE_STREAM_GLOBAL_WS_PATH;
|
||||
const pathname = parseRequestPathname(rawUrl);
|
||||
@@ -87,23 +236,61 @@ export function createMessageStreamWsRuntime({
|
||||
const requestedLastEventId = requestUrl.searchParams.get('lastEventId')?.trim() || '';
|
||||
const requestedDirectory = requestUrl.searchParams.get('directory')?.trim() || '';
|
||||
|
||||
if (isGlobalStream) {
|
||||
const pingInterval = setInterval(() => {
|
||||
if (socket.readyState !== 1) {
|
||||
return;
|
||||
}
|
||||
|
||||
try {
|
||||
socket.ping();
|
||||
} catch {
|
||||
}
|
||||
}, heartbeatIntervalMs);
|
||||
|
||||
const heartbeatInterval = setInterval(() => {
|
||||
if (!globalHub.isConnected()) {
|
||||
return;
|
||||
}
|
||||
|
||||
sendMessageStreamWsEvent(socket, { type: 'openchamber:heartbeat', timestamp: Date.now() }, { directory: 'global' });
|
||||
}, heartbeatIntervalMs);
|
||||
|
||||
socket.on('close', () => {
|
||||
clearInterval(pingInterval);
|
||||
clearInterval(heartbeatInterval);
|
||||
globalClients.delete(socket);
|
||||
globalClientLastEventIds.delete(socket);
|
||||
globalReadyClients.delete(socket);
|
||||
wsClients.delete(socket);
|
||||
stopGlobalHubIfUnused();
|
||||
});
|
||||
|
||||
socket.on('error', () => {
|
||||
void 0;
|
||||
});
|
||||
|
||||
globalClients.add(socket);
|
||||
globalClientLastEventIds.set(socket, requestedLastEventId);
|
||||
globalHub.start();
|
||||
if (globalHub.isConnected()) {
|
||||
markGlobalClientReady(socket, requestedLastEventId);
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
const controller = new AbortController();
|
||||
let currentUpstreamAbortController = null;
|
||||
let upstreamConnected = false;
|
||||
let streamReady = false;
|
||||
let lastEventId = requestedLastEventId;
|
||||
let reader = null;
|
||||
const cleanup = () => {
|
||||
if (!controller.signal.aborted) {
|
||||
controller.abort();
|
||||
}
|
||||
if (currentUpstreamAbortController && !currentUpstreamAbortController.signal.aborted) {
|
||||
currentUpstreamAbortController.abort();
|
||||
}
|
||||
reader?.stop();
|
||||
wsClients.delete(socket);
|
||||
};
|
||||
|
||||
const wait = (ms) => new Promise((resolve) => setTimeout(resolve, ms));
|
||||
|
||||
const pingInterval = setInterval(() => {
|
||||
if (socket.readyState !== 1) {
|
||||
return;
|
||||
@@ -135,25 +322,11 @@ export function createMessageStreamWsRuntime({
|
||||
});
|
||||
|
||||
const run = async () => {
|
||||
const forwardBlock = (block) => {
|
||||
if (!block) {
|
||||
return;
|
||||
}
|
||||
|
||||
const envelope = parseSseEventEnvelope(block);
|
||||
const payload = envelope?.payload ?? null;
|
||||
if (!payload) {
|
||||
return;
|
||||
}
|
||||
|
||||
const forwardEvent = ({ envelope, payload }) => {
|
||||
const directory = isGlobalStream
|
||||
? (typeof envelope?.directory === 'string' && envelope.directory.length > 0 ? envelope.directory : 'global')
|
||||
: (requestedDirectory || envelope?.directory || 'global');
|
||||
|
||||
if (typeof envelope?.eventId === 'string' && envelope.eventId.length > 0) {
|
||||
lastEventId = envelope.eventId;
|
||||
}
|
||||
|
||||
sendMessageStreamWsEvent(socket, payload, {
|
||||
directory,
|
||||
eventId: typeof envelope?.eventId === 'string' && envelope.eventId.length > 0 ? envelope.eventId : undefined,
|
||||
@@ -165,149 +338,85 @@ export function createMessageStreamWsRuntime({
|
||||
};
|
||||
|
||||
try {
|
||||
while (!controller.signal.aborted) {
|
||||
let targetUrl;
|
||||
try {
|
||||
targetUrl = new URL(buildOpenCodeUrl(isGlobalStream ? '/global/event' : '/event', ''));
|
||||
} catch {
|
||||
sendMessageStreamWsFrame(socket, { type: 'error', message: 'OpenCode service unavailable' });
|
||||
socket.close(1011, 'OpenCode service unavailable');
|
||||
return;
|
||||
let buildUrlFailed = false;
|
||||
const closeWithInitialError = ({ message, closeReason = message, triggerHealthCheckFor = null }) => {
|
||||
sendMessageStreamWsFrame(socket, { type: 'error', message });
|
||||
socket.close(1011, closeReason);
|
||||
if (triggerHealthCheckFor === true || (triggerHealthCheckFor && shouldTriggerUpstreamHealthCheck(triggerHealthCheckFor))) {
|
||||
triggerHealthCheck?.();
|
||||
}
|
||||
reader?.stop();
|
||||
cleanup();
|
||||
};
|
||||
|
||||
if (!isGlobalStream && requestedDirectory) {
|
||||
targetUrl.searchParams.set('directory', requestedDirectory);
|
||||
}
|
||||
reader = createUpstreamSseReader({
|
||||
initialLastEventId: requestedLastEventId,
|
||||
signal: controller.signal,
|
||||
stallTimeoutMs: upstreamStallTimeoutMs,
|
||||
reconnectDelayMs: upstreamReconnectDelayMs,
|
||||
fetchImpl,
|
||||
buildUrl: () => {
|
||||
buildUrlFailed = false;
|
||||
let targetUrl;
|
||||
try {
|
||||
targetUrl = new URL(buildOpenCodeUrl('/event', ''));
|
||||
} catch {
|
||||
buildUrlFailed = true;
|
||||
throw new Error('OpenCode service unavailable');
|
||||
}
|
||||
|
||||
const headers = {
|
||||
Accept: 'text/event-stream',
|
||||
'Cache-Control': 'no-cache',
|
||||
Connection: 'keep-alive',
|
||||
...getOpenCodeAuthHeaders(),
|
||||
};
|
||||
if (lastEventId) {
|
||||
headers['Last-Event-ID'] = lastEventId;
|
||||
}
|
||||
if (requestedDirectory) {
|
||||
targetUrl.searchParams.set('directory', requestedDirectory);
|
||||
}
|
||||
|
||||
const upstreamController = new AbortController();
|
||||
currentUpstreamAbortController = upstreamController;
|
||||
const abortUpstream = () => upstreamController.abort();
|
||||
controller.signal.addEventListener('abort', abortUpstream, { once: true });
|
||||
return targetUrl;
|
||||
},
|
||||
getHeaders: getOpenCodeAuthHeaders,
|
||||
onConnect() {
|
||||
if (!streamReady) {
|
||||
sendMessageStreamWsFrame(socket, {
|
||||
type: 'ready',
|
||||
scope: 'directory',
|
||||
});
|
||||
streamReady = true;
|
||||
}
|
||||
|
||||
let upstream;
|
||||
try {
|
||||
upstream = await fetchImpl(targetUrl.toString(), {
|
||||
headers,
|
||||
signal: upstreamController.signal,
|
||||
});
|
||||
} catch {
|
||||
controller.signal.removeEventListener('abort', abortUpstream);
|
||||
currentUpstreamAbortController = null;
|
||||
upstreamConnected = true;
|
||||
},
|
||||
onDisconnect() {
|
||||
upstreamConnected = false;
|
||||
},
|
||||
onEvent: forwardEvent,
|
||||
onError(error) {
|
||||
if (controller.signal.aborted) {
|
||||
return;
|
||||
}
|
||||
if (!streamReady) {
|
||||
sendMessageStreamWsFrame(socket, { type: 'error', message: 'Failed to connect to OpenCode event stream' });
|
||||
socket.close(1011, 'Failed to connect to OpenCode event stream');
|
||||
triggerHealthCheck?.();
|
||||
return;
|
||||
}
|
||||
await wait(upstreamReconnectDelayMs);
|
||||
continue;
|
||||
}
|
||||
|
||||
if (!upstream.ok || !upstream.body) {
|
||||
controller.signal.removeEventListener('abort', abortUpstream);
|
||||
currentUpstreamAbortController = null;
|
||||
upstreamConnected = false;
|
||||
if (!streamReady) {
|
||||
sendMessageStreamWsFrame(socket, {
|
||||
type: 'error',
|
||||
message: `OpenCode event stream unavailable (${upstream.status})`,
|
||||
if (error?.type === 'upstream_unavailable') {
|
||||
closeWithInitialError({
|
||||
message: `OpenCode event stream unavailable (${error.status})`,
|
||||
closeReason: 'OpenCode event stream unavailable',
|
||||
triggerHealthCheckFor: error.response,
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
closeWithInitialError({
|
||||
message: buildUrlFailed ? 'OpenCode service unavailable' : 'Failed to connect to OpenCode event stream',
|
||||
closeReason: buildUrlFailed ? 'OpenCode service unavailable' : 'Failed to connect to OpenCode event stream',
|
||||
triggerHealthCheckFor: !buildUrlFailed,
|
||||
});
|
||||
socket.close(1011, 'OpenCode event stream unavailable');
|
||||
if (shouldTriggerUpstreamHealthCheck(upstream)) {
|
||||
triggerHealthCheck?.();
|
||||
}
|
||||
return;
|
||||
}
|
||||
await wait(upstreamReconnectDelayMs);
|
||||
continue;
|
||||
}
|
||||
|
||||
if (!streamReady) {
|
||||
sendMessageStreamWsFrame(socket, {
|
||||
type: 'ready',
|
||||
scope: isGlobalStream ? 'global' : 'directory',
|
||||
});
|
||||
streamReady = true;
|
||||
|
||||
if (isGlobalStream) {
|
||||
wsClients.add(socket);
|
||||
if (error?.type === 'stream_error') {
|
||||
console.warn('Message stream WS proxy error:', error.error);
|
||||
}
|
||||
}
|
||||
},
|
||||
});
|
||||
|
||||
upstreamConnected = true;
|
||||
const decoder = new TextDecoder();
|
||||
const reader = upstream.body.getReader();
|
||||
let buffer = '';
|
||||
let upstreamAbortReason = null;
|
||||
let stallTimer = null;
|
||||
const resetStallTimer = () => {
|
||||
if (stallTimer) {
|
||||
clearTimeout(stallTimer);
|
||||
}
|
||||
stallTimer = setTimeout(() => {
|
||||
upstreamAbortReason = 'upstream_stalled';
|
||||
upstreamConnected = false;
|
||||
upstreamController.abort();
|
||||
}, upstreamStallTimeoutMs);
|
||||
};
|
||||
|
||||
resetStallTimer();
|
||||
|
||||
try {
|
||||
while (true) {
|
||||
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) {
|
||||
const block = buffer.slice(0, separatorIndex);
|
||||
buffer = buffer.slice(separatorIndex + 2);
|
||||
forwardBlock(block);
|
||||
separatorIndex = buffer.indexOf('\n\n');
|
||||
}
|
||||
}
|
||||
|
||||
if (buffer.trim().length > 0) {
|
||||
forwardBlock(buffer.trim());
|
||||
}
|
||||
} catch (error) {
|
||||
if (!controller.signal.aborted && upstreamAbortReason !== 'upstream_stalled') {
|
||||
console.warn('Message stream WS proxy error:', error);
|
||||
}
|
||||
} finally {
|
||||
if (stallTimer) {
|
||||
clearTimeout(stallTimer);
|
||||
}
|
||||
upstreamConnected = false;
|
||||
currentUpstreamAbortController = null;
|
||||
controller.signal.removeEventListener('abort', abortUpstream);
|
||||
}
|
||||
|
||||
if (controller.signal.aborted) {
|
||||
return;
|
||||
}
|
||||
|
||||
await wait(upstreamReconnectDelayMs);
|
||||
}
|
||||
await reader.start();
|
||||
} catch (error) {
|
||||
if (!controller.signal.aborted) {
|
||||
console.warn('Message stream WS proxy error:', error);
|
||||
@@ -367,6 +476,11 @@ export function createMessageStreamWsRuntime({
|
||||
wsServer,
|
||||
async close() {
|
||||
server.off('upgrade', upgradeHandler);
|
||||
unsubscribeGlobalEvent();
|
||||
unsubscribeGlobalStatus();
|
||||
if (ownsGlobalHub) {
|
||||
globalHub.stop();
|
||||
}
|
||||
|
||||
try {
|
||||
for (const client of wsServer.clients) {
|
||||
|
||||
Reference in New Issue
Block a user