Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
88 changes: 74 additions & 14 deletions apps/dashboard/lib/activity/client-stream.ts
Original file line number Diff line number Diff line change
Expand Up @@ -13,11 +13,17 @@ export interface ActivityStreamConnectionOptions {
onEvent: (event: ActivityEvent) => void;
onFallback: () => void;
onReady?: () => void;
reconnectBaseDelayMs?: number;
reconnectJitterRatio?: number;
reconnectMaxDelayMs?: number;
url?: string;
}

const DEFAULT_CONNECTION_TIMEOUT_MS = 10_000;
const DEFAULT_HEARTBEAT_TIMEOUT_MS = 45_000;
const DEFAULT_RECONNECT_BASE_DELAY_MS = 1_000;
const DEFAULT_RECONNECT_JITTER_RATIO = 0.25;
const DEFAULT_RECONNECT_MAX_DELAY_MS = 30_000;

export function connectActivityStream({
connectionTimeoutMs = DEFAULT_CONNECTION_TIMEOUT_MS,
Expand All @@ -26,12 +32,18 @@ export function connectActivityStream({
onEvent,
onFallback,
onReady = () => {},
reconnectBaseDelayMs = DEFAULT_RECONNECT_BASE_DELAY_MS,
reconnectJitterRatio = DEFAULT_RECONNECT_JITTER_RATIO,
reconnectMaxDelayMs = DEFAULT_RECONNECT_MAX_DELAY_MS,
url = "/api/activity/stream",
}: ActivityStreamConnectionOptions): () => void {
let source: ActivityEventSourceLike | null = null;
let stopped = false;
let fallbackStarted = false;
let ready = false;
let connectionReady = false;
let lastEventId: string | null = null;
let reconnectAttempt = 0;
let reconnectTimer: ReturnType<typeof setTimeout> | null = null;
let watchdog: ReturnType<typeof setTimeout> | null = null;

const clearWatchdog = () => {
Expand All @@ -40,29 +52,38 @@ export function connectActivityStream({
watchdog = null;
};

const armWatchdog = (timeoutMs: number) => {
const clearReconnectTimer = () => {
if (reconnectTimer === null) return;
clearTimeout(reconnectTimer);
reconnectTimer = null;
};

const armWatchdog = (timeoutMs: number, onTimeout: () => void) => {
clearWatchdog();
watchdog = setTimeout(startFallback, timeoutMs);
watchdog = setTimeout(onTimeout, timeoutMs);
};

const markAlive = () => {
if (stopped || fallbackStarted) return;
armWatchdog(heartbeatTimeoutMs);
armWatchdog(heartbeatTimeoutMs, scheduleReconnect);
};

const onActivity = ((rawEvent: Event) => {
markAlive();
const event = parseActivityEvent(rawEvent);
if (event) onEvent(event);
if (!event) return;
lastEventId = event.id;
onEvent(event);
}) as EventListener;

const onHeartbeat = (() => {
markAlive();
}) as EventListener;

const onReadyEvent = (() => {
const firstReady = !ready;
ready = true;
const firstReady = !connectionReady;
connectionReady = true;
reconnectAttempt = 0;
markAlive();
if (firstReady) onReady();
}) as EventListener;
Expand All @@ -78,30 +99,69 @@ export function connectActivityStream({
source = null;
};

const streamUrl = () => {
if (!lastEventId) return url;
const separator = url.includes("?") ? "&" : "?";
return `${url}${separator}lastEventId=${encodeURIComponent(lastEventId)}`;
};

const startFallback = () => {
if (stopped || fallbackStarted) return;
fallbackStarted = true;
clearReconnectTimer();
detach();
onFallback();
};

const reconnectDelay = () => {
const exponentialDelay = Math.min(
reconnectMaxDelayMs,
reconnectBaseDelayMs * 2 ** reconnectAttempt
);
const jitter = exponentialDelay * reconnectJitterRatio * Math.random();
reconnectAttempt += 1;
return Math.round(exponentialDelay + jitter);
};

const scheduleReconnect = () => {
if (stopped || fallbackStarted) return;
if (!connectionReady) {
startFallback();
return;
}

detach();
clearReconnectTimer();
reconnectTimer = setTimeout(connect, reconnectDelay());
};

const onError = (() => {
startFallback();
scheduleReconnect();
}) as EventListener;

try {
source = createEventSource(url);
const connect = () => {
if (stopped || fallbackStarted) return;
connectionReady = false;

try {
source = createEventSource(streamUrl());
} catch {
startFallback();
return;
}

source.addEventListener("activity", onActivity);
source.addEventListener("error", onError);
source.addEventListener("heartbeat", onHeartbeat);
source.addEventListener("ready", onReadyEvent);
armWatchdog(connectionTimeoutMs);
} catch {
startFallback();
}
armWatchdog(connectionTimeoutMs, startFallback);
};

connect();

return () => {
stopped = true;
clearReconnectTimer();
detach();
};
}
Expand Down
60 changes: 43 additions & 17 deletions apps/dashboard/test/activity-stream.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -75,32 +75,50 @@ describe("activity SSE delivery", () => {
assert.equal(getActivitySubscriberCount(), initialSubscribers);
});

test("client connector accepts activity and falls back exactly once on stream error", () => {
const source = new FakeEventSource();
test("client connector accepts activity and reconnects with backoff on stream error", async () => {
const sources: FakeEventSource[] = [];
const received: string[] = [];
let fallbackCount = 0;
let readyCount = 0;
const event = makeActivityEvent({ id: "evt_client_stream_001" });

const disconnect = connectActivityStream({
createEventSource: () => source,
createEventSource: (url) => {
const source = new FakeEventSource();
sources.push(source);
if (sources.length === 2) {
assert.match(url, /lastEventId=evt_client_stream_001/);
}
return source;
},
onEvent: (activity) => received.push(activity.id),
onFallback: () => {
fallbackCount += 1;
},
onReady: () => {
readyCount += 1;
},
reconnectBaseDelayMs: 10,
reconnectJitterRatio: 0,
reconnectMaxDelayMs: 10,
});

source.emit("ready", "{}");
source.emit("activity", JSON.stringify(event));
source.emit("activity", "not-json");
source.emit("error");
source.emit("error");
sources[0].emit("ready", "{}");
sources[0].emit("activity", JSON.stringify(event));
sources[0].emit("activity", "not-json");
sources[0].emit("error");
sources[0].emit("error");

await delay(30);
sources[1].emit("ready", "{}");

assert.deepEqual(received, [event.id]);
assert.equal(fallbackCount, 1);
assert.equal(source.closeCount, 1);
assert.equal(fallbackCount, 0);
assert.equal(readyCount, 2);
assert.equal(sources[0].closeCount, 1);

disconnect();
assert.equal(source.closeCount, 1);
assert.equal(sources[1].closeCount, 1);
});

test("ready handshake reconciles the REST snapshot after subscription", () => {
Expand Down Expand Up @@ -146,24 +164,32 @@ describe("activity SSE delivery", () => {
assert.equal(source.closeCount, 1);
});

test("client falls back when a ready stream stops sending heartbeats", async () => {
const source = new FakeEventSource();
test("client reconnects when a ready stream stops sending heartbeats", async () => {
const sources: FakeEventSource[] = [];
let fallbackCount = 0;

connectActivityStream({
connectionTimeoutMs: 100,
createEventSource: () => source,
createEventSource: () => {
const source = new FakeEventSource();
sources.push(source);
return source;
},
heartbeatTimeoutMs: 10,
onEvent: () => {},
onFallback: () => {
fallbackCount += 1;
},
reconnectBaseDelayMs: 10,
reconnectJitterRatio: 0,
reconnectMaxDelayMs: 10,
});

source.emit("ready", "{}");
sources[0].emit("ready", "{}");
await delay(30);
assert.equal(fallbackCount, 1);
assert.equal(source.closeCount, 1);
assert.equal(fallbackCount, 0);
assert.equal(sources[0].closeCount, 1);
assert.equal(sources.length, 2);
});

test("server disconnects a stream whose bounded output queue fills", async () => {
Expand Down
Loading