mirror of
https://github.com/Ark0N/Codeman.git
synced 2026-09-30 20:49:41 +02:00
feat(multiuser): phase 4, event fan-out + stream scoping
Scopes real-time streams and the init snapshot so a multi-user client only receives what it owns. No-op in single-user mode (identity-less clients). - WS terminal (ws-routes): owner gate after the session lookup. A non-admin may only attach to their own session (close 4003); the global auth hook already ran on the upgrade and decorated req.authUser, so an unauthenticated upgrade never reaches the handler. - SSE (sse-stream-manager): per-client identity stored at addClient; broadcast() and the terminal-batch flush both enforce a routing hint via canDeliver(). WebServer.broadcast auto-derives the hint (deriveSseHint): session-scoped event families resolve the owner from the payload's session id (fail closed when the owner can't be resolved), machine-level families (docker/tunnel/update/system/ cron) + host-plan telemetry are admin-only, everything else stays global. Raw terminal bytes resolve the owner once and are withheld from non-owners. - getLightState is filtered per connection AFTER the shared cache (sessions, respawnStatus, subagents, workflowRuns by owner; scheduledRuns + planUsage admin-only); applied to both the SSE init snapshot and GET /api/status. - file-routes: getKnownSessionWorkingDir + getSessionAttachmentHistory (the preview/thumbnail/history helpers that bypass findSessionOrFail) now owner-check the session, closing a cross-user file-read path. - GET /api/search: harvestSources is owner-scoped. Deferred to a follow-up (documented in docs/multi-user-plan.md): away-digest + subagent/workflow REST list scoping, push-subscription identity + routing, per-user screenshot subdirs. The live-event versions of these are already routed by the SSE hint; only the on-demand REST aggregates remain global for admins-only follow-up. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
@@ -18,7 +18,7 @@ export interface ConfigPort {
|
||||
getClaudeModeConfig(): Promise<{ claudeMode?: ClaudeMode; allowedTools?: string }>;
|
||||
getTerminalHistoryConfig(): Promise<TerminalHistoryConfig>;
|
||||
getDefaultClaudeMdPath(): Promise<string | undefined>;
|
||||
getLightState(): unknown;
|
||||
getLightState(identity?: { username: string; role: 'admin' | 'user' }): unknown;
|
||||
getLightSessionsState(): unknown[];
|
||||
startTranscriptWatcher(sessionId: string, transcriptPath: string): void;
|
||||
stopTranscriptWatcher(sessionId: string): void;
|
||||
|
||||
@@ -22,7 +22,8 @@ import { generateFirstPageThumbnail } from '../../document-thumbnailer.js';
|
||||
import { getOfficePreviewPdfPath, getPreviewPdfDownloadName } from '../../document-preview-cache.js';
|
||||
import { sanitizeAttachmentHistoryItem } from '../../session-attachment-history.js';
|
||||
import { isBlockedAttachmentPath, loadAttachmentGuardConfig } from '../../config/attachment-guard.js';
|
||||
import { findSessionOrFail, validateSessionFilePath } from '../route-helpers.js';
|
||||
import { canAccessOwned, findSessionOrFail, getAuthUser, validateSessionFilePath } from '../route-helpers.js';
|
||||
import type { FastifyRequest } from 'fastify';
|
||||
import type { SessionAttachmentHistoryItem, SessionState } from '../../types/session.js';
|
||||
import { isSensitivePath } from '../sensitive-path.js';
|
||||
import { SseEvent } from '../sse-events.js';
|
||||
@@ -227,13 +228,17 @@ async function serveThumbnail(reply: FastifyReply, resolvedPath: string, extensi
|
||||
function getKnownSessionWorkingDir(
|
||||
ctx: SessionPort & ConfigPort,
|
||||
sessionId: string,
|
||||
reply: FastifyReply
|
||||
reply: FastifyReply,
|
||||
req: FastifyRequest
|
||||
): string | undefined {
|
||||
// Multi-user: a non-admin may only reach their OWN session's files. A foreign
|
||||
// (or missing) session is reported identically as 404 so existence isn't leaked.
|
||||
const user = getAuthUser(req);
|
||||
const liveSession = ctx.sessions.get(sessionId);
|
||||
if (liveSession) return liveSession.workingDir;
|
||||
if (liveSession && canAccessOwned(user, liveSession.owner)) return liveSession.workingDir;
|
||||
|
||||
const stored = ctx.store.getSession(sessionId);
|
||||
if (stored) return stored.workingDir;
|
||||
if (stored && canAccessOwned(user, (stored as { owner?: string }).owner)) return stored.workingDir;
|
||||
|
||||
reply.code(404).send(createErrorResponse(ApiErrorCode.NOT_FOUND, `Session ${sessionId} not found`));
|
||||
return undefined;
|
||||
@@ -261,10 +266,13 @@ function appendDownloadFlag(url: string): string {
|
||||
|
||||
function getSessionAttachmentHistory(
|
||||
ctx: SessionPort & ConfigPort,
|
||||
sessionId: string
|
||||
sessionId: string,
|
||||
req: FastifyRequest
|
||||
): { workingDir: string; history: SessionAttachmentHistoryItem[] } | undefined {
|
||||
const user = getAuthUser(req);
|
||||
const liveSession = ctx.sessions.get(sessionId);
|
||||
if (liveSession) {
|
||||
if (!canAccessOwned(user, liveSession.owner)) return undefined;
|
||||
return {
|
||||
workingDir: liveSession.workingDir,
|
||||
history: liveSession.getAttachmentHistoryForPersist() ?? liveSession.attachmentHistory ?? [],
|
||||
@@ -272,7 +280,7 @@ function getSessionAttachmentHistory(
|
||||
}
|
||||
|
||||
const stored = ctx.store.getSession(sessionId) as StoredSessionWithPrivateAttachmentHistory | undefined;
|
||||
if (!stored) return undefined;
|
||||
if (!stored || !canAccessOwned(user, (stored as { owner?: string }).owner)) return undefined;
|
||||
|
||||
return {
|
||||
workingDir: stored.workingDir,
|
||||
@@ -766,7 +774,7 @@ export function registerFileRoutes(app: FastifyInstance, ctx: SessionPort & Even
|
||||
// each entry to current metadata + routes. External entries are re-registered.
|
||||
app.get('/api/sessions/:id/attachments', async (req, reply) => {
|
||||
const { id } = req.params as { id: string };
|
||||
const sessionHistory = getSessionAttachmentHistory(ctx, id);
|
||||
const sessionHistory = getSessionAttachmentHistory(ctx, id, req);
|
||||
if (!sessionHistory) {
|
||||
reply.code(404).send(createErrorResponse(ApiErrorCode.NOT_FOUND, `Session ${id} not found`));
|
||||
return;
|
||||
@@ -794,7 +802,7 @@ export function registerFileRoutes(app: FastifyInstance, ctx: SessionPort & Even
|
||||
// size/mtime as the underlying file is rewritten).
|
||||
app.get('/api/sessions/:id/attachments/:attachmentId', async (req, reply) => {
|
||||
const { id, attachmentId } = req.params as { id: string; attachmentId: string };
|
||||
const workingDir = getKnownSessionWorkingDir(ctx, id, reply);
|
||||
const workingDir = getKnownSessionWorkingDir(ctx, id, reply, req);
|
||||
if (!workingDir) return;
|
||||
const record = getAttachmentOr404(reply, id, attachmentId);
|
||||
if (!record) return;
|
||||
@@ -850,7 +858,7 @@ export function registerFileRoutes(app: FastifyInstance, ctx: SessionPort & Even
|
||||
// convert server-side; PDF/PNG/text redirect to the raw route.
|
||||
app.get('/api/sessions/:id/attachments/:attachmentId/preview', async (req, reply) => {
|
||||
const { id, attachmentId } = req.params as { id: string; attachmentId: string };
|
||||
const workingDir = getKnownSessionWorkingDir(ctx, id, reply);
|
||||
const workingDir = getKnownSessionWorkingDir(ctx, id, reply, req);
|
||||
if (!workingDir) return;
|
||||
const record = getAttachmentOr404(reply, id, attachmentId);
|
||||
if (!record) return;
|
||||
@@ -870,7 +878,7 @@ export function registerFileRoutes(app: FastifyInstance, ctx: SessionPort & Even
|
||||
// Serve a first-page thumbnail of a registered attachment by id.
|
||||
app.get('/api/sessions/:id/attachments/:attachmentId/thumbnail', async (req, reply) => {
|
||||
const { id, attachmentId } = req.params as { id: string; attachmentId: string };
|
||||
const workingDir = getKnownSessionWorkingDir(ctx, id, reply);
|
||||
const workingDir = getKnownSessionWorkingDir(ctx, id, reply, req);
|
||||
if (!workingDir) return;
|
||||
const record = getAttachmentOr404(reply, id, attachmentId);
|
||||
if (!record) return;
|
||||
@@ -884,7 +892,7 @@ export function registerFileRoutes(app: FastifyInstance, ctx: SessionPort & Even
|
||||
app.get('/api/sessions/:id/file-preview', async (req, reply) => {
|
||||
const { id } = req.params as { id: string };
|
||||
const { path: filePath } = req.query as { path?: string };
|
||||
const workingDir = getKnownSessionWorkingDir(ctx, id, reply);
|
||||
const workingDir = getKnownSessionWorkingDir(ctx, id, reply, req);
|
||||
if (!workingDir) return;
|
||||
|
||||
if (!filePath) {
|
||||
@@ -912,7 +920,7 @@ export function registerFileRoutes(app: FastifyInstance, ctx: SessionPort & Even
|
||||
app.get('/api/sessions/:id/file-thumbnail', async (req, reply) => {
|
||||
const { id } = req.params as { id: string };
|
||||
const { path: filePath } = req.query as { path?: string };
|
||||
const workingDir = getKnownSessionWorkingDir(ctx, id, reply);
|
||||
const workingDir = getKnownSessionWorkingDir(ctx, id, reply, req);
|
||||
if (!workingDir) return;
|
||||
|
||||
if (!filePath) {
|
||||
|
||||
@@ -23,7 +23,7 @@
|
||||
*/
|
||||
|
||||
import { FastifyInstance } from 'fastify';
|
||||
import { parseBody } from '../route-helpers.js';
|
||||
import { canAccessOwned, getAuthUser, parseBody } from '../route-helpers.js';
|
||||
import { SearchQuerySchema } from '../schemas.js';
|
||||
import {
|
||||
searchSources,
|
||||
@@ -62,13 +62,14 @@ interface SessionLike {
|
||||
* Harvest the three source arrays from the live in-memory stores. Reads only
|
||||
* bounded, already-loaded data — no disk I/O, no terminal buffers.
|
||||
*/
|
||||
function harvestSources(ctx: SessionPort & InfraPort): SearchSources {
|
||||
function harvestSources(ctx: SessionPort & InfraPort, canSee?: (owner?: string) => boolean): SearchSources {
|
||||
const sessions: SessionSearchInput[] = [];
|
||||
const events: EventSearchInput[] = [];
|
||||
const files: FileSearchInput[] = [];
|
||||
|
||||
for (const raw of ctx.sessions.values()) {
|
||||
const s = raw as unknown as SessionLike;
|
||||
const s = raw as unknown as SessionLike & { owner?: string };
|
||||
if (canSee && !canSee(s.owner)) continue; // multi-user ownership scope
|
||||
const sessionName = s.name ?? '';
|
||||
const timestamp = s.lastActivityAt ?? s.createdAt ?? 0;
|
||||
|
||||
@@ -96,7 +97,8 @@ function harvestSources(ctx: SessionPort & InfraPort): SearchSources {
|
||||
|
||||
// Events: from the live run-summary trackers, keyed by session id.
|
||||
for (const [sessionId, tracker] of ctx.runSummaryTrackers) {
|
||||
const session = ctx.sessions.get(sessionId) as unknown as SessionLike | undefined;
|
||||
const session = ctx.sessions.get(sessionId) as unknown as (SessionLike & { owner?: string }) | undefined;
|
||||
if (canSee && !canSee(session?.owner)) continue; // multi-user ownership scope
|
||||
const sessionName = session?.name ?? '';
|
||||
const summary = tracker.getSummary();
|
||||
// Newest events are most relevant; cap the per-session harvest.
|
||||
@@ -120,6 +122,8 @@ export function registerSearchRoutes(app: FastifyInstance, ctx: SessionPort & In
|
||||
app.get('/api/search', async (req) => {
|
||||
// Zod-validate the query. parseBody throws a structured 400 on failure.
|
||||
const { q, types, limit } = parseBody(SearchQuerySchema, req.query);
|
||||
const user = getAuthUser(req);
|
||||
const canSee = (owner?: string) => canAccessOwned(user, owner);
|
||||
|
||||
const allowed: Set<SearchSourceType> | null = types
|
||||
? new Set(
|
||||
@@ -130,7 +134,7 @@ export function registerSearchRoutes(app: FastifyInstance, ctx: SessionPort & In
|
||||
)
|
||||
: null;
|
||||
|
||||
const sources = harvestSources(ctx);
|
||||
const sources = harvestSources(ctx, canSee);
|
||||
|
||||
// Apply the optional source-type filter before searching so excluded
|
||||
// sources never contribute to (or consume budget in) the result set.
|
||||
|
||||
@@ -139,7 +139,7 @@ export function registerSystemRoutes(
|
||||
|
||||
// ========== Status ==========
|
||||
|
||||
app.get('/api/status', async () => ctx.getLightState());
|
||||
app.get('/api/status', async (req) => ctx.getLightState(req.authUser));
|
||||
|
||||
// ========== Tunnel ==========
|
||||
|
||||
|
||||
@@ -35,6 +35,7 @@ import type { SessionPort } from '../ports/session-port.js';
|
||||
import { MAX_INPUT_LENGTH } from '../../config/terminal-limits.js';
|
||||
import { isAllowedRequestHost, isAllowedRequestOrigin, type HostPolicy } from '../network-auth-policy.js';
|
||||
import { WsConnectionRegistry } from '../ws-connection-registry.js';
|
||||
import { canAccessOwned, getAuthUser } from '../route-helpers.js';
|
||||
|
||||
/** Micro-batch interval for terminal output (ms). Short enough for low latency,
|
||||
* long enough to group Ink's rapid cursor-up redraw sequences into single frames. */
|
||||
@@ -93,6 +94,17 @@ export function registerWsRoutes(app: FastifyInstance, ctx: SessionPort, getHost
|
||||
return;
|
||||
}
|
||||
|
||||
// Multi-user owner gate: writing to this socket injects keystrokes into the
|
||||
// agent, so a non-admin may only attach to their OWN session. The global auth
|
||||
// hook already ran on the upgrade request and decorated req.authUser (an
|
||||
// unauthenticated upgrade never reaches here — the hook 401s the handshake).
|
||||
// findSessionOrFail throws an HTTP-shaped error, so the check is inlined here
|
||||
// as a 4003 close. No-op in single-user mode (canAccessOwned returns true).
|
||||
if (!canAccessOwned(getAuthUser(req), session.owner)) {
|
||||
socket.close(4003, 'Forbidden');
|
||||
return;
|
||||
}
|
||||
|
||||
// Structured transport logging — surfaces WS open/close/timeout churn so the
|
||||
// tunnel-flap behavior (COD-134) is observable in the server logs. Fastify is
|
||||
// configured logger:false, so we log via console (→ journald under systemd).
|
||||
|
||||
+97
-4
@@ -334,6 +334,7 @@ export class WebServer extends EventEmitter {
|
||||
const session = this.sessions.get(sessionId);
|
||||
return session ? this.getSessionStateWithRespawn(session) : null;
|
||||
},
|
||||
resolveSessionOwner: (sessionId) => this.sessions.get(sessionId)?.owner,
|
||||
},
|
||||
this.cleanup
|
||||
);
|
||||
@@ -795,12 +796,12 @@ export class WebServer extends EventEmitter {
|
||||
// Track tunnel clients — cloudflared proxies locally so req.ip is always
|
||||
// 127.0.0.1; detect tunnel traffic via Cf-Connecting-Ip header instead.
|
||||
const isRemote = !!req.headers['cf-connecting-ip'];
|
||||
this.sse.addClient(reply, sessionFilter, isRemote, clientId);
|
||||
this.sse.addClient(reply, sessionFilter, isRemote, clientId, req.authUser);
|
||||
|
||||
// Send initial state
|
||||
// Use light state for SSE init to avoid sending 2MB+ terminal buffers
|
||||
// Buffers are fetched on-demand when switching tabs
|
||||
this.sse.sendSSE(reply, SseEvent.Init, this.getLightState());
|
||||
this.sse.sendSSE(reply, SseEvent.Init, this.getLightState(req.authUser));
|
||||
// Flush Cloudflare tunnel buffer with padding — ensures the init event
|
||||
// (and any immediately following events) are delivered without proxy delay.
|
||||
this.sse.sendPadding(reply);
|
||||
@@ -1751,7 +1752,49 @@ export class WebServer extends EventEmitter {
|
||||
* Get lightweight state for SSE init - excludes full terminal buffers
|
||||
* to prevent browser freezes. Terminal buffers are fetched on-demand.
|
||||
*/
|
||||
private getLightState() {
|
||||
private getLightState(identity?: import('../types/user.js').AuthUser) {
|
||||
const base = this.computeLightState();
|
||||
// Multi-user: filter the shared cached blob per connection identity (the plan's
|
||||
// "filter AFTER the cache" approach). No-op for admins / single-user.
|
||||
if (isMultiUserMode() && identity && identity.role !== 'admin') {
|
||||
return this.filterLightStateForUser(base, identity.username);
|
||||
}
|
||||
return base;
|
||||
}
|
||||
|
||||
/** Shallow-filter the light-state blob to what a non-admin user may see. */
|
||||
private filterLightStateForUser(base: Record<string, unknown>, username: string): Record<string, unknown> {
|
||||
const ownedIds = new Set<string>();
|
||||
const ownedClaudeIds = new Set<string>();
|
||||
for (const [id, s] of this.sessions) {
|
||||
if (s.owner === username) {
|
||||
ownedIds.add(id);
|
||||
if (s.claudeSessionId) ownedClaudeIds.add(s.claudeSessionId);
|
||||
}
|
||||
}
|
||||
const sessions = Array.isArray(base.sessions)
|
||||
? (base.sessions as Array<{ owner?: string }>).filter((s) => s.owner === username)
|
||||
: base.sessions;
|
||||
const respawnStatus: Record<string, unknown> = {};
|
||||
for (const [id, v] of Object.entries((base.respawnStatus as Record<string, unknown>) ?? {})) {
|
||||
if (ownedIds.has(id)) respawnStatus[id] = v;
|
||||
}
|
||||
const bySession = (arr: unknown, key: 'sessionId' | 'sessionUuid') =>
|
||||
Array.isArray(arr)
|
||||
? (arr as Array<Record<string, unknown>>).filter((x) => ownedClaudeIds.has(String(x[key])))
|
||||
: arr;
|
||||
return {
|
||||
...base,
|
||||
sessions,
|
||||
respawnStatus,
|
||||
scheduledRuns: [], // legacy ScheduledRun has no owner yet → admin-only
|
||||
subagents: bySession(base.subagents, 'sessionId'),
|
||||
workflowRuns: bySession(base.workflowRuns, 'sessionUuid'),
|
||||
planUsage: null, // host-plan telemetry is admin-only
|
||||
};
|
||||
}
|
||||
|
||||
private computeLightState() {
|
||||
const now = Date.now();
|
||||
if (this.cachedLightState && now - this.cachedLightState.timestamp < WebServer.LIGHT_STATE_CACHE_TTL_MS) {
|
||||
return this.cachedLightState.data;
|
||||
@@ -1796,7 +1839,57 @@ export class WebServer extends EventEmitter {
|
||||
this.cachedLightState = null;
|
||||
this.cachedSessionsList = null;
|
||||
}
|
||||
this.sse.broadcast(event, data);
|
||||
// Multi-user: derive an ownership routing hint so an event only reaches the
|
||||
// clients entitled to it (no-op in single-user — hint stays undefined).
|
||||
this.sse.broadcast(event, data, isMultiUserMode() ? this.deriveSseHint(event, data) : undefined);
|
||||
}
|
||||
|
||||
/**
|
||||
* Map an SSE event + payload to a routing hint (multi-user). Session-scoped
|
||||
* families resolve the owner from a sessionId in the payload (fail closed if it
|
||||
* can't be resolved); machine-level families are admin-only; host-plan telemetry
|
||||
* is admin-only; everything else stays global. Default is fail-closed for the
|
||||
* session-scoped prefixes so a missed field starves rather than leaks.
|
||||
*/
|
||||
private deriveSseHint(event: string, data: unknown): import('./sse-stream-manager.js').SseRoutingHint | undefined {
|
||||
// Machine-level / host-wide: admins only.
|
||||
if (
|
||||
event.startsWith('docker:') ||
|
||||
event.startsWith('tunnel:') ||
|
||||
event.startsWith('update:') ||
|
||||
event.startsWith('system:') ||
|
||||
event.startsWith('cron:') ||
|
||||
event === SseEvent.SessionStatusTelemetry
|
||||
) {
|
||||
return { adminOnly: true };
|
||||
}
|
||||
// Session-scoped families: resolve the owner from the payload's session id.
|
||||
const SESSION_PREFIXES = [
|
||||
'session:',
|
||||
'ralph:',
|
||||
'respawn:',
|
||||
'subagent:',
|
||||
'workflow:',
|
||||
'attachment:',
|
||||
'task:',
|
||||
'mux:',
|
||||
'transcript:',
|
||||
'plan:',
|
||||
'orchestrator:',
|
||||
'hook:',
|
||||
'image:',
|
||||
'scheduled:',
|
||||
'team:',
|
||||
'case:',
|
||||
];
|
||||
if (SESSION_PREFIXES.some((p) => event.startsWith(p))) {
|
||||
const d = (data ?? {}) as { sessionId?: string; id?: string; session?: { id?: string } };
|
||||
const sessionId = d.sessionId ?? d.id ?? d.session?.id;
|
||||
const owner = sessionId ? this.sessions.get(sessionId)?.owner : undefined;
|
||||
return { owner, sessionScoped: true };
|
||||
}
|
||||
// Unrecognized / genuinely global events (connection status, needsRefresh): all.
|
||||
return undefined;
|
||||
}
|
||||
|
||||
private batchTerminalData(sessionId: string, data: string): void {
|
||||
|
||||
@@ -17,6 +17,7 @@
|
||||
|
||||
import type { FastifyReply } from 'fastify';
|
||||
import type { BackgroundTask } from '../session.js';
|
||||
import type { AuthUser } from '../types.js';
|
||||
import { CleanupManager, StaleExpirationMap } from '../utils/index.js';
|
||||
import { SseEvent } from './sse-events.js';
|
||||
import {
|
||||
@@ -38,6 +39,26 @@ const SSE_PADDING = ':' + 'p'.repeat(SSE_PADDING_SIZE) + '\n';
|
||||
interface SseStreamManagerDeps {
|
||||
/** Get session state with respawn info for session:updated broadcasts */
|
||||
getSessionStateWithRespawn(sessionId: string): unknown;
|
||||
/** Resolve a session's owner (multi-user) for SSE routing; undefined = unknown. */
|
||||
resolveSessionOwner?(sessionId: string): string | undefined;
|
||||
}
|
||||
|
||||
/**
|
||||
* Optional per-broadcast routing hint (multi-user). Resolved by WebServer.broadcast
|
||||
* before delegation. When absent, an event is delivered to all clients (global).
|
||||
*/
|
||||
export interface SseRoutingHint {
|
||||
/** Deliver only to this session's owner (+ admins). */
|
||||
owner?: string;
|
||||
/** Deliver only to admins (machine-level events: docker builds, tunnel, update). */
|
||||
adminOnly?: boolean;
|
||||
/** Deliver only to this exact user (+ admins). */
|
||||
username?: string;
|
||||
/**
|
||||
* The event is session-scoped but the owner could not be resolved — non-admins
|
||||
* are starved (fail closed) rather than leaked to.
|
||||
*/
|
||||
sessionScoped?: boolean;
|
||||
}
|
||||
|
||||
export class SseStreamManager {
|
||||
@@ -50,6 +71,8 @@ export class SseStreamManager {
|
||||
private sseClients: Map<FastifyReply, Set<string> | null> = new Map();
|
||||
/** Optional client-supplied IDs → reply, for live filter updates without reconnecting */
|
||||
private sseClientsById: Map<string, FastifyReply> = new Map();
|
||||
/** Per-client identity (multi-user); absent for single-user clients → no filtering. */
|
||||
private sseClientIdentity: Map<FastifyReply, AuthUser> = new Map();
|
||||
/** SSE clients connecting from non-localhost (i.e. through tunnel) */
|
||||
private remoteSseClients: Set<FastifyReply> = new Set();
|
||||
/** Clients with backpressure — skip writes until 'drain' fires */
|
||||
@@ -105,8 +128,15 @@ export class SseStreamManager {
|
||||
this._isTunnelActive = active;
|
||||
}
|
||||
|
||||
addClient(reply: FastifyReply, sessionFilter: Set<string> | null, isRemote: boolean, clientId?: string): void {
|
||||
addClient(
|
||||
reply: FastifyReply,
|
||||
sessionFilter: Set<string> | null,
|
||||
isRemote: boolean,
|
||||
clientId?: string,
|
||||
identity?: AuthUser
|
||||
): void {
|
||||
this.sseClients.set(reply, sessionFilter);
|
||||
if (identity) this.sseClientIdentity.set(reply, identity);
|
||||
if (isRemote) {
|
||||
this.remoteSseClients.add(reply);
|
||||
}
|
||||
@@ -117,6 +147,7 @@ export class SseStreamManager {
|
||||
this.sseClients.delete(prev);
|
||||
this.remoteSseClients.delete(prev);
|
||||
this.backpressuredClients.delete(prev);
|
||||
this.sseClientIdentity.delete(prev);
|
||||
}
|
||||
this.sseClientsById.set(clientId, reply);
|
||||
}
|
||||
@@ -126,12 +157,31 @@ export class SseStreamManager {
|
||||
this.sseClients.delete(reply);
|
||||
this.remoteSseClients.delete(reply);
|
||||
this.backpressuredClients.delete(reply);
|
||||
this.sseClientIdentity.delete(reply);
|
||||
// Clear any clientId mappings pointing at this reply
|
||||
for (const [id, r] of this.sseClientsById) {
|
||||
if (r === reply) this.sseClientsById.delete(id);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Whether an SSE event carrying `hint` may be delivered to `reply`. Clients with
|
||||
* no identity (single-user) always receive everything. Admins receive everything.
|
||||
* A non-admin receives an event only when the hint targets them (owner/username)
|
||||
* or the event is unrouted/global; session-scoped events with an unresolved owner
|
||||
* are withheld (fail closed).
|
||||
*/
|
||||
private canDeliver(reply: FastifyReply, hint?: SseRoutingHint): boolean {
|
||||
const identity = this.sseClientIdentity.get(reply);
|
||||
if (!identity || identity.role === 'admin') return true;
|
||||
if (!hint) return true;
|
||||
if (hint.adminOnly) return false;
|
||||
if (hint.username !== undefined) return hint.username === identity.username;
|
||||
if (hint.owner !== undefined) return hint.owner === identity.username;
|
||||
if (hint.sessionScoped) return false; // session-scoped but owner unknown → fail closed
|
||||
return true;
|
||||
}
|
||||
|
||||
/**
|
||||
* Update an existing client's session subscription filter without forcing
|
||||
* an SSE reconnect. Returns true if the client was found and updated.
|
||||
@@ -197,7 +247,7 @@ export class SseStreamManager {
|
||||
|
||||
// ========== Broadcasting ==========
|
||||
|
||||
broadcast(event: string, data: unknown): void {
|
||||
broadcast(event: string, data: unknown, hint?: SseRoutingHint): void {
|
||||
// Skip serialization entirely when no clients are listening
|
||||
if (this.sseClients.size === 0) return;
|
||||
|
||||
@@ -224,6 +274,8 @@ export class SseStreamManager {
|
||||
// active session's terminal output. Terminal events bypass this method
|
||||
// entirely (see flushSessionTerminalBatch — it applies the filter).
|
||||
for (const [client] of this.sseClients) {
|
||||
// Multi-user ownership routing (no-op for identity-less single-user clients).
|
||||
if (!this.canDeliver(client, hint)) continue;
|
||||
this.sendSSEPreformatted(client, message);
|
||||
}
|
||||
}
|
||||
@@ -314,9 +366,15 @@ export class SseStreamManager {
|
||||
// terminal data is high-frequency and latency-sensitive.
|
||||
const padding = this._isTunnelActive ? SSE_PADDING : '';
|
||||
const message = `event: session:terminal\ndata: {"id":"${sessionId}","data":${escapedData}}\n\n` + padding;
|
||||
// Raw terminal bytes are the highest-value payload: resolve the session owner
|
||||
// ONCE and withhold the batch from any non-admin who is not the owner (fail
|
||||
// closed if the owner is unknown). No-op for identity-less single-user clients.
|
||||
const owner = this.deps.resolveSessionOwner?.(sessionId);
|
||||
const termHint: SseRoutingHint = { owner, sessionScoped: true };
|
||||
for (const [client, filter] of this.sseClients) {
|
||||
// Skip clients that have a session filter and aren't subscribed to this session
|
||||
if (filter && !filter.has(sessionId)) continue;
|
||||
if (!this.canDeliver(client, termHint)) continue;
|
||||
this.sendSSEPreformatted(client, message);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user