Files
openchamber/packages/web/server/lib/opencode/watcher.js
T

116 lines
2.8 KiB
JavaScript
Raw Normal View History

2026-04-24 11:58:16 +03:00
import { createUpstreamSseReader } from '../event-stream/upstream-reader.js';
export const createOpenCodeWatcherRuntime = (deps) => {
const {
waitForOpenCodePort,
buildOpenCodeUrl,
getOpenCodeAuthHeaders,
onPayload,
2026-04-24 11:58:16 +03:00
fetchImpl = fetch,
upstreamStallTimeoutMs,
upstreamReconnectDelayMs = 1000,
globalEventHub = null,
} = deps;
let abortController = null;
2026-04-24 11:58:16 +03:00
let reader = null;
let unsubscribeEvent = null;
let unsubscribeStatus = null;
const unwrapGlobalEventPayload = (eventData) => {
if (!eventData || typeof eventData !== 'object') {
return null;
}
if (eventData.payload && typeof eventData.payload === 'object') {
return eventData.payload;
}
return eventData;
};
const start = async () => {
if (abortController) {
return;
}
await waitForOpenCodePort();
abortController = new AbortController();
const signal = abortController.signal;
2026-04-24 11:58:16 +03:00
if (globalEventHub) {
unsubscribeEvent = globalEventHub.subscribeEvent((event) => {
const payload = unwrapGlobalEventPayload(event.payload);
if (!payload || typeof payload !== 'object') {
return;
}
onPayload(payload);
});
unsubscribeStatus = globalEventHub.subscribeStatus((status) => {
if (signal.aborted) {
return;
}
if (status.type === 'connect') {
console.log('[PushWatcher] connected');
2026-04-24 11:58:16 +03:00
return;
}
2026-04-24 11:58:16 +03:00
if (status.type === 'error' || status.type === 'initial-error') {
console.warn('[PushWatcher] disconnected', status.error?.error?.message ?? status.error?.message ?? status.error);
}
});
globalEventHub.start();
return;
}
2026-04-24 11:58:16 +03:00
reader = createUpstreamSseReader({
signal,
buildUrl: () => buildOpenCodeUrl('/global/event', ''),
getHeaders: getOpenCodeAuthHeaders,
fetchImpl,
stallTimeoutMs: upstreamStallTimeoutMs,
reconnectDelayMs: upstreamReconnectDelayMs,
onConnect() {
console.log('[PushWatcher] connected');
},
onEvent(event) {
const payload = unwrapGlobalEventPayload(event.payload);
if (!payload || typeof payload !== 'object') {
return;
}
onPayload(payload);
},
onError(error) {
if (signal.aborted) {
return;
}
console.warn('[PushWatcher] disconnected', error?.error?.message ?? error?.message ?? error);
},
});
2026-04-24 11:58:16 +03:00
void reader.start();
};
const stop = () => {
if (!abortController) {
return;
}
try {
abortController.abort();
2026-04-24 11:58:16 +03:00
reader?.stop();
unsubscribeEvent?.();
unsubscribeStatus?.();
} catch {
}
2026-04-24 11:58:16 +03:00
reader = null;
unsubscribeEvent = null;
unsubscribeStatus = null;
abortController = null;
};
return {
start,
stop,
};
};