CodeIssuesDiscussionsWikiPull RequestsProjectsCommitsActionsReleasesContributorsPulse● GatesSecuritySettingsDeploymentsPipelineInsightsAgents✨ Explain✨ Ask AI✨ Workspace✨ Spec✨ Tests▓ Debt Map✨ NL Search🏛 Archaeology
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 | /**
* Server-Sent Events (SSE) infrastructure.
*
* Manages live connections so gate runs, notifications, and CI updates
* stream to the browser in real time instead of requiring page refreshes.
*/
type SSEChannel = "gate" | "notification" | "pr" | "ci";
interface SSEClient {
id: string;
controller: ReadableStreamDefaultController;
userId?: string;
channels: Set<string>; // e.g. "gate:repoId", "pr:prId", "notification:userId"
connectedAt: number;
}
const clients = new Map<string, SSEClient>();
let clientIdCounter = 0;
function nextClientId(): string {
return `sse_${++clientIdCounter}_${Date.now().toString(36)}`;
}
/**
* Create an SSE Response for a Hono handler.
* The caller subscribes to channels, and this function returns a streaming Response.
*/
export function createSSEStream(
channels: string[],
userId?: string
): Response {
const clientId = nextClientId();
const stream = new ReadableStream({
start(controller) {
const client: SSEClient = {
id: clientId,
controller,
userId,
channels: new Set(channels),
connectedAt: Date.now(),
};
clients.set(clientId, client);
// Send initial connection event
const data = `event: connected\ndata: ${JSON.stringify({ clientId, channels })}\n\n`;
controller.enqueue(new TextEncoder().encode(data));
},
cancel() {
clients.delete(clientId);
},
});
return new Response(stream, {
headers: {
"Content-Type": "text/event-stream",
"Cache-Control": "no-cache",
Connection: "keep-alive",
"X-Accel-Buffering": "no",
},
});
}
/**
* Broadcast an event to all clients subscribed to a channel.
*/
export function broadcast(
channel: string,
event: string,
data: unknown
): number {
const encoded = new TextEncoder();
const payload = `event: ${event}\ndata: ${JSON.stringify(data)}\n\n`;
const bytes = encoded.encode(payload);
let sent = 0;
for (const [id, client] of clients) {
if (client.channels.has(channel)) {
try {
client.controller.enqueue(bytes);
sent++;
} catch {
clients.delete(id);
}
}
}
return sent;
}
/**
* Send an event to a specific user (by userId) across all their connections.
*/
export function sendToUser(
userId: string,
event: string,
data: unknown
): number {
const encoded = new TextEncoder();
const payload = `event: ${event}\ndata: ${JSON.stringify(data)}\n\n`;
const bytes = encoded.encode(payload);
let sent = 0;
for (const [id, client] of clients) {
if (client.userId === userId) {
try {
client.controller.enqueue(bytes);
sent++;
} catch {
clients.delete(id);
}
}
}
return sent;
}
/**
* Get active connection count (for monitoring).
*/
export function getActiveConnections(): number {
return clients.size;
}
/**
* Clean up stale connections (call periodically).
*/
export function cleanupStaleConnections(maxAgeMs = 30 * 60 * 1000): number {
const cutoff = Date.now() - maxAgeMs;
let removed = 0;
for (const [id, client] of clients) {
if (client.connectedAt < cutoff) {
try {
client.controller.close();
} catch {}
clients.delete(id);
removed++;
}
}
return removed;
}
// Periodic cleanup every 5 minutes
setInterval(() => cleanupStaleConnections(), 5 * 60 * 1000);
|