Pre-launch — Gluecron is in final validation. Public signups and git hosting for non-owner users open after launch review.
CodeIssuesDiscussionsWikiPull RequestsProjectsCommitsActionsReleasesContributorsPulse● GatesSecuritySettingsDeploymentsPipelineInsightsAgents✨ Explain✨ Ask AI✨ Workspace✨ Spec✨ Tests▓ Debt Map✨ NL Search🏛 Archaeology
claude/adoring-hopper-5x74bqclaude/affectionate-feynman-ykrf1hclaude/architecture-audit-design-wxprenclaude/build-status-update-3MXsfclaude/charming-meitner-mllb5rclaude/compare-gate-gluecron-s4mFQclaude/confident-faraday-tikcwbclaude/continue-work-XMTlIclaude/crontech-gluecron-deploy-7MIECclaude/crontech-platform-setup-SeKfwclaude/design-2026claude/ecstatic-ptolemy-jMdigclaude/enhance-github-integration-QNHdGclaude/fix-aa-loop-issue-PonMQclaude/fix-actions-and-processclaude/fix-desktop-errors-XqoW8claude/fix-red-workflowsclaude/fix-website-access-6FKJNclaude/gatetest-integration-hardeningclaude/github-audit-improvements-bDFr9claude/gluecron-launch-status-FoMRlclaude/hopeful-lamport-olfCTclaude/issue-to-pr-and-protectionsclaude/jolly-heisenberg-2sg1Qclaude/launch-preparation-QmTb6claude/new-session-xk1l7claude/plan-platform-architecture-kkN4yclaude/platform-analysis-roadmap-1nUGLclaude/platform-launch-assessment-8dWV8claude/polish-platform-release-AeDrUclaude/resume-previous-work-KzyLwclaude/review-crontech-handoff-qYEVqclaude/review-project-completeness-lHhS2claude/review-readme-docs-ulqPKclaude/serene-edison-rj87weclaude/setup-multi-repo-dev-BCwNQclaude/ship-fixes-and-tests-Jvz1cclaude/site-audit-competitive-pctlwgclaude/site-migration-vercel-XstpKclaude/standalone-product-repos-XHFTDcopilot/feat-smart-empty-states-keyboard-first-enhancementcopilot/feat-smart-morning-digest-review-context-restorecopilot/fix-and-process-workflowscopilot/update-ai-powered-code-reviewfeat/debt-mapfeat/push-policy-codeowners-hardeningfeat/smart-digest-contextfeat/stage-impactfeat/t1-secret-migrationfeat/u-polishfeat/w-self-hostfeat/w2-claude-configfix/agent-journey-orphan-sweepgatetest/auto-fix-1776586424172gatetest/auto-fix-1776586534814gatetest/auto-fix-1776590685143gatetest/auto-fix-1776590808199mainops/redeploy-retriggerstyle/dxt-cta-themeworktree-agent-a3377aad30d55da26worktree-agent-a7ef607b7ee1d6c74
sse.ts3.3 KB · 148 lines
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);