Split message stream websocket bridges
This commit is contained in:
@@ -6,9 +6,11 @@ This module contains the OpenChamber message-stream WebSocket protocol and runti
|
|||||||
## Entrypoints and structure
|
## Entrypoints and structure
|
||||||
- `packages/web/server/lib/event-stream/index.js`: public entrypoint re-exporting protocol and runtime helpers.
|
- `packages/web/server/lib/event-stream/index.js`: public entrypoint re-exporting protocol and runtime helpers.
|
||||||
- `packages/web/server/lib/event-stream/global-hub.js`: shared global upstream SSE hub for server-side subscribers and browser WS fan-out.
|
- `packages/web/server/lib/event-stream/global-hub.js`: shared global upstream SSE hub for server-side subscribers and browser WS fan-out.
|
||||||
|
- `packages/web/server/lib/event-stream/global-ws-bridge.js`: browser-facing global WS bridge that subscribes clients to the shared global hub.
|
||||||
|
- `packages/web/server/lib/event-stream/directory-ws-bridge.js`: browser-facing per-directory WS bridge that owns one scoped upstream reader per connection.
|
||||||
- `packages/web/server/lib/event-stream/protocol.js`: path constants, SSE envelope parsing, and WebSocket frame serialization helpers.
|
- `packages/web/server/lib/event-stream/protocol.js`: path constants, SSE envelope parsing, and WebSocket frame serialization helpers.
|
||||||
- `packages/web/server/lib/event-stream/upstream-reader.js`: reusable upstream SSE reader with event-id tracking, stall recovery, and reconnect handling.
|
- `packages/web/server/lib/event-stream/upstream-reader.js`: reusable upstream SSE reader with event-id tracking, stall recovery, and reconnect handling.
|
||||||
- `packages/web/server/lib/event-stream/runtime.js`: WebSocket server runtime, upgrade handling, shared global upstream hub, per-directory upstream reader orchestration, and global event broadcasting.
|
- `packages/web/server/lib/event-stream/runtime.js`: thin WebSocket server runtime for upgrade handling and path dispatch to the global/directory bridges.
|
||||||
- `packages/web/server/lib/event-stream/protocol.test.js`: unit tests for protocol helpers.
|
- `packages/web/server/lib/event-stream/protocol.test.js`: unit tests for protocol helpers.
|
||||||
- `packages/web/server/lib/event-stream/upstream-reader.test.js`: unit tests for upstream SSE reader behavior.
|
- `packages/web/server/lib/event-stream/upstream-reader.test.js`: unit tests for upstream SSE reader behavior.
|
||||||
- `packages/web/server/lib/event-stream/runtime.test.js`: unit tests for runtime-side broadcaster behavior.
|
- `packages/web/server/lib/event-stream/runtime.test.js`: unit tests for runtime-side broadcaster behavior.
|
||||||
@@ -44,10 +46,11 @@ This module contains the OpenChamber message-stream WebSocket protocol and runti
|
|||||||
- Global synthetic events such as `openchamber:session-status`, `openchamber:session-activity`, `openchamber:notification`, and `openchamber:heartbeat` are preserved on the WS path, but heartbeat frames are emitted only while an upstream SSE stream is actively attached.
|
- Global synthetic events such as `openchamber:session-status`, `openchamber:session-activity`, `openchamber:notification`, and `openchamber:heartbeat` are preserved on the WS path, but heartbeat frames are emitted only while an upstream SSE stream is actively attached.
|
||||||
- Global UI broadcasts are fan-out capable across both SSE and WS clients.
|
- Global UI broadcasts are fan-out capable across both SSE and WS clients.
|
||||||
- The reusable upstream reader centralizes SSE fetch/parsing/reconnect behavior for the WS runtime and OpenCode watcher. Additional event consumers should move to it only with parity tests for their lifecycle and error semantics.
|
- The reusable upstream reader centralizes SSE fetch/parsing/reconnect behavior for the WS runtime and OpenCode watcher. Additional event consumers should move to it only with parity tests for their lifecycle and error semantics.
|
||||||
|
- Browser transport concerns live in the WS bridge modules; server-side global stream ownership lives in `global-hub.js`.
|
||||||
|
|
||||||
## Notes for contributors
|
## Notes for contributors
|
||||||
- Keep protocol helpers pure and small so they can be unit tested without spinning up a server.
|
- Keep protocol helpers pure and small so they can be unit tested without spinning up a server.
|
||||||
- Keep runtime wiring in this module instead of `packages/web/server/index.js` unless the logic is strictly route-local glue.
|
- Keep `runtime.js` focused on WebSocket upgrade and endpoint dispatch. Put global browser-client lifecycle in `global-ws-bridge.js`, directory stream lifecycle in `directory-ws-bridge.js`, and upstream stream sharing in `global-hub.js`.
|
||||||
- Do not change upstream OpenCode transport assumptions here; OpenCode remains SSE-based.
|
- Do not change upstream OpenCode transport assumptions here; OpenCode remains SSE-based.
|
||||||
- Keep global replay bounded; do not turn it into an unbounded event log.
|
- Keep global replay bounded; do not turn it into an unbounded event log.
|
||||||
|
|
||||||
|
|||||||
@@ -0,0 +1,185 @@
|
|||||||
|
import { sendMessageStreamWsEvent, sendMessageStreamWsFrame } from './protocol.js';
|
||||||
|
import { createUpstreamSseReader } from './upstream-reader.js';
|
||||||
|
|
||||||
|
function shouldTriggerUpstreamHealthCheck(upstream) {
|
||||||
|
if (!upstream) {
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
|
||||||
|
if (!upstream.body) {
|
||||||
|
return upstream.ok || upstream.status >= 500;
|
||||||
|
}
|
||||||
|
|
||||||
|
return upstream.status >= 500;
|
||||||
|
}
|
||||||
|
|
||||||
|
export function acceptDirectoryMessageStreamWsConnection({
|
||||||
|
socket,
|
||||||
|
requestedLastEventId,
|
||||||
|
requestedDirectory,
|
||||||
|
buildOpenCodeUrl,
|
||||||
|
getOpenCodeAuthHeaders,
|
||||||
|
processForwardedEventPayload,
|
||||||
|
wsClients,
|
||||||
|
triggerHealthCheck,
|
||||||
|
heartbeatIntervalMs,
|
||||||
|
upstreamStallTimeoutMs,
|
||||||
|
upstreamReconnectDelayMs,
|
||||||
|
fetchImpl,
|
||||||
|
}) {
|
||||||
|
const controller = new AbortController();
|
||||||
|
let upstreamConnected = false;
|
||||||
|
let streamReady = false;
|
||||||
|
let reader = null;
|
||||||
|
|
||||||
|
const cleanup = () => {
|
||||||
|
if (!controller.signal.aborted) {
|
||||||
|
controller.abort();
|
||||||
|
}
|
||||||
|
reader?.stop();
|
||||||
|
wsClients.delete(socket);
|
||||||
|
};
|
||||||
|
|
||||||
|
const pingInterval = setInterval(() => {
|
||||||
|
if (socket.readyState !== 1) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
try {
|
||||||
|
socket.ping();
|
||||||
|
} catch {
|
||||||
|
}
|
||||||
|
}, heartbeatIntervalMs);
|
||||||
|
|
||||||
|
const heartbeatInterval = setInterval(() => {
|
||||||
|
if (!upstreamConnected) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
sendMessageStreamWsEvent(socket, { type: 'openchamber:heartbeat', timestamp: Date.now() }, { directory: 'global' });
|
||||||
|
}, heartbeatIntervalMs);
|
||||||
|
|
||||||
|
socket.on('close', () => {
|
||||||
|
clearInterval(pingInterval);
|
||||||
|
clearInterval(heartbeatInterval);
|
||||||
|
upstreamConnected = false;
|
||||||
|
cleanup();
|
||||||
|
});
|
||||||
|
|
||||||
|
socket.on('error', () => {
|
||||||
|
void 0;
|
||||||
|
});
|
||||||
|
|
||||||
|
const run = async () => {
|
||||||
|
const forwardEvent = ({ envelope, payload }) => {
|
||||||
|
const directory = requestedDirectory || envelope?.directory || 'global';
|
||||||
|
|
||||||
|
sendMessageStreamWsEvent(socket, payload, {
|
||||||
|
directory,
|
||||||
|
eventId: typeof envelope?.eventId === 'string' && envelope.eventId.length > 0 ? envelope.eventId : undefined,
|
||||||
|
});
|
||||||
|
|
||||||
|
processForwardedEventPayload(payload, (syntheticPayload) => {
|
||||||
|
sendMessageStreamWsEvent(socket, syntheticPayload, { directory: 'global' });
|
||||||
|
});
|
||||||
|
};
|
||||||
|
|
||||||
|
try {
|
||||||
|
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();
|
||||||
|
};
|
||||||
|
|
||||||
|
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');
|
||||||
|
}
|
||||||
|
|
||||||
|
if (requestedDirectory) {
|
||||||
|
targetUrl.searchParams.set('directory', requestedDirectory);
|
||||||
|
}
|
||||||
|
|
||||||
|
return targetUrl;
|
||||||
|
},
|
||||||
|
getHeaders: getOpenCodeAuthHeaders,
|
||||||
|
onConnect() {
|
||||||
|
if (!streamReady) {
|
||||||
|
sendMessageStreamWsFrame(socket, {
|
||||||
|
type: 'ready',
|
||||||
|
scope: 'directory',
|
||||||
|
});
|
||||||
|
streamReady = true;
|
||||||
|
}
|
||||||
|
|
||||||
|
upstreamConnected = true;
|
||||||
|
},
|
||||||
|
onDisconnect() {
|
||||||
|
upstreamConnected = false;
|
||||||
|
},
|
||||||
|
onEvent: forwardEvent,
|
||||||
|
onError(error) {
|
||||||
|
if (controller.signal.aborted) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
if (!streamReady) {
|
||||||
|
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,
|
||||||
|
});
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
if (error?.type === 'stream_error') {
|
||||||
|
console.warn('Message stream WS proxy error:', error.error);
|
||||||
|
}
|
||||||
|
},
|
||||||
|
});
|
||||||
|
|
||||||
|
await reader.start();
|
||||||
|
} catch (error) {
|
||||||
|
if (!controller.signal.aborted) {
|
||||||
|
console.warn('Message stream WS proxy error:', error);
|
||||||
|
sendMessageStreamWsFrame(socket, { type: 'error', message: 'Message stream proxy error' });
|
||||||
|
socket.close(1011, 'Message stream proxy error');
|
||||||
|
}
|
||||||
|
} finally {
|
||||||
|
cleanup();
|
||||||
|
try {
|
||||||
|
if (socket.readyState === 1 || socket.readyState === 0) {
|
||||||
|
socket.close();
|
||||||
|
}
|
||||||
|
} catch {
|
||||||
|
}
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
void run();
|
||||||
|
}
|
||||||
@@ -0,0 +1,206 @@
|
|||||||
|
import { sendMessageStreamWsEvent, sendMessageStreamWsFrame } from './protocol.js';
|
||||||
|
|
||||||
|
function shouldTriggerUpstreamHealthCheck(upstream) {
|
||||||
|
if (!upstream) {
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
|
||||||
|
if (!upstream.body) {
|
||||||
|
return upstream.ok || upstream.status >= 500;
|
||||||
|
}
|
||||||
|
|
||||||
|
return upstream.status >= 500;
|
||||||
|
}
|
||||||
|
|
||||||
|
export function createGlobalMessageStreamWsBridge({
|
||||||
|
globalHub,
|
||||||
|
ownsGlobalHub,
|
||||||
|
wsClients,
|
||||||
|
processForwardedEventPayload,
|
||||||
|
triggerHealthCheck,
|
||||||
|
heartbeatIntervalMs,
|
||||||
|
}) {
|
||||||
|
const clients = new Set();
|
||||||
|
const clientLastEventIds = new Map();
|
||||||
|
const readyClients = new Set();
|
||||||
|
|
||||||
|
const removeClient = (socket) => {
|
||||||
|
clients.delete(socket);
|
||||||
|
clientLastEventIds.delete(socket);
|
||||||
|
readyClients.delete(socket);
|
||||||
|
wsClients.delete(socket);
|
||||||
|
};
|
||||||
|
|
||||||
|
const replayEvents = (socket, requestedLastEventId) => {
|
||||||
|
for (const entry of globalHub.replayAfter(requestedLastEventId)) {
|
||||||
|
const sent = sendMessageStreamWsEvent(socket, entry.payload, {
|
||||||
|
directory: entry.directory,
|
||||||
|
eventId: entry.eventId,
|
||||||
|
});
|
||||||
|
if (!sent) {
|
||||||
|
removeClient(socket);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
const markReady = (socket, requestedLastEventId) => {
|
||||||
|
if (socket.readyState !== 1) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
const sent = sendMessageStreamWsFrame(socket, {
|
||||||
|
type: 'ready',
|
||||||
|
scope: 'global',
|
||||||
|
});
|
||||||
|
if (!sent) {
|
||||||
|
removeClient(socket);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
readyClients.add(socket);
|
||||||
|
wsClients.add(socket);
|
||||||
|
replayEvents(socket, requestedLastEventId);
|
||||||
|
};
|
||||||
|
|
||||||
|
const stopHubIfUnused = () => {
|
||||||
|
if (ownsGlobalHub && clients.size === 0) {
|
||||||
|
globalHub.stop();
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
const closeClientsWithInitialError = ({ message, closeReason = message, triggerHealthCheckFor = null }) => {
|
||||||
|
for (const socket of Array.from(clients)) {
|
||||||
|
sendMessageStreamWsFrame(socket, { type: 'error', message });
|
||||||
|
try {
|
||||||
|
socket.close(1011, closeReason);
|
||||||
|
} catch {
|
||||||
|
}
|
||||||
|
removeClient(socket);
|
||||||
|
}
|
||||||
|
|
||||||
|
if (triggerHealthCheckFor === true || (triggerHealthCheckFor && shouldTriggerUpstreamHealthCheck(triggerHealthCheckFor))) {
|
||||||
|
triggerHealthCheck?.();
|
||||||
|
}
|
||||||
|
|
||||||
|
if (ownsGlobalHub) {
|
||||||
|
globalHub.stop();
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
const unsubscribeEvent = globalHub.subscribeEvent(({ payload, directory, eventId }) => {
|
||||||
|
for (const socket of Array.from(clients)) {
|
||||||
|
if (!readyClients.has(socket)) {
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
const sent = sendMessageStreamWsEvent(socket, payload, {
|
||||||
|
directory,
|
||||||
|
eventId,
|
||||||
|
});
|
||||||
|
if (!sent) {
|
||||||
|
removeClient(socket);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
processForwardedEventPayload(payload, (syntheticPayload) => {
|
||||||
|
for (const socket of Array.from(clients)) {
|
||||||
|
if (!readyClients.has(socket)) {
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
const sent = sendMessageStreamWsEvent(socket, syntheticPayload, { directory: 'global' });
|
||||||
|
if (!sent) {
|
||||||
|
removeClient(socket);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
const unsubscribeStatus = globalHub.subscribeStatus((status) => {
|
||||||
|
if (status.type === 'connect') {
|
||||||
|
for (const socket of Array.from(clients)) {
|
||||||
|
if (!readyClients.has(socket)) {
|
||||||
|
markReady(socket, clientLastEventIds.get(socket) ?? '');
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
if (status.type === 'initial-error') {
|
||||||
|
const error = status.error;
|
||||||
|
if (error?.type === 'upstream_unavailable') {
|
||||||
|
closeClientsWithInitialError({
|
||||||
|
message: `OpenCode event stream unavailable (${error.status})`,
|
||||||
|
closeReason: 'OpenCode event stream unavailable',
|
||||||
|
triggerHealthCheckFor: error.response,
|
||||||
|
});
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
closeClientsWithInitialError({
|
||||||
|
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 accept = (socket, { requestedLastEventId = '' } = {}) => {
|
||||||
|
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);
|
||||||
|
removeClient(socket);
|
||||||
|
stopHubIfUnused();
|
||||||
|
});
|
||||||
|
|
||||||
|
socket.on('error', () => {
|
||||||
|
void 0;
|
||||||
|
});
|
||||||
|
|
||||||
|
clients.add(socket);
|
||||||
|
clientLastEventIds.set(socket, requestedLastEventId);
|
||||||
|
globalHub.start();
|
||||||
|
if (globalHub.isConnected()) {
|
||||||
|
markReady(socket, requestedLastEventId);
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
const close = () => {
|
||||||
|
unsubscribeEvent();
|
||||||
|
unsubscribeStatus();
|
||||||
|
if (ownsGlobalHub) {
|
||||||
|
globalHub.stop();
|
||||||
|
}
|
||||||
|
for (const socket of Array.from(clients)) {
|
||||||
|
removeClient(socket);
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
return {
|
||||||
|
accept,
|
||||||
|
close,
|
||||||
|
};
|
||||||
|
}
|
||||||
@@ -6,27 +6,15 @@ import {
|
|||||||
MESSAGE_STREAM_GLOBAL_WS_PATH,
|
MESSAGE_STREAM_GLOBAL_WS_PATH,
|
||||||
MESSAGE_STREAM_WS_HEARTBEAT_INTERVAL_MS,
|
MESSAGE_STREAM_WS_HEARTBEAT_INTERVAL_MS,
|
||||||
sendMessageStreamWsEvent,
|
sendMessageStreamWsEvent,
|
||||||
sendMessageStreamWsFrame,
|
|
||||||
} from './protocol.js';
|
} from './protocol.js';
|
||||||
import { createGlobalMessageStreamHub } from './global-hub.js';
|
import { createGlobalMessageStreamHub } from './global-hub.js';
|
||||||
|
import { createGlobalMessageStreamWsBridge } from './global-ws-bridge.js';
|
||||||
|
import { acceptDirectoryMessageStreamWsConnection } from './directory-ws-bridge.js';
|
||||||
import {
|
import {
|
||||||
DEFAULT_UPSTREAM_RECONNECT_DELAY_MS,
|
DEFAULT_UPSTREAM_RECONNECT_DELAY_MS,
|
||||||
DEFAULT_UPSTREAM_STALL_TIMEOUT_MS,
|
DEFAULT_UPSTREAM_STALL_TIMEOUT_MS,
|
||||||
createUpstreamSseReader,
|
|
||||||
} from './upstream-reader.js';
|
} from './upstream-reader.js';
|
||||||
|
|
||||||
function shouldTriggerUpstreamHealthCheck(upstream) {
|
|
||||||
if (!upstream) {
|
|
||||||
return true;
|
|
||||||
}
|
|
||||||
|
|
||||||
if (!upstream.body) {
|
|
||||||
return upstream.ok || upstream.status >= 500;
|
|
||||||
}
|
|
||||||
|
|
||||||
return upstream.status >= 500;
|
|
||||||
}
|
|
||||||
|
|
||||||
export function createGlobalUiEventBroadcaster({
|
export function createGlobalUiEventBroadcaster({
|
||||||
sseClients,
|
sseClients,
|
||||||
wsClients,
|
wsClients,
|
||||||
@@ -91,143 +79,15 @@ export function createMessageStreamWsRuntime({
|
|||||||
upstreamReconnectDelayMs,
|
upstreamReconnectDelayMs,
|
||||||
});
|
});
|
||||||
|
|
||||||
const globalClients = new Set();
|
const globalBridge = createGlobalMessageStreamWsBridge({
|
||||||
const globalClientLastEventIds = new Map();
|
globalHub,
|
||||||
const globalReadyClients = new Set();
|
ownsGlobalHub,
|
||||||
|
wsClients,
|
||||||
const replayGlobalEvents = (socket, requestedLastEventId) => {
|
processForwardedEventPayload,
|
||||||
for (const entry of globalHub.replayAfter(requestedLastEventId)) {
|
triggerHealthCheck,
|
||||||
const sent = sendMessageStreamWsEvent(socket, entry.payload, {
|
heartbeatIntervalMs,
|
||||||
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) => {
|
wsServer.on('connection', (socket, req) => {
|
||||||
const rawUrl = typeof req?.url === 'string' ? req.url : MESSAGE_STREAM_GLOBAL_WS_PATH;
|
const rawUrl = typeof req?.url === 'string' ? req.url : MESSAGE_STREAM_GLOBAL_WS_PATH;
|
||||||
const pathname = parseRequestPathname(rawUrl);
|
const pathname = parseRequestPathname(rawUrl);
|
||||||
@@ -237,204 +97,26 @@ export function createMessageStreamWsRuntime({
|
|||||||
const requestedDirectory = requestUrl.searchParams.get('directory')?.trim() || '';
|
const requestedDirectory = requestUrl.searchParams.get('directory')?.trim() || '';
|
||||||
|
|
||||||
if (isGlobalStream) {
|
if (isGlobalStream) {
|
||||||
const pingInterval = setInterval(() => {
|
globalBridge.accept(socket, {
|
||||||
if (socket.readyState !== 1) {
|
requestedLastEventId,
|
||||||
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;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
const controller = new AbortController();
|
acceptDirectoryMessageStreamWsConnection({
|
||||||
let upstreamConnected = false;
|
socket,
|
||||||
let streamReady = false;
|
requestedLastEventId,
|
||||||
let reader = null;
|
requestedDirectory,
|
||||||
const cleanup = () => {
|
buildOpenCodeUrl,
|
||||||
if (!controller.signal.aborted) {
|
getOpenCodeAuthHeaders,
|
||||||
controller.abort();
|
processForwardedEventPayload,
|
||||||
}
|
wsClients,
|
||||||
reader?.stop();
|
triggerHealthCheck,
|
||||||
wsClients.delete(socket);
|
heartbeatIntervalMs,
|
||||||
};
|
upstreamStallTimeoutMs,
|
||||||
|
upstreamReconnectDelayMs,
|
||||||
const pingInterval = setInterval(() => {
|
fetchImpl,
|
||||||
if (socket.readyState !== 1) {
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
|
|
||||||
try {
|
|
||||||
socket.ping();
|
|
||||||
} catch {
|
|
||||||
}
|
|
||||||
}, heartbeatIntervalMs);
|
|
||||||
|
|
||||||
const heartbeatInterval = setInterval(() => {
|
|
||||||
if (!upstreamConnected) {
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
|
|
||||||
sendMessageStreamWsEvent(socket, { type: 'openchamber:heartbeat', timestamp: Date.now() }, { directory: 'global' });
|
|
||||||
}, heartbeatIntervalMs);
|
|
||||||
|
|
||||||
socket.on('close', () => {
|
|
||||||
clearInterval(pingInterval);
|
|
||||||
clearInterval(heartbeatInterval);
|
|
||||||
upstreamConnected = false;
|
|
||||||
cleanup();
|
|
||||||
});
|
});
|
||||||
|
|
||||||
socket.on('error', () => {
|
|
||||||
void 0;
|
|
||||||
});
|
|
||||||
|
|
||||||
const run = async () => {
|
|
||||||
const forwardEvent = ({ envelope, payload }) => {
|
|
||||||
const directory = isGlobalStream
|
|
||||||
? (typeof envelope?.directory === 'string' && envelope.directory.length > 0 ? envelope.directory : 'global')
|
|
||||||
: (requestedDirectory || envelope?.directory || 'global');
|
|
||||||
|
|
||||||
sendMessageStreamWsEvent(socket, payload, {
|
|
||||||
directory,
|
|
||||||
eventId: typeof envelope?.eventId === 'string' && envelope.eventId.length > 0 ? envelope.eventId : undefined,
|
|
||||||
});
|
|
||||||
|
|
||||||
processForwardedEventPayload(payload, (syntheticPayload) => {
|
|
||||||
sendMessageStreamWsEvent(socket, syntheticPayload, { directory: 'global' });
|
|
||||||
});
|
|
||||||
};
|
|
||||||
|
|
||||||
try {
|
|
||||||
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();
|
|
||||||
};
|
|
||||||
|
|
||||||
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');
|
|
||||||
}
|
|
||||||
|
|
||||||
if (requestedDirectory) {
|
|
||||||
targetUrl.searchParams.set('directory', requestedDirectory);
|
|
||||||
}
|
|
||||||
|
|
||||||
return targetUrl;
|
|
||||||
},
|
|
||||||
getHeaders: getOpenCodeAuthHeaders,
|
|
||||||
onConnect() {
|
|
||||||
if (!streamReady) {
|
|
||||||
sendMessageStreamWsFrame(socket, {
|
|
||||||
type: 'ready',
|
|
||||||
scope: 'directory',
|
|
||||||
});
|
|
||||||
streamReady = true;
|
|
||||||
}
|
|
||||||
|
|
||||||
upstreamConnected = true;
|
|
||||||
},
|
|
||||||
onDisconnect() {
|
|
||||||
upstreamConnected = false;
|
|
||||||
},
|
|
||||||
onEvent: forwardEvent,
|
|
||||||
onError(error) {
|
|
||||||
if (controller.signal.aborted) {
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
|
|
||||||
if (!streamReady) {
|
|
||||||
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,
|
|
||||||
});
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
|
|
||||||
if (error?.type === 'stream_error') {
|
|
||||||
console.warn('Message stream WS proxy error:', error.error);
|
|
||||||
}
|
|
||||||
},
|
|
||||||
});
|
|
||||||
|
|
||||||
await reader.start();
|
|
||||||
} catch (error) {
|
|
||||||
if (!controller.signal.aborted) {
|
|
||||||
console.warn('Message stream WS proxy error:', error);
|
|
||||||
sendMessageStreamWsFrame(socket, { type: 'error', message: 'Message stream proxy error' });
|
|
||||||
socket.close(1011, 'Message stream proxy error');
|
|
||||||
}
|
|
||||||
} finally {
|
|
||||||
cleanup();
|
|
||||||
try {
|
|
||||||
if (socket.readyState === 1 || socket.readyState === 0) {
|
|
||||||
socket.close();
|
|
||||||
}
|
|
||||||
} catch {
|
|
||||||
}
|
|
||||||
}
|
|
||||||
};
|
|
||||||
|
|
||||||
void run();
|
|
||||||
});
|
});
|
||||||
|
|
||||||
const upgradeHandler = (req, socket, head) => {
|
const upgradeHandler = (req, socket, head) => {
|
||||||
@@ -476,11 +158,7 @@ export function createMessageStreamWsRuntime({
|
|||||||
wsServer,
|
wsServer,
|
||||||
async close() {
|
async close() {
|
||||||
server.off('upgrade', upgradeHandler);
|
server.off('upgrade', upgradeHandler);
|
||||||
unsubscribeGlobalEvent();
|
globalBridge.close();
|
||||||
unsubscribeGlobalStatus();
|
|
||||||
if (ownsGlobalHub) {
|
|
||||||
globalHub.stop();
|
|
||||||
}
|
|
||||||
|
|
||||||
try {
|
try {
|
||||||
for (const client of wsServer.clients) {
|
for (const client of wsServer.clients) {
|
||||||
|
|||||||
Reference in New Issue
Block a user