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 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 | /**
* SSE endpoint: `GET /live-events/:topic`.
*
* Topic format: `repo:{repoId}`, `pr:{prId}`, `user:{userId}`. The regex
* `^[a-z]+:[a-zA-Z0-9\-]+$` is enforced; anything else is a 400.
*
* Auth / authorization:
* - Runs behind softAuth so we have the viewer (or null).
* - For `repo:{repoId}` topics, we do a cheap DB check that the viewer has
* read access via `resolveRepoAccess`. `pr:` and `user:` topics currently
* only require a valid topic string — when we add PR-level privacy we'll
* extend this handler in place.
*
* Transport:
* - `text/event-stream` with keep-alive + nginx-friendly `X-Accel-Buffering`.
* - We write `id:` / `event:` / `data:` blocks per SSEEvent and send a
* `: ping` comment every 25s to keep intermediaries from timing out.
* - On stream close we unsubscribe and clear the heartbeat timer.
*/
import { Hono } from "hono";
import { eq } from "drizzle-orm";
import { softAuth } from "../middleware/auth";
import type { AuthEnv } from "../middleware/auth";
import { db } from "../db";
import { repositories } from "../db/schema";
import { resolveRepoAccess, satisfiesAccess } from "../middleware/repo-access";
import { subscribe, type SSEEvent } from "../lib/sse";
const app = new Hono<AuthEnv>();
/**
* Topic shape — `kind:id(:segment)*`. The first colon separates the kind
* (lowercase, used for the read-gate dispatch) from the id; subsequent
* colon-segments are scoping suffixes the publisher chose, e.g.
* `repo:<uuid>:issue:7`. Each segment is alphanumerics + dash so the
* URL path stays predictable.
*/
const TOPIC_RE = /^[a-z]+:[a-zA-Z0-9\-]+(?::[a-zA-Z0-9\-]+)*$/;
const HEARTBEAT_MS = 25_000;
app.get("/live-events/:topic", softAuth, async (c) => {
const topic = c.req.param("topic");
if (!topic || !TOPIC_RE.test(topic)) {
return c.json({ error: "Invalid topic" }, 400);
}
const user = c.get("user") ?? null;
// Topic is `kind:primaryId(:scope)*`. Slice on the first two colons so a
// multi-segment topic like `repo:<uuid>:issue:7` resolves to
// kind = "repo", primaryId = "<uuid>"
// and the trailing `:issue:7` is treated as scoping that the publisher
// chose (the broadcaster is keyed on the full topic string, so the suffix
// is preserved across publish/subscribe).
const firstColon = topic.indexOf(":");
const secondColon = topic.indexOf(":", firstColon + 1);
const kind = topic.slice(0, firstColon);
const primaryId =
secondColon === -1
? topic.slice(firstColon + 1)
: topic.slice(firstColon + 1, secondColon);
// For repo topics, gate on read access. Other topic kinds pass through.
if (kind === "repo") {
try {
const [repo] = await db
.select({ id: repositories.id, isPrivate: repositories.isPrivate })
.from(repositories)
.where(eq(repositories.id, primaryId))
.limit(1);
if (!repo) {
return c.json({ error: "Not found" }, 404);
}
const access = await resolveRepoAccess({
repoId: repo.id,
userId: user?.id ?? null,
isPublic: !repo.isPrivate,
});
if (!satisfiesAccess(access, "read")) {
return c.json({ error: "Forbidden" }, 403);
}
} catch {
return c.json({ error: "Not found" }, 404);
}
}
const encoder = new TextEncoder();
const stream = new ReadableStream<Uint8Array>({
start(controller) {
let closed = false;
const safeEnqueue = (chunk: string) => {
if (closed) return;
try {
controller.enqueue(encoder.encode(chunk));
} catch {
// Controller already closed — mark local state so we stop trying.
closed = true;
}
};
// Initial comment flushes headers on some proxies.
safeEnqueue(": open\n\n");
const unsubscribe = subscribe(topic, (event: SSEEvent) => {
let payload = "";
if (event.id !== undefined) payload += `id: ${event.id}\n`;
if (event.event !== undefined) payload += `event: ${event.event}\n`;
const data =
typeof event.data === "string"
? event.data
: JSON.stringify(event.data);
// SSE `data:` lines must not contain raw newlines — split if present.
for (const line of data.split("\n")) {
payload += `data: ${line}\n`;
}
payload += "\n";
safeEnqueue(payload);
});
const heartbeat = setInterval(() => {
safeEnqueue(": ping\n\n");
}, HEARTBEAT_MS);
const cleanup = () => {
if (closed) return;
closed = true;
clearInterval(heartbeat);
unsubscribe();
try {
controller.close();
} catch {
// Already closed — nothing to do.
}
};
// Client-side abort (navigation, tab close) surfaces via the request's
// AbortSignal. Bun's fetch-style request exposes this on `c.req.raw`.
const signal = c.req.raw.signal;
if (signal) {
if (signal.aborted) {
cleanup();
} else {
signal.addEventListener("abort", cleanup, { once: true });
}
}
},
});
return new Response(stream, {
status: 200,
headers: {
"Content-Type": "text/event-stream; charset=utf-8",
"Cache-Control": "no-cache, no-transform",
Connection: "keep-alive",
"X-Accel-Buffering": "no",
},
});
});
export default app;
|