diff --git a/src/backend/clients/dynamodb/DDBClient.ts b/src/backend/clients/dynamodb/DDBClient.ts index c9010e5f8..7bee880df 100644 --- a/src/backend/clients/dynamodb/DDBClient.ts +++ b/src/backend/clients/dynamodb/DDBClient.ts @@ -511,6 +511,7 @@ export class DDBClient extends PuterClient { expression: string, expressionValues?: Record, expressionNames?: Record, + options?: { condition?: string }, ) { const hasValues = !!expressionValues && Object.keys(expressionValues).length > 0; @@ -531,6 +532,9 @@ export class DDBClient extends PuterClient { ...(hasNames ? { ExpressionAttributeNames: expressionNames } : {}), + ...(options?.condition + ? { ConditionExpression: options.condition } + : {}), ReturnValues: 'ALL_NEW', ReturnConsumedCapacity: 'TOTAL', }), diff --git a/src/backend/clients/event/types.ts b/src/backend/clients/event/types.ts index 4c2b9ebf7..31f2d3642 100644 --- a/src/backend/clients/event/types.ts +++ b/src/backend/clients/event/types.ts @@ -474,14 +474,30 @@ export type EventMap = { * * The dispatch hot path never reads the counter — it is the broadcast that * invalidates, which is what keeps an unsubscribed write at zero Redis - * commands. + * commands. Hence `outer.pubsub.*`: the cache is a per-process map, the + * webhook reaches one node per peer region, and only the Redis fan reaches + * that node's siblings — a node the bump never reached has nothing that + * expires. */ - 'outer.events.generationBumped': { + 'outer.pubsub.events.generationBumped': { userId: number; generation: number; /** Whether the table changed; only then does a peer need to re-read it. */ durable: boolean; }; + /** + * A user's sockets moved between regions, so every node must drop what it + * cached about where they are. Bumped on a first connect, a last + * disconnect, and every repair — never on a timer, which is what keeps + * presence reads proportional to session churn rather than to event volume. + * `outer.pubsub.*` for the same reason as the generation bump: a + * per-process cache that only the Redis fan reaches on a region's other + * nodes. + */ + 'outer.pubsub.events.presenceBumped': { + userId: number; + generation: number; + }; 'outer.fs.write-hash': { hash: string; uuid: string }; /** * Cache keys the KV read cache must stop serving, because the entries diff --git a/src/backend/controllers/broadcast/BroadcastController.ts b/src/backend/controllers/broadcast/BroadcastController.ts index c4a100007..7b15b1a5c 100644 --- a/src/backend/controllers/broadcast/BroadcastController.ts +++ b/src/backend/controllers/broadcast/BroadcastController.ts @@ -20,6 +20,7 @@ import type { Request, Response } from 'express'; import { Controller, Post } from '../../core/http/decorators.js'; import type { BroadcastService } from '../../services/broadcast/BroadcastService.js'; +import type { ForwardBatch } from '../../services/events/forwardQueue.js'; import { PuterController } from '../types.js'; /** @@ -55,18 +56,11 @@ export class BroadcastController extends PuterController { return; } - const headerOnce = (name: string): string | undefined => { - const value = req.headers[name]; - if (Array.isArray(value)) return value[0]; - return value; - }; - - const result = await broadcast.verifyAndEmit(req.rawBody, req.body, { - peerId: headerOnce('x-broadcast-peer-id'), - timestamp: headerOnce('x-broadcast-timestamp'), - nonce: headerOnce('x-broadcast-nonce'), - signature: headerOnce('x-broadcast-signature'), - }); + const result = await broadcast.verifyAndEmit( + req.rawBody, + req.body, + broadcastHeaders(req), + ); if (result.ok) { res.status(200).json({ ok: true, ...(result.info ?? {}) }); @@ -76,4 +70,57 @@ export class BroadcastController extends PuterController { error: { message: result.message ?? 'Bad request' }, }); } + + /** + * Event deliveries one region addressed at this one, on the same signed + * channel as the webhook above but carrying socket traffic rather than bus + * events. The answer names the pairs this region holds no socket for, which + * is what lets the sender correct its presence row. + */ + @Post('/events', { subdomain: '*' }) + async events(req: Request, res: Response): Promise { + const broadcast = this.services.broadcast as unknown as + | BroadcastService + | undefined; + if (!broadcast) { + res.status(503).json({ + error: { message: 'Broadcast service not registered' }, + }); + return; + } + + const verified = await broadcast.verifySignedRequest( + req.rawBody, + broadcastHeaders(req), + ); + if (!verified.ok) { + res.status(verified.status ?? 400).json({ + error: { message: verified.message ?? 'Bad request' }, + }); + return; + } + if (verified.info?.ignored) { + res.status(200).json({ ok: true, ...verified.info }); + return; + } + + const reply = await this.services.eventForward.receive( + (req.body ?? {}) as ForwardBatch, + ); + res.status(200).json({ ok: true, ...reply }); + } } + +const broadcastHeaders = (req: Request) => { + const headerOnce = (name: string): string | undefined => { + const value = req.headers[name]; + if (Array.isArray(value)) return value[0]; + return value; + }; + return { + peerId: headerOnce('x-broadcast-peer-id'), + timestamp: headerOnce('x-broadcast-timestamp'), + nonce: headerOnce('x-broadcast-nonce'), + signature: headerOnce('x-broadcast-signature'), + }; +}; diff --git a/src/backend/services/broadcast/BroadcastService.ts b/src/backend/services/broadcast/BroadcastService.ts index 30b060dcd..481572c1f 100644 --- a/src/backend/services/broadcast/BroadcastService.ts +++ b/src/backend/services/broadcast/BroadcastService.ts @@ -55,6 +55,17 @@ interface IncomingHeaders { signature: string | undefined; } +// -- Addressed sends ------------------------------------------------- + +/** Connections one peer may hold open, and how long an idle one survives. */ +const ADDRESSED_MAX_SOCKETS_PER_PEER = 16; +const ADDRESSED_MAX_FREE_SOCKETS_PER_PEER = 4; +const ADDRESSED_SOCKET_IDLE_MS = 30_000; +const ADDRESSED_REQUEST_TIMEOUT_MS = 5_000; + +/** What a deployment with no configured identity calls itself. */ +const DEFAULT_REGION_ID = 'local'; + // -- Service --------------------------------------------------------- /** @@ -104,6 +115,18 @@ export class BroadcastService extends PuterService { #webhookHostHeader: string | null = null; /** Self-signed certs are common between Puter nodes — accept them. */ #webhookHttpsAgent = new HttpsAgent({ rejectUnauthorized: false }); + /** + * The agent addressed sends use. Separate because they run at event rates, + * where a fresh handshake per request outweighs the payload; the pool is + * bounded so a peer that stops answering cannot accumulate connections. + */ + #addressedHttpsAgent = new HttpsAgent({ + rejectUnauthorized: false, + keepAlive: true, + maxSockets: ADDRESSED_MAX_SOCKETS_PER_PEER, + maxFreeSockets: ADDRESSED_MAX_FREE_SOCKETS_PER_PEER, + timeout: ADDRESSED_SOCKET_IDLE_MS, + }); #redisSub: ReturnType | null = null; // -- Lifecycle --------------------------------------------------- @@ -148,13 +171,6 @@ export class BroadcastService extends PuterService { body: unknown, headers: IncomingHeaders, ): Promise { - if (!rawBody) { - return { - ok: false, - status: 400, - message: 'Missing or invalid body', - }; - } if (!body || typeof body !== 'object') { return { ok: false, status: 400, message: 'Invalid JSON body' }; } @@ -170,6 +186,75 @@ export class BroadcastService extends PuterService { }; } + const verified = await this.verifySignedRequest(rawBody, headers); + if (!verified.ok || verified.info?.ignored) return verified; + + await this.#emitIncomingEventsSequentially(incomingEvents); + return { ok: true }; + } + + // -- Addressed peer channel -------------------------------------- + + /** + * What this deployment calls itself. Peers name each other by this id, and + * a deployment that configured none is on its own. + */ + get regionId(): string { + return this.#resolveLocalPeerId() ?? DEFAULT_REGION_ID; + } + + /** Peers this node may address directly, by id. */ + get addressablePeers(): string[] { + return this.#webhookPeers + .map((peer) => this.#resolvePeerIdOf(peer)) + .filter((id): id is string => id !== null); + } + + /** + * One signed POST to one named peer, on the same credentials and replay + * protection the all-peers fan uses, over the keep-alive agent. Returns the + * peer's parsed answer; throws when it could not be reached or refused — a + * caller must be able to tell "the peer said no" from "we never heard", + * because only the first authorises acting on the peer's behalf. + */ + async postToPeer( + peerId: string, + path: string, + payload: unknown, + ): Promise { + const peer = this.#peersByKey[peerId]; + if (!peer) throw new Error(`unknown broadcast peer ${peerId}`); + const response = await this.#post( + peer, + path, + JSON.stringify(payload), + this.#addressedHttpsAgent, + ADDRESSED_REQUEST_TIMEOUT_MS, + ); + try { + return JSON.parse(String(response ?? '')); + } catch { + return null; + } + } + + /** + * Verify an inbound signed POST without acting on its body: the half of + * `verifyAndEmit` an addressed channel needs, since what it carries is not + * an event-bus payload. + */ + async verifySignedRequest( + rawBody: Buffer | undefined, + headers: IncomingHeaders, + ): Promise { + if (!rawBody) { + return { + ok: false, + status: 400, + message: 'Missing or invalid body', + }; + } + const peerId = headers.peerId; if (!peerId) { return { @@ -235,7 +320,6 @@ export class BroadcastService extends PuterService { }; } - await this.#emitIncomingEventsSequentially(incomingEvents); return { ok: true }; } @@ -400,9 +484,29 @@ export class BroadcastService extends PuterService { peer: IBroadcastPeerConfig, events: BroadcastEvent[], ): Promise { + await this.#post( + peer, + null, + JSON.stringify({ events }), + this.#webhookHttpsAgent, + 15_000, + ); + } + + /** + * One signed POST to one peer. `path` replaces the configured webhook path + * for an addressed channel, and is `null` for the webhook itself. + */ + async #post( + peer: IBroadcastPeerConfig, + path: string | null, + rawBody: string, + agent: HttpsAgent, + timeoutMs: number, + ): Promise { const peerId = this.#resolvePeerIdOf(peer); if (!peerId) return; - const requestUrl = this.#normalizeWebhookUrl(peer.webhook_url); + const requestUrl = this.#normalizeWebhookUrl(peer.webhook_url, path); const mySecret = this.#self()?.secret; if (!requestUrl || !mySecret) return; @@ -411,7 +515,6 @@ export class BroadcastService extends PuterService { const nextNonce = await this.#nextOutgoingNonce(peerId); const timestamp = Math.floor(Date.now() / 1000); - const rawBody = JSON.stringify({ events }); const payloadToSign = `${timestamp}.${nextNonce}.${rawBody}`; const signature = createHmac('sha256', mySecret) .update(payloadToSign) @@ -433,15 +536,13 @@ export class BroadcastService extends PuterService { url: requestUrl, headers, data: rawBody, - timeout: 15_000, + timeout: timeoutMs, // We translate non-2xx into a thrown error ourselves so we // can log the response body on failure. validateStatus: () => true, responseType: 'text', transformResponse: (value: unknown) => value, - ...(requestUrl.startsWith('https:') - ? { httpsAgent: this.#webhookHttpsAgent } - : {}), + ...(requestUrl.startsWith('https:') ? { httpsAgent: agent } : {}), }); if (response.status < 200 || response.status >= 300) { @@ -452,6 +553,7 @@ export class BroadcastService extends PuterService { `Webhook POST failed: ${response.status} ${response.statusText}`, ); } + return response.data as string | undefined; } // -- Inbound helpers --------------------------------------------- @@ -623,7 +725,10 @@ export class BroadcastService extends PuterService { return this.config.broadcast ?? {}; } - #normalizeWebhookUrl(url: string | undefined): string | null { + #normalizeWebhookUrl( + url: string | undefined, + path: string | null = null, + ): string | null { if (typeof url !== 'string' || url.trim() === '') return null; const trimmed = url.trim(); let parsed: URL; @@ -637,6 +742,7 @@ export class BroadcastService extends PuterService { // Coerce protocol so a misconfigured `http://...` peer URL still // gets sent over our preferred transport. parsed.protocol = `${this.#webhookProtocol}:`; + if (path !== null) parsed.pathname = path; return parsed.toString(); } diff --git a/src/backend/services/events/EventForwardService.ts b/src/backend/services/events/EventForwardService.ts new file mode 100644 index 000000000..0e61d5b3f --- /dev/null +++ b/src/backend/services/events/EventForwardService.ts @@ -0,0 +1,592 @@ +/* + * Copyright (C) 2024-present Puter Technologies Inc. + * + * This file is part of Puter. + * + * Puter is free software: you can redistribute it and/or modify + * it under the terms of the GNU Affero General Public License as published + * by the Free Software Foundation, either version 3 of the License, or + * (at your option) any later version. + * + * This program is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + * GNU Affero General Public License for more details. + * + * You should have received a copy of the GNU Affero General Public License + * along with this program. If not, see . + */ + +import type { Actor } from '../../core/actor.js'; +import { + PRESENCE_NO_APP, + type PresenceRow, +} from '../../stores/events/PresenceStore.js'; +import { + appSocketRoom, + type SocketSpecifier, +} from '../socket/SocketService.js'; +import { PuterService } from '../types.js'; +import { + FORWARD_MAX_QUEUED, + PeerForwardQueue, + type ForwardBatch, + type ForwardDelivery, + type ForwardItem, + type ForwardReply, +} from './forwardQueue.js'; +import { PresenceCache, remoteRegions } from './presenceCache.js'; +import type { DeliverableEvent, GapMarker } from './registry.js'; + +/** + * Getting a socket delivery to the region that holds the socket. + * + * Cross-region traffic here is proportional to matched socket deliveries with a + * live client somewhere else — never to event volume, and never to connected + * population. Three things hold that line: + * + * - **Presence says where to send**, and is read through a generation-keyed + * region-local cache, so a busy subscription against a settled row reads the + * table once. + * - **Nothing is broadcast on the chance someone is listening.** An empty row is + * no hop at all, which is the common case, and a deployment with no peers + * configured never gets as far as reading one. + * - **Wrong rows are corrected on the read path.** A peer that answers "no + * socket" from its own registry is authoritative, and its region is removed + * with a write conditional on the version that was read. A peer that times + * out is not: a timeout is ambiguous, and evicting on it would blackhole a + * healthy region for the length of a partition and beyond it. + * + * No delivery state crosses a region. A `single`'s lease, retry counter and + * queue live where it was emitted; a peer holding the socket relays the + * client's ack home, and that is the whole of what it knows. + */ + +/** Path peers accept addressed event batches on. */ +export const FORWARD_WEBHOOK_PATH = '/broadcast/events'; + +/** Identities held so a delivery need not look one up to address a row. */ +const UUID_CACHE_MAX = 10_000; + +/** Where a delivery is put down once it reaches the region holding the socket. */ +export const forwardTarget = ( + userId: number, + appUid: string | null, +): SocketSpecifier => + appUid ? { room: appSocketRoom(userId, appUid) } : { room: String(userId) }; + +/** The presence pair a socket, or a subscription row, belongs to. */ +export const presenceApp = (appUid: string | null | undefined): string => + appUid ?? PRESENCE_NO_APP; + +/** One delivery, in the terms this service addresses it by. */ +export interface ForwardableDelivery { + holderUserId: number; + appUid: string | null; + subId: string; + event: DeliverableEvent; + ackRequired?: true; + ackId?: string; +} + +/** What stands in for events a full queue could not carry across. */ +const overflowGap = (event: DeliverableEvent): GapMarker => ({ + id: event.id, + subject: event.subject, + op: 'gap', + reason: 'backlog_overflow', + ts: event.ts, +}); + +export class EventForwardService extends PuterService { + readonly #cache = new PresenceCache(); + /** UserId → uuid, which is what the presence row is keyed by. */ + readonly #uuids = new Map(); + readonly #leaveTimers = new Map>(); + /** Peer → subscriptions with a gap marker still queued for it. */ + readonly #pendingMarkers = new Map>(); + #queue: PeerForwardQueue | null = null; + #draining = false; + + /** + * How long a region waits before taking itself out of a row. Reloads, flaky + * networks and rolling deploys all reconnect inside it and write nothing at + * all, which is the point — the write is an optimisation, since lazy repair + * is what actually keeps rows honest. + */ + static LEAVE_DELAY_MIN_MS = 30_000; + static LEAVE_DELAY_MAX_MS = 60_000; + + /** Deliveries held for one peer before the oldest are shed with markers. */ + static MAX_QUEUED = FORWARD_MAX_QUEUED; + + override onServerStart(): void { + this.clients.event.on( + 'outer.pubsub.events.presenceBumped', + (_key, data, meta) => { + // Our own emit reaches local listeners too, and that half has + // already been applied. + if (!(meta as { from_outside?: boolean })?.from_outside) return; + const { userId } = (data ?? {}) as { userId?: number }; + if (typeof userId === 'number') this.#cache.bump(userId); + }, + ); + } + + override async onServerPrepareShutdown(): Promise { + // A drain is not a disconnect: every socket on this node is going at + // once, and writing a row for each would be the burst the transition- + // only rule exists to avoid. Lazy repair corrects what is left. + this.#draining = true; + for (const timer of this.#leaveTimers.values()) clearTimeout(timer); + this.#leaveTimers.clear(); + await this.#queue?.flushAll(); + this.#queue?.stop(); + } + + /** + * Whether anything here has work to do. A deployment with no peers has + * nowhere to forward and nobody to read its rows, so it takes no part in + * presence at all — not a write, not a read, not a timer. + */ + get active(): boolean { + return ( + this.config.events?.enabled === true && + this.services.broadcast.addressablePeers.length > 0 + ); + } + + /** What this deployment calls itself in a presence row. */ + get region(): string { + return this.services.broadcast.regionId; + } + + // -- Presence transitions ---------------------------------------- + + /** + * One more connection for the pair. Only the one that crosses zero in this + * region writes: every reconnect after it is a counter increment. + */ + async noteConnect(actor: Actor): Promise { + if (!this.active) return; + const pair = this.#pairOf(actor); + if (!pair) return; + + // A reconnect inside the window cancels the write the disconnect owed + // — and owes none of its own, because this region never left the row. + // A reload, a flaky network and a rolling deploy all land here. + const pending = this.#leaveTimers.get(pair.key); + if (pending) { + clearTimeout(pending); + this.#leaveTimers.delete(pair.key); + } + + const count = await this.stores.presence.addConnection( + pair.userId, + pair.appUid, + ); + if (count !== 1 || pending) return; + await this.stores.presence.join( + pair.userUuid, + pair.appUid, + this.region, + ); + await this.#bump(pair.userId); + } + + /** + * One connection gone. The last one owes the region's removal, but not yet: + * a reload is a disconnect followed immediately by a connect, and both + * writes are avoidable. + */ + async noteDisconnect(actor: Actor): Promise { + if (!this.active) return; + const pair = this.#pairOf(actor); + if (!pair) return; + + const count = await this.stores.presence.removeConnection( + pair.userId, + pair.appUid, + ); + if (count > 0 || this.#draining) return; + + const delay = + EventForwardService.LEAVE_DELAY_MIN_MS + + Math.random() * + Math.max( + 0, + EventForwardService.LEAVE_DELAY_MAX_MS - + EventForwardService.LEAVE_DELAY_MIN_MS, + ); + const timer = setTimeout(() => { + this.#leaveTimers.delete(pair.key); + void this.#leaveIfStillGone(pair).catch((err: unknown) => { + console.warn('[events] presence leave failed', err); + }); + }, delay); + timer.unref?.(); + this.#leaveTimers.set(pair.key, timer); + } + + async #leaveIfStillGone(pair: PresencePair): Promise { + if (this.#draining) return; + if ( + await this.stores.presence.holdsConnection(pair.userId, pair.appUid) + ) + return; + const row = await this.stores.presence.read(pair.userUuid, pair.appUid); + if (!(this.region in row.regions)) return; + await this.stores.presence.leave( + pair.userUuid, + pair.appUid, + this.region, + row.version, + ); + await this.#bump(pair.userId); + } + + // -- Forwarding -------------------------------------------------- + + /** + * Hand one `broadcast` delivery to every other region holding a socket for + * the pair. The emitting region has already delivered its own copy and + * never asks presence about itself. + */ + async fanOut(delivery: ForwardableDelivery): Promise { + if (!this.active) return; + const regions = await this.regionsFor( + delivery.holderUserId, + delivery.appUid, + ); + for (const region of regions) this.#send(region, delivery); + } + + /** + * The region to try for one `single` attempt, or `null` when the candidates + * are spent. Ordered most-recently-connected first, which is the socket + * most likely to still be there. + */ + async candidateRegion( + holderUserId: number, + appUid: string | null, + attempt: number, + ): Promise { + if (!this.active) return null; + const regions = await this.regionsFor(holderUserId, appUid); + return regions[attempt] ?? null; + } + + /** Send one `single` to a named region, carrying what settles it. */ + handOff(region: string, delivery: ForwardableDelivery): void { + this.#send(region, delivery); + } + + /** + * Relay a client's settle to the region that owns the lease. Nothing is + * settled locally: the queue it belongs to is not here. + */ + relayAck( + region: string, + userId: number, + subId: string, + entryId: string, + ): void { + if (!this.active) return; + if (!this.services.broadcast.addressablePeers.includes(region)) return; + this.#queueFor().push(region, { + kind: 'ack', + userId, + subId, + entryId, + }); + } + + /** Whether a region name is a peer this node can address. */ + isPeer(region: string): boolean { + return this.services.broadcast.addressablePeers.includes(region); + } + + /** + * Regions other than this one holding a socket for the pair, read through + * the generation-keyed cache. + */ + async regionsFor( + holderUserId: number, + appUid: string | null, + ): Promise { + const app = presenceApp(appUid); + const cached = this.#cache.read(holderUserId, app); + if (cached) return remoteRegions(cached, this.region); + + const userUuid = await this.#uuidOf(holderUserId); + if (!userUuid) return []; + + const epoch = this.#cache.generationOf(holderUserId); + let row: PresenceRow; + try { + row = await this.stores.presence.read(userUuid, app); + } catch (err) { + console.warn('[events] presence read failed', err); + return []; + } + this.#cache.write(holderUserId, app, epoch, row); + return remoteRegions(row, this.region); + } + + // -- Inbound ----------------------------------------------------- + + /** + * Apply one peer's batch. Deliveries go out over this region's own sockets; + * acks settle where the lease actually lives, which is here. The reply + * names the pairs this region holds nothing for — the only signal that lets + * the sender edit the row. + */ + async receive(batch: ForwardBatch): Promise { + const noSocket: Array<{ userId: number; appUid: string | null }> = []; + const checked = new Set(); + + for (const item of batch.items ?? []) { + if (item.kind === 'ack') { + await this.services.events + .settleRelayedAck(item.userId, item.subId, item.entryId) + .catch((err: unknown) => { + console.warn('[events] relayed ack failed', err); + }); + continue; + } + if (item.kind !== 'delivery') continue; + + await this.services.events.deliverForwarded(item); + + const key = `${item.userId}|${presenceApp(item.appUid)}`; + if (checked.has(key)) continue; + checked.add(key); + if (!(await this.#holdsSocket(item))) + noSocket.push({ userId: item.userId, appUid: item.appUid }); + } + + return noSocket.length > 0 ? { noSocket } : {}; + } + + /** + * Whether this region has anywhere to put deliveries for the pair. The + * per-region connection counter is the honest answer: the socket registry + * on one node says nothing about the others. + */ + async #holdsSocket(item: ForwardDelivery): Promise { + if (this.services.socket.has(forwardTarget(item.userId, item.appUid))) + return true; + // Still in the row on purpose: a pair inside its disconnect window is + // one this region expects back, and has not written itself out for. + if (this.#leaveTimers.has(`${item.userId}|${presenceApp(item.appUid)}`)) + return true; + try { + return await this.stores.presence.holdsConnection( + item.userId, + presenceApp(item.appUid), + ); + } catch { + // Unable to tell is not "definitely not": a repair has to be + // affirmative, so this reports a socket rather than inviting one. + return true; + } + } + + // -- Transport --------------------------------------------------- + + #send(region: string, delivery: ForwardableDelivery): void { + if (!this.isPeer(region)) return; + const item: ForwardDelivery = { + kind: 'delivery', + userId: delivery.holderUserId, + appUid: delivery.appUid, + subId: delivery.subId, + event: delivery.event, + ...(delivery.ackRequired + ? { + ackRequired: true as const, + ackId: delivery.ackId, + origin: this.region, + } + : {}), + }; + this.#queueFor().push(region, item); + } + + #queueFor(): PeerForwardQueue { + this.#queue ??= new PeerForwardQueue({ + maxQueued: EventForwardService.MAX_QUEUED, + send: (peerId, items) => this.#ship(peerId, items), + onOverflow: (peerId, dropped) => this.#overflowed(peerId, dropped), + }); + return this.#queue; + } + + async #ship(peerId: string, items: ForwardItem[]): Promise { + // Shipped or lost, these are out of the queue either way, so the next + // shed for their subscriptions may queue a marker again. + this.#forgetMarkers(peerId, items); + const batch: ForwardBatch = { from: this.region, items }; + const reply = (await this.services.broadcast.postToPeer( + peerId, + FORWARD_WEBHOOK_PATH, + batch, + )) as ForwardReply | null; + + for (const missing of reply?.noSocket ?? []) + await this.#repair(peerId, missing.userId, missing.appUid); + } + + /** + * Take a region out of a row it is no longer in, on the strength of that + * region saying so. Conditional on the version that was read: a connect + * racing this repair bumps it, and the repair loses harmlessly rather than + * blackholing a socket that just arrived. + */ + async #repair( + region: string, + userId: number, + appUid: string | null, + ): Promise { + const app = presenceApp(appUid); + if (!this.#cache.claimRepair(userId, app, region)) return; + + const userUuid = await this.#uuidOf(userId); + if (!userUuid) return; + const row = this.#cache.read(userId, app); + if (!row || !(region in row.regions)) return; + + try { + const applied = await this.stores.presence.leave( + userUuid, + app, + region, + row.version, + ); + if (!applied) return; + this.#cache.forget(userId, app, region); + await this.#bump(userId); + } catch (err) { + console.warn('[events] presence repair failed', err); + } + } + + /** + * Deliveries the queue could not hold. Never a log line on its own: each + * subscription that lost events gets a marker in their place — handed back + * to the queue, which puts it where the shed just made room — and the + * region being this far behind is worth someone's attention. + * + * One marker per (peer, subscription) at a time. A subscription still + * losing events while its marker waits has already been told, and a second + * marker would take a slot from a delivery; a marker that is itself shed + * frees the slot for the next one. + */ + #overflowed(peerId: string, dropped: ForwardItem[]): ForwardItem[] { + const pending = this.#pendingMarkersFor(peerId); + const markers: ForwardItem[] = []; + for (const item of dropped) { + if (item.kind !== 'delivery') continue; + if (item.event.op === 'gap') { + pending.delete(item.subId); + continue; + } + // A `single` loses nothing here: its lease is still running, and + // the next attempt is what the expiry is for. + if (item.ackRequired || pending.has(item.subId)) continue; + pending.add(item.subId); + markers.push({ ...item, event: overflowGap(item.event) }); + } + + this.clients.alarm.create( + 'events_forward_overflow', + 'Events queued for another region were dropped to stay inside the queue bound', + { peerId, dropped: dropped.length }, + 'warning', + { dedup: true }, + ); + return markers; + } + + #pendingMarkersFor(peerId: string): Set { + let pending = this.#pendingMarkers.get(peerId); + if (!pending) { + pending = new Set(); + this.#pendingMarkers.set(peerId, pending); + } + return pending; + } + + #forgetMarkers(peerId: string, items: ForwardItem[]): void { + const pending = this.#pendingMarkers.get(peerId); + if (!pending) return; + for (const item of items) + if (item.kind === 'delivery' && item.event.op === 'gap') + pending.delete(item.subId); + } + + // -- Generation -------------------------------------------------- + + /** + * Say that presence moved. One `INCR` and one broadcast, the same shape the + * watched-token set and the permission cache use — which is what makes + * reads scale with transitions rather than with events. + */ + async #bump(userId: number): Promise { + try { + const generation = + await this.stores.presence.bumpGeneration(userId); + this.#cache.bump(userId, generation); + this.clients.event.emit( + 'outer.pubsub.events.presenceBumped', + { userId, generation }, + {}, + ); + } catch (err) { + this.#cache.bump(userId); + console.warn('[events] presence generation bump failed', err); + } + } + + // -- Plumbing ---------------------------------------------------- + + #pairOf(actor: Actor): PresencePair | null { + const userId = actor.user?.id; + const userUuid = actor.user?.uuid; + if (typeof userId !== 'number' || !userUuid) return null; + const appUid = presenceApp(actor.effectiveApp?.uid ?? null); + this.#rememberUuid(userId, userUuid); + return { userId, userUuid, appUid, key: `${userId}|${appUid}` }; + } + + async #uuidOf(userId: number): Promise { + const held = this.#uuids.get(userId); + if (held) return held; + try { + const user = await this.stores.user.getById(userId); + const uuid = user?.uuid; + if (!uuid) return null; + this.#rememberUuid(userId, uuid); + return uuid; + } catch { + return null; + } + } + + #rememberUuid(userId: number, uuid: string): void { + this.#uuids.delete(userId); + this.#uuids.set(userId, uuid); + while (this.#uuids.size > UUID_CACHE_MAX) { + const oldest = this.#uuids.keys().next(); + if (oldest.done) break; + this.#uuids.delete(oldest.value); + } + } +} + +interface PresencePair { + userId: number; + userUuid: string; + appUid: string; + key: string; +} diff --git a/src/backend/services/events/EventsService.test.ts b/src/backend/services/events/EventsService.test.ts index 0e1539ec9..e53faa78a 100644 --- a/src/backend/services/events/EventsService.test.ts +++ b/src/backend/services/events/EventsService.test.ts @@ -77,7 +77,7 @@ type GenerationBumpHandler = ( ) => void; const remoteGenerationBumpHandler = (): GenerationBumpHandler | undefined => eventBus.on.mock.calls.find( - ([key]: [string]) => key === 'outer.events.generationBumped', + ([key]: [string]) => key === 'outer.pubsub.events.generationBumped', )?.[1] as GenerationBumpHandler | undefined; const COUNTED = new Set([ @@ -268,6 +268,18 @@ const buildService = ( permission: permissionStore, } as never, { + eventForward: { + // A deployment with no peers has nowhere to forward to, which + // is what every test here is. + region: 'local', + isPeer: () => false, + noteConnect: async () => undefined, + noteDisconnect: async () => undefined, + candidateRegion: async () => null, + fanOut: async () => undefined, + handOff: () => undefined, + relayAck: () => undefined, + }, socket: { send: vi.fn(async (spec: { socket?: string }, _key, data) => { outbox.push({ @@ -683,7 +695,7 @@ describe('cross-process invalidation', () => { const handler = remoteGenerationBumpHandler(); expect(handler).toBeDefined(); handler?.( - 'outer.events.generationBumped', + 'outer.pubsub.events.generationBumped', { userId, generation }, { from_outside: true }, ); @@ -702,7 +714,7 @@ describe('cross-process invalidation', () => { // No `from_outside`: this is what the local half of our own emit // looks like, and it must not force a redundant re-check. remoteGenerationBumpHandler()?.( - 'outer.events.generationBumped', + 'outer.pubsub.events.generationBumped', { userId, generation: 1 }, {}, ); @@ -716,14 +728,14 @@ describe('cross-process invalidation', () => { const handler = remoteGenerationBumpHandler(); handler?.( - 'outer.events.generationBumped', + 'outer.pubsub.events.generationBumped', { userId, generation: 1, durable: false }, { from_outside: true }, ); expect(cold).not.toHaveBeenCalled(); handler?.( - 'outer.events.generationBumped', + 'outer.pubsub.events.generationBumped', { userId, generation: 2, durable: true }, { from_outside: true }, ); @@ -766,7 +778,7 @@ describe('cross-process invalidation', () => { }); remoteGenerationBumpHandler()?.( - 'outer.events.generationBumped', + 'outer.pubsub.events.generationBumped', { userId, generation: 1 }, // behind this process's own recorded generation { from_outside: true }, ); diff --git a/src/backend/services/events/EventsService.ts b/src/backend/services/events/EventsService.ts index 2b96d66dd..ea00c3a98 100644 --- a/src/backend/services/events/EventsService.ts +++ b/src/backend/services/events/EventsService.ts @@ -107,6 +107,8 @@ import { type SubscriptionGrant, } from './authorization.js'; import { DeliveryCoalescer } from './coalescer.js'; +import { forwardTarget } from './EventForwardService.js'; +import type { ForwardDelivery } from './forwardQueue.js'; import { DELIVERY_USAGE_TYPES, EVENTS_COSTS, @@ -192,6 +194,12 @@ export interface UnsubscribeRequest { export interface AckRequest { subId?: unknown; id?: unknown; + /** + * The region that owns the lease, echoed back from the delivery. Absent + * when it was emitted where the client is connected, which is everything a + * single-region deployment ever sees. + */ + origin?: unknown; } /** Body of `POST /events/subscribe`. */ @@ -304,6 +312,12 @@ export interface DeliveryEnvelope { ackRequired?: true; /** The handle `events.ack` names — the delivery, not the event. */ ackId?: string; + /** + * Which region holds the lease, present only when that is not the region + * the client is connected to. Echoed back with the ack so whichever region + * receives it knows where to send it. + */ + origin?: string; } /** @@ -315,6 +329,12 @@ interface AddressedDelivery { envelope: DeliveryEnvelope; /** False for a row that asked for its handler and no socket copy. */ socket: boolean; + /** + * Whether the socket copy also goes to the other regions holding one. True + * for `broadcast`, which every connected subscriber gets; a `single` picks + * exactly one region itself and never fans. + */ + remote?: boolean; worker?: WorkerInvocation; meter: DeliveryMeter; /** @@ -851,7 +871,7 @@ export class EventsService extends PuterService { ); this.clients.event.on( - 'outer.events.generationBumped', + 'outer.pubsub.events.generationBumped', (_key, data, meta) => { // Our own emit reaches local listeners too, and that half has // already been applied. @@ -971,11 +991,67 @@ export class EventsService extends PuterService { }); }) as (...args: never[]) => void); + // Presence is per (user, app) and per region, so it is the connection + // rather than the subscription that moves it — a client that has not + // subscribed yet is still somewhere, and that is what a peer needs. + void this.services.eventForward.noteConnect(actor).catch((err) => { + console.warn('[events] presence connect failed', err); + }); + socket.once('disconnect', (() => { void this.reapSocket(userId, socket.id); + void this.services.eventForward + .noteDisconnect(actor) + .catch((err) => { + console.warn('[events] presence disconnect failed', err); + }); }) as (...args: never[]) => void); } + /** + * Put down one delivery another region handed us. Nothing is metered or + * re-checked: the region that emitted it did both, and this one holds no + * state about it at all. + */ + async deliverForwarded(item: ForwardDelivery): Promise { + if (!this.enabled) return; + const envelope: DeliveryEnvelope = { + subId: item.subId, + event: item.event, + ...(item.ackRequired + ? { + ackRequired: true as const, + ackId: item.ackId, + origin: item.origin, + } + : {}), + }; + await this.services.socket.send( + forwardTarget(item.userId, item.appUid), + EVENTS_DELIVERY_CHANNEL, + envelope, + ); + } + + /** + * Settle a delivery whose client acked it in another region. The same + * settle a local ack performs — including the guard that a suspended row + * hands out nothing more — minus the checks the connected region already + * made against the client. + */ + async settleRelayedAck( + holderUserId: number, + subId: string, + entryId: string, + ): Promise { + if (!this.enabled || !subId || !entryId) return; + const row = await this.stores.durableSubscription.getBySubId(subId); + if (!row || row.holderUserId !== holderUserId) return; + + await this.stores.pendingDelivery.settle(subId, entryId); + if (row.suspendedAt === null) await this.#drain(row); + } + async #answer( ack: unknown, run: () => Promise, @@ -1276,6 +1352,16 @@ export class EventsService extends PuterService { ) throw unknownSubscription(); + // The lease lives where the event was emitted, which is not always + // where the client ended up connected. Nothing about the delivery is + // held here to settle, so the ack goes home rather than being applied. + const origin = String(request?.origin ?? ''); + const forward = this.services.eventForward; + if (origin && origin !== forward.region && forward.isPeer(origin)) { + forward.relayAck(origin, holderUserId, subId, entryId); + return; + } + await this.stores.pendingDelivery.settle(subId, entryId); // A suspended row is not delivered to: settling what it already handed // out must not be the trigger that hands out the next one. @@ -2165,6 +2251,9 @@ export class EventsService extends PuterService { target: deliveryTarget(row), envelope: { subId: row.subId, event }, socket: targets.includes('socket'), + // A session row is addressed at one connection, which is the + // one that made it and is therefore here. + remote: row.socketId === undefined, worker: this.#workerInvocation(row, event), meter: meterFor(row), // `broadcast` is one send per delivery — no retry to dedup. @@ -2406,7 +2495,8 @@ export class EventsService extends PuterService { if (!targetsOf(row).includes('socket')) continue; void this.#send({ target: deliveryTarget(row), - socket: true, + socket: targetsOf(row).includes('socket'), + remote: row.socketId === undefined, envelope: { subId: row.subId, event: marker }, meter: meterFor(row), // A gap marker is never billed regardless — `#delivered`'s own @@ -2850,30 +2940,53 @@ export class EventsService extends PuterService { const targets = targetsOf(row); const target = deliveryTarget(row); - // Only this region's own connections are visible here; the ones other - // regions hold arrive with presence. - const overSocket = + // Candidates in order: this region's own connection, then the regions + // presence names, most recently connected first. The attempt counter + // spans them, so two attempts is two candidates however they are split. + const socketsLeft = targets.includes('socket') && - claimed.socketAttempts < SINGLE_SOCKET_ATTEMPTS && - this.services.socket.has(target); + claimed.socketAttempts < SINGLE_SOCKET_ATTEMPTS; + const here = socketsLeft && this.services.socket.has(target); + const region = + socketsLeft && !(here && claimed.socketAttempts === 0) + ? await this.services.eventForward.candidateRegion( + row.holderUserId, + row.appUid, + claimed.socketAttempts - (here ? 1 : 0), + ) + : null; - if (overSocket) { + if (here || region) { await this.stores.pendingDelivery.recordSocketAttempt( row.subId, claimed.entryId, ); - await this.#send({ - target, - socket: true, - envelope: { + const envelope: DeliveryEnvelope = { + subId: row.subId, + event: claimed.event, + ackRequired: true, + ackId: claimed.entryId, + }; + const bill = await this.#firstAttempt(row.subId, claimed); + if (region) { + this.services.eventForward.handOff(region, { + holderUserId: row.holderUserId, + appUid: row.appUid, subId: row.subId, event: claimed.event, ackRequired: true, ackId: claimed.entryId, - }, - meter, - bill: await this.#firstAttempt(row.subId, claimed), - }); + }); + this.#delivered(envelope, meter, bill); + } else { + await this.#send({ + target, + socket: true, + envelope, + meter, + bill, + }); + } return false; } @@ -3108,6 +3221,7 @@ export class EventsService extends PuterService { await this.#send({ target: delivery.target, socket: delivery.socket, + remote: delivery.remote, meter: delivery.meter, envelope: { subId: delivery.envelope.subId, @@ -3141,6 +3255,20 @@ export class EventsService extends PuterService { } catch (err) { console.warn('[events] socket send failed', err); } + + // Every connected subscriber gets a `broadcast`, and the ones this + // region cannot reach are wherever presence says they are. + if (delivery.remote) + void this.services.eventForward + .fanOut({ + holderUserId: delivery.meter.holderUserId, + appUid: delivery.meter.appUid, + subId: delivery.envelope.subId, + event: delivery.envelope.event, + }) + .catch((err: unknown) => { + console.warn('[events] forward failed', err); + }); } let invoked = false; @@ -3404,7 +3532,7 @@ export class EventsService extends PuterService { this.#cache.bump(userId, generation); try { this.clients.event.emit( - 'outer.events.generationBumped', + 'outer.pubsub.events.generationBumped', { userId, generation, durable }, {}, ); diff --git a/src/backend/services/events/durable.integration.test.ts b/src/backend/services/events/durable.integration.test.ts index cb0285097..9d5e861f4 100644 --- a/src/backend/services/events/durable.integration.test.ts +++ b/src/backend/services/events/durable.integration.test.ts @@ -529,7 +529,7 @@ describe('what a write costs once a durable row exists', () => { const after = `${anchor}/post-peer-${uuidv4().slice(0, 8)}.txt`; const reads = await countTableReads(async () => { await env.server.clients.event.emitAndWait( - 'outer.events.generationBumped', + 'outer.pubsub.events.generationBumped', { userId, generation: generation + 1, durable: true }, { from_outside: true }, ); @@ -553,7 +553,7 @@ describe('what a write costs once a durable row exists', () => { const path = `${anchor}/peer-session-${uuidv4().slice(0, 8)}.txt`; const reads = await countTableReads(async () => { await env.server.clients.event.emitAndWait( - 'outer.events.generationBumped', + 'outer.pubsub.events.generationBumped', { userId, generation: generation + 1, durable: false }, { from_outside: true }, ); diff --git a/src/backend/services/events/forwardQueue.test.ts b/src/backend/services/events/forwardQueue.test.ts new file mode 100644 index 000000000..492023d15 --- /dev/null +++ b/src/backend/services/events/forwardQueue.test.ts @@ -0,0 +1,275 @@ +/* + * Copyright (C) 2024-present Puter Technologies Inc. + * + * This file is part of Puter. + * + * Puter is free software: you can redistribute it and/or modify + * it under the terms of the GNU Affero General Public License as published + * by the Free Software Foundation, either version 3 of the License, or + * (at your option) any later version. + * + * This program is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + * GNU Affero General Public License for more details. + * + * You should have received a copy of the GNU Affero General Public License + * along with this program. If not, see . + */ + +import { afterEach, describe, expect, it, vi } from 'vitest'; +import { PeerForwardQueue, type ForwardItem } from './forwardQueue.js'; + +/** + * One request per window per peer, never one per event, and nothing lost + * without somebody being told. + */ + +const queues: PeerForwardQueue[] = []; + +const delivery = (subId: string, id = subId): ForwardItem => ({ + kind: 'delivery', + userId: 7, + appUid: 'app-a', + subId, + event: { + id, + subject: 'fs:/u7/Documents', + op: 'write', + uid: 'node-1', + path: '/u7/Documents/notes.txt', + self: true, + seq: 0, + ts: 1, + }, +}); + +const makeQueue = ( + over: Partial[0]> = {}, +) => { + const sent: Array<{ peerId: string; items: ForwardItem[] }> = []; + const overflowed: Array<{ peerId: string; dropped: ForwardItem[] }> = []; + const queue = new PeerForwardQueue({ + flushMs: 10, + send: async (peerId, items) => { + sent.push({ peerId, items }); + }, + onOverflow: (peerId, dropped) => overflowed.push({ peerId, dropped }), + ...over, + }); + queues.push(queue); + return { queue, sent, overflowed }; +}; + +const settled = ( + sent: Array<{ peerId: string; items: ForwardItem[] }>, + count = 1, +): Promise => + vi.waitFor(() => expect(sent.length).toBeGreaterThanOrEqual(count), { + timeout: 2_000, + interval: 5, + }); + +afterEach(() => { + while (queues.length) queues.pop()?.stop(); +}); + +describe('the addressed peer queue', () => { + it('carries a window of events in one request', async () => { + const { queue, sent } = makeQueue(); + + for (let i = 0; i < 8; i++) queue.push('east', delivery(`sub-${i}`)); + await settled(sent); + + expect(sent).toHaveLength(1); + expect(sent[0].peerId).toBe('east'); + expect(sent[0].items).toHaveLength(8); + }); + + it('keeps each peer to its own batch', async () => { + const { queue, sent } = makeQueue(); + + queue.push('east', delivery('a')); + queue.push('south', delivery('b')); + await settled(sent, 2); + + expect(sent.map((call) => call.peerId).sort()).toEqual([ + 'east', + 'south', + ]); + for (const call of sent) expect(call.items).toHaveLength(1); + }); + + it('ships early once the item bound trips', async () => { + const { queue, sent } = makeQueue({ maxItems: 3, flushMs: 60_000 }); + + for (let i = 0; i < 3; i++) queue.push('east', delivery(`sub-${i}`)); + await settled(sent); + + expect(sent[0].items).toHaveLength(3); + }); + + it('ships early once the byte bound trips', async () => { + const { queue, sent } = makeQueue({ maxBytes: 200, flushMs: 60_000 }); + + queue.push('east', delivery('a')); + queue.push('east', delivery('b')); + await settled(sent); + + expect(sent[0].items.length).toBeGreaterThan(0); + }); + + it('reports what it could not hold rather than dropping it quietly', async () => { + const { queue, overflowed } = makeQueue({ + maxQueued: 2, + flushMs: 60_000, + }); + + queue.push('east', delivery('a')); + queue.push('east', delivery('b')); + queue.push('east', delivery('c')); + queue.push('east', delivery('d')); + + expect(overflowed).toHaveLength(2); + expect(overflowed[0].peerId).toBe('east'); + expect( + overflowed.flatMap((call) => + call.dropped.map((item) => + item.kind === 'delivery' ? item.subId : item.kind, + ), + ), + ).toEqual(['a', 'b']); + }); + + it('does not lose the rest of a window to one failed send', async () => { + const failures: string[] = []; + const { queue, sent } = makeQueue({ + send: async (peerId, items) => { + if (failures.length === 0) { + failures.push(peerId); + throw new Error('peer unreachable'); + } + sentAfter.push({ peerId, items }); + }, + }); + const sentAfter: Array<{ peerId: string; items: ForwardItem[] }> = []; + + queue.push('east', delivery('a')); + await vi.waitFor(() => expect(failures).toHaveLength(1)); + + queue.push('east', delivery('b')); + await settled(sentAfter); + + expect(sentAfter[0].items).toHaveLength(1); + expect(sent).toHaveLength(0); + }); + + it('gives what arrived mid-flight the next window instead of the current one', async () => { + let release: (() => void) | null = null; + const sent: Array<{ peerId: string; items: ForwardItem[] }> = []; + const { queue } = makeQueue({ + maxItems: 1, + send: async (peerId, items) => { + sent.push({ peerId, items }); + if (sent.length === 1) + await new Promise((resolve) => { + release = resolve; + }); + }, + }); + + queue.push('east', delivery('a')); + await vi.waitFor(() => expect(release).not.toBeNull()); + queue.push('east', delivery('b')); + release!(); + + await settled(sent, 2); + expect( + sent[0].items.map((item) => (item as { subId: string }).subId), + ).toEqual(['a']); + expect( + sent[1].items.map((item) => (item as { subId: string }).subId), + ).toEqual(['b']); + }); + + it('sends nothing more once it is stopped', async () => { + const { queue, sent } = makeQueue(); + queue.stop(); + + queue.push('east', delivery('a')); + await new Promise((resolve) => setTimeout(resolve, 40)); + + expect(sent).toHaveLength(0); + }); + it('queues what the overflow handler returns without shedding again', async () => { + let overflows = 0; + const { queue, sent } = makeQueue({ + maxQueued: 2, + onOverflow: (_peerId, dropped) => { + overflows++; + return dropped.map((item) => ({ + ...(item as Extract), + event: { + ...(item as Extract) + .event, + op: 'gap', + }, + })) as ForwardItem[]; + }, + }); + + queue.push('east', delivery('a')); + queue.push('east', delivery('b')); + queue.push('east', delivery('c')); + await settled(sent); + + // One shed, one marker for it — not a marker per item down the queue. + expect(overflows).toBe(1); + expect( + sent[0].items.map((item) => + item.kind === 'delivery' + ? `${item.subId}:${item.event.op}` + : item.kind, + ), + ).toEqual(['b:write', 'c:write', 'a:gap']); + }); + + it('sheds deliveries before the markers that report them', async () => { + const { queue, sent, overflowed } = makeQueue({ + maxQueued: 2, + onOverflow: (_peerId, dropped) => { + overflowed.push({ peerId: 'east', dropped }); + return dropped.map((item) => ({ + ...(item as Extract), + event: { + ...(item as Extract) + .event, + op: 'gap', + }, + })) as ForwardItem[]; + }, + }); + + queue.push('east', delivery('a')); + queue.push('east', delivery('b')); + queue.push('east', delivery('c')); // sheds a, queues a:gap + queue.push('east', delivery('d')); // sheds b and c, never the marker + await settled(sent); + + expect( + overflowed.flatMap((call) => + call.dropped.map((item) => + item.kind === 'delivery' ? item.subId : item.kind, + ), + ), + ).toEqual(['a', 'b', 'c']); + expect( + sent[0].items.some( + (item) => + item.kind === 'delivery' && + item.event.op === 'gap' && + item.subId === 'a', + ), + ).toBe(true); + }); +}); diff --git a/src/backend/services/events/forwardQueue.ts b/src/backend/services/events/forwardQueue.ts new file mode 100644 index 000000000..76878449f --- /dev/null +++ b/src/backend/services/events/forwardQueue.ts @@ -0,0 +1,269 @@ +/* + * Copyright (C) 2024-present Puter Technologies Inc. + * + * This file is part of Puter. + * + * Puter is free software: you can redistribute it and/or modify + * it under the terms of the GNU Affero General Public License as published + * by the Free Software Foundation, either version 3 of the License, or + * (at your option) any later version. + * + * This program is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + * GNU Affero General Public License for more details. + * + * You should have received a copy of the GNU Affero General Public License + * along with this program. If not, see . + */ + +import type { DeliverableEvent } from './registry.js'; + +/** + * The addressed peer queue: one batch per peer per window, never one request + * per event. + * + * Separate from the all-peers generation fan because almost nothing about that + * queue suits this traffic. Its window is sized for cache invalidations and is + * far too long for a delivery a client is waiting on; its transport opens a + * fresh connection per send, which is invisible at that cadence and dominant at + * this one; and it drops the oldest entries on overflow behind a log line, + * which would break the one promise this system makes about loss — that it is + * always observable. Here an overflow is reported to the subscriptions that + * lost events, and alarms. + * + * Batching also amortises the two fixed costs of a signed send, the shared + * counter that allocates the anti-replay nonce and the signature over the body, + * across everything in the window instead of paying both per event. + */ + +// -- Wire ------------------------------------------------------------ + +/** One event handed to whichever region holds the socket for it. */ +export interface ForwardDelivery { + kind: 'delivery'; + /** Who the delivery is for, and which app's room it lands in. */ + userId: number; + appUid: string | null; + subId: string; + event: DeliverableEvent; + /** Set on a `single`: the far side asks the client to settle it. */ + ackRequired?: true; + ackId?: string; + /** The region holding the lease, which is where the ack has to end up. */ + origin?: string; +} + +/** A client's settle, on its way back to the region that owns the lease. */ +export interface ForwardAck { + kind: 'ack'; + userId: number; + subId: string; + entryId: string; +} + +export type ForwardItem = ForwardDelivery | ForwardAck; + +/** One batch, as a peer receives it. */ +export interface ForwardBatch { + /** Sending region, so the receiver can address a reply. */ + from: string; + items: ForwardItem[]; +} + +/** + * What a peer answers. `noSocket` names the pairs it holds no connection for, + * read off its own socket registry — the only signal that authorises a repair. + */ +export interface ForwardReply { + noSocket?: Array<{ userId: number; appUid: string | null }>; +} + +// -- Bounds ---------------------------------------------------------- + +/** + * How long items wait for company. Deliveries are already debounced an order of + * magnitude longer upstream, so nothing perceives this. + */ +export const FORWARD_FLUSH_MS = 25; + +/** Items and bytes one batch carries, whichever bound trips first. */ +export const FORWARD_MAX_ITEMS = 200; +export const FORWARD_MAX_BYTES = 256 * 1024; + +/** Items held for one peer before the oldest are shed with markers. */ +export const FORWARD_MAX_QUEUED = 5_000; + +export interface ForwardQueueOptions { + flushMs?: number; + maxItems?: number; + maxBytes?: number; + maxQueued?: number; + /** Ships one batch. Rejection is a lost batch, not a retry. */ + send: (peerId: string, items: ForwardItem[]) => Promise; + /** + * Items the queue could not hold. Never silent. Whatever it returns is + * queued in their place — the markers that say something was lost — and is + * not itself held to the bound, so a shed cannot cascade. + */ + onOverflow: ( + peerId: string, + dropped: ForwardItem[], + ) => ForwardItem[] | void; +} + +interface PeerQueue { + items: ForwardItem[]; + bytes: number; + timer: ReturnType | null; + flushing: boolean; +} + +export class PeerForwardQueue { + readonly #peers = new Map(); + readonly #options: Required< + Omit + > & + Pick; + #stopped = false; + + constructor(options: ForwardQueueOptions) { + this.#options = { + flushMs: options.flushMs ?? FORWARD_FLUSH_MS, + maxItems: options.maxItems ?? FORWARD_MAX_ITEMS, + maxBytes: options.maxBytes ?? FORWARD_MAX_BYTES, + maxQueued: options.maxQueued ?? FORWARD_MAX_QUEUED, + send: options.send, + onOverflow: options.onOverflow, + }; + } + + /** Queue one item for one peer, flushing early once a bound trips. */ + push(peerId: string, item: ForwardItem): void { + if (this.#stopped) return; + const queue = this.#queueFor(peerId); + queue.items.push(item); + queue.bytes += sizeOf(item); + + const { maxQueued, maxItems, maxBytes, flushMs } = this.#options; + if (queue.items.length > maxQueued) { + const dropped = shed(queue, queue.items.length - maxQueued); + // Appended here, past the bound check: a replacement sent back + // through `push` would trip the bound it just made room under and + // shed the next item, and so on down the whole queue. + const replacements = this.#options.onOverflow(peerId, dropped); + if (Array.isArray(replacements)) + for (const item of replacements) { + queue.items.push(item); + queue.bytes += sizeOf(item); + } + } + + if ( + !queue.flushing && + (queue.items.length >= maxItems || queue.bytes >= maxBytes) + ) { + void this.flush(peerId); + return; + } + if (queue.timer) return; + queue.timer = setTimeout(() => { + queue.timer = null; + void this.flush(peerId); + }, flushMs); + queue.timer.unref?.(); + } + + /** Ship whatever is held for one peer. */ + async flush(peerId: string): Promise { + const queue = this.#peers.get(peerId); + if (!queue || queue.flushing || queue.items.length === 0) return; + if (queue.timer) { + clearTimeout(queue.timer); + queue.timer = null; + } + + queue.flushing = true; + const batch = queue.items; + queue.items = []; + queue.bytes = 0; + try { + await this.#options.send(peerId, batch); + } catch (err) { + // A lost batch is a lost at-most-once delivery, and a `single` + // still has its lease to pace the next attempt. Re-queuing would + // hand the client an event whose moment has passed. + console.warn(`[events] forward to ${peerId} failed`, err); + } finally { + queue.flushing = false; + // Anything that arrived mid-flush gets the next window. + if (queue.items.length > 0 && !queue.timer && !this.#stopped) { + queue.timer = setTimeout(() => { + queue.timer = null; + void this.flush(peerId); + }, this.#options.flushMs); + queue.timer.unref?.(); + } + } + } + + async flushAll(): Promise { + await Promise.all([...this.#peers.keys()].map((id) => this.flush(id))); + } + + stop(): void { + this.#stopped = true; + for (const queue of this.#peers.values()) { + if (queue.timer) clearTimeout(queue.timer); + queue.timer = null; + } + } + + #queueFor(peerId: string): PeerQueue { + const held = this.#peers.get(peerId); + if (held) return held; + const queue: PeerQueue = { + items: [], + bytes: 0, + timer: null, + flushing: false, + }; + this.#peers.set(peerId, queue); + return queue; + } +} + +const sizeOf = (item: ForwardItem): number => { + try { + return JSON.stringify(item).length; + } catch { + return 0; + } +}; + +const isGapMarker = (item: ForwardItem): boolean => + item.kind === 'delivery' && item.event.op === 'gap'; + +/** + * Take `count` items off the old end, deliveries before markers: a marker is + * the one thing saying events were lost, so it is the last thing to go. Bytes + * come off per dropped item rather than being re-summed over the queue, which + * at the bound would walk every held item for every drop. + */ +const shed = (queue: PeerQueue, count: number): ForwardItem[] => { + const dropped: ForwardItem[] = []; + for (let i = 0; i < queue.items.length && dropped.length < count; ) { + if (isGapMarker(queue.items[i])) { + i++; + continue; + } + dropped.push(queue.items[i]); + queue.items.splice(i, 1); + } + // Nothing but markers left to give: the oldest of those go too. + while (dropped.length < count && queue.items.length > 0) + dropped.push(queue.items.shift() as ForwardItem); + for (const item of dropped) queue.bytes -= sizeOf(item); + if (queue.bytes < 0) queue.bytes = 0; + return dropped; +}; diff --git a/src/backend/services/events/forwarding.test.ts b/src/backend/services/events/forwarding.test.ts new file mode 100644 index 000000000..8ea9b5b52 --- /dev/null +++ b/src/backend/services/events/forwarding.test.ts @@ -0,0 +1,999 @@ +/* + * Copyright (C) 2024-present Puter Technologies Inc. + * + * This file is part of Puter. + * + * Puter is free software: you can redistribute it and/or modify + * it under the terms of the GNU Affero General Public License as published + * by the Free Software Foundation, either version 3 of the License, or + * (at your option) any later version. + * + * This program is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + * GNU Affero General Public License for more details. + * + * You should have received a copy of the GNU Affero General Public License + * along with this program. If not, see . + */ + +/** + * Two deployments, one replicated table, and everything that has to hold + * between them. + * + * Each has its own Redis — pending deliveries, leases and connection counts are + * region-local by construction and nothing here replicates them — and shares + * only the presence row, which is the whole of what one deployment tells + * another. Peers are simulated the way the broadcast tests do it: the outbound + * call is intercepted at the transport and handed straight to the other + * instance, so both sides of every exchange are real code. + */ + +import MockRedis from 'ioredis-mock'; +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; +import type { Actor } from '../../core/actor.js'; +import { EventSubscriptionStore } from '../../stores/events/EventSubscriptionStore.js'; +import { PendingDeliveryStore } from '../../stores/events/PendingDeliveryStore.js'; +import { + PresenceStore, + presenceItemKey, + PRESENCE_NO_APP, +} from '../../stores/events/PresenceStore.js'; +import type { + DurableSubscription, + SubscriptionTarget, +} from '../../stores/events/types.js'; +import type { FSEntry } from '../../stores/fs/FSEntry.js'; +import type { IConfig } from '../../types.js'; +import { EventForwardService } from './EventForwardService.js'; +import { EventsService, type DeliveryEnvelope } from './EventsService.js'; +import type { ForwardBatch, ForwardDelivery } from './forwardQueue.js'; +import { fsAnchorToken } from './subjects.js'; +import type { + WorkerInvocation, + WorkerInvocationOutcome, +} from './workerSeam.js'; + +// -- The replicated table -------------------------------------------- + +interface StoredRow { + regions: Record; + version: number; +} + +let table: Map; +let tableReads: number; +let tableWrites: number; + +/** + * The reserved-item path, as the key-value store exposes it. Stubbed at the + * store boundary so both regions share one row; the real path against the real + * table is covered by the presence integration suite. + */ +const kvStub = () => ({ + getReservedItem: async (key: string) => { + tableReads++; + const row = table.get(key); + return row + ? { regions: { ...row.regions }, version: row.version } + : null; + }, + setReservedEntry: async ( + key: string, + _field: string, + entry: string, + value: unknown, + ) => { + tableWrites++; + const row = table.get(key) ?? { regions: {}, version: 0 }; + const next = { + regions: { ...row.regions, [entry]: value as number }, + version: row.version + 1, + }; + table.set(key, next); + return next.version; + }, + removeReservedEntry: async ( + key: string, + _field: string, + entry: string, + expectedVersion: number, + ) => { + tableWrites++; + const row = table.get(key); + if (!row || row.version !== expectedVersion) return false; + const { [entry]: _gone, ...rest } = row.regions; + table.set(key, { regions: rest, version: row.version + 1 }); + return true; + }, +}); + +// -- Regions ---------------------------------------------------------- + +interface Region { + name: string; + redis: InstanceType; + presence: PresenceStore; + pending: PendingDeliveryStore; + subscriptions: EventSubscriptionStore; + forward: EventForwardService; + events: EventsService; + /** Every addressed POST this region made, whether or not it arrived. */ + posts: Array<{ peerId: string; batch: ForwardBatch }>; + /** Rooms this region terminates a socket for. */ + rooms: Set; + sent: DeliveryEnvelope[]; + invoked: WorkerInvocation[]; + alarms: ReturnType; + /** Peers whose POSTs never come back — a timeout, not a refusal. */ + unreachable: Set; + handlers: Map void>; +} + +let regions: Map; +let seq = 0; +let userId = 0; +let rows: Map; +let workerOutcome: WorkerInvocationOutcome; + +const anchorUid = () => `docs-${seq}`; +const anchorPath = () => `/u${userId}/Documents`; +const ancestors = () => [{ uid: anchorUid(), path: anchorPath() }]; + +const entry = (over: Partial = {}): FSEntry => + ({ + uid: `file-${seq}`, + uuid: `file-${seq}`, + path: `${anchorPath()}/notes.txt`, + userId, + isDir: false, + ...over, + }) as FSEntry; + +const actorFor = (appUid: string | null = null): Actor => + ({ + user: { id: userId, uuid: `user-${userId}`, username: `u${userId}` }, + effectiveApp: appUid ? { uid: appUid } : null, + app: appUid ? { uid: appUid } : null, + }) as unknown as Actor; + +const durableRow = ( + over: Partial = {}, +): DurableSubscription => ({ + durable: true, + subId: `sub-${seq}-${over.subId ?? 'one'}`, + holderUserId: userId, + ownerUserId: userId, + subject: `fs:${anchorPath()}`, + token: fsAnchorToken(anchorUid()), + anchorUid: anchorUid(), + anchorPath: anchorPath(), + match: null, + op: null, + appUid: null, + permission: 'list', + delivery: 'broadcast', + targets: ['socket'] as SubscriptionTarget[], + handlerName: null, + context: null, + expiresAt: null, + suspendedAt: null, + suspendedReason: null, + createdAt: Math.floor(Date.now() / 1000), + ...over, +}); + +const makeRegion = ( + name: string, + peers: string[], + config: Partial = {}, +): Region => { + // A keyspace per region, because that is what a region is here: leases, + // pending queues and connection counts are local by construction, and a + // test that quietly shared them would prove the opposite of the claim. + const redis = new MockRedis.Cluster(['redis://localhost:6379'], { + redisOptions: { keyPrefix: `${name}-${seq}:` }, + }); + const region: Region = { + name, + redis, + posts: [], + rooms: new Set(), + sent: [], + invoked: [], + alarms: vi.fn(), + unreachable: new Set(), + handlers: new Map(), + } as unknown as Region; + + const fullConfig = { + events: { enabled: true }, + ...config, + } as IConfig; + + const clients = { + redis, + alarm: { create: region.alarms }, + event: { + on: (key: string, handler: unknown) => { + region.handlers.set( + key, + handler as ( + key: string, + data: unknown, + meta: unknown, + ) => void, + ); + }, + // The broadcast channel, simulated: an `outer.*` emit reaches every + // peer tagged `from_outside`, and never its own emitter. + emit: (key: string, data: unknown, meta: object) => { + for (const other of regions.values()) { + if (other === region) continue; + other.handlers.get(key)?.(key, data, { + ...meta, + from_outside: true, + }); + } + }, + }, + } as never; + + const presence = new PresenceStore(fullConfig, clients, { + kv: kvStub(), + } as never); + const pending = new PendingDeliveryStore(fullConfig, clients, {} as never); + const subscriptions = new EventSubscriptionStore( + fullConfig, + clients, + {} as never, + ); + + const socket = { + has: (spec: { room?: string | number; socket?: string }) => + region.rooms.has(String(spec.room ?? spec.socket ?? '')), + send: vi.fn(async (_spec: unknown, _key: string, data: unknown) => { + region.sent.push(data as DeliveryEnvelope); + }), + }; + + const broadcast = { + get regionId() { + return name; + }, + get addressablePeers() { + return peers; + }, + postToPeer: async (peerId: string, _path: string, payload: unknown) => { + const batch = payload as ForwardBatch; + region.posts.push({ peerId, batch }); + if (region.unreachable.has(peerId)) + throw new Error('peer did not answer'); + const target = regions.get(peerId); + if (!target) throw new Error(`no such peer ${peerId}`); + return target.forward.receive(batch); + }, + }; + + const services: Record = { + broadcast, + socket, + fs: { getAncestorChain: async () => ancestors() }, + acl: { + check: async () => true, + getSafeAclError: async () => ({ + status: 404, + message: 'Subject does not exist', + fields: { code: 'subject_does_not_exist' }, + }), + }, + notification: { notify: vi.fn() }, + metering: { + bufferIncrementUsages: () => undefined, + hasAnyUsageCached: async () => true, + }, + }; + + const stores = { + kv: kvStub(), + presence, + pendingDelivery: pending, + eventSubscription: subscriptions, + durableSubscription: { + warmRegion: async () => false, + getBySubId: async (subId: string) => rows.get(subId) ?? null, + remove: async (row: DurableSubscription) => { + rows.delete(row.subId); + return { userId: row.holderUserId, generation: 1 }; + }, + suspend: async () => [{ userId, generation: 1 }], + }, + fsEntry: { + getEntryByUuid: async (uid: string) => + uid === anchorUid() + ? entry({ uid: anchorUid(), isDir: true }) + : null, + getEntryByPath: async () => null, + getEntryById: async () => null, + }, + user: { getById: async (id: number) => ({ id, uuid: `user-${id}` }) }, + app: { getByUid: async (uid: string) => ({ uid, id: 1 }) }, + permission: { getCacheGeneration: async () => 1 }, + } as never; + + region.presence = presence; + region.pending = pending; + region.subscriptions = subscriptions; + region.forward = new EventForwardService( + fullConfig, + clients, + stores, + services as never, + ); + region.events = new EventsService( + fullConfig, + clients, + stores, + services as never, + ); + services.eventForward = region.forward; + services.events = region.events; + + region.events.worker = { + invoke: async (invocation: WorkerInvocation) => { + region.invoked.push(invocation); + return workerOutcome; + }, + }; + region.forward.onServerStart(); + regions.set(name, region); + return region; +}; + +// -- Helpers ---------------------------------------------------------- + +const rowFor = (appUid: string = PRESENCE_NO_APP): StoredRow | null => + table.get(presenceItemKey(`user-${userId}`, appUid)) ?? null; + +const register = async ( + region: Region, + over: Partial = {}, +): Promise => { + const row = durableRow(over); + rows.set(row.subId, row); + await region.subscriptions.cacheDurable([row]); + region.events.invalidateUser(userId); + return row; +}; + +const dispatch = (region: Region, node = entry()): Promise => + region.events.dispatchFs('fs.write.file', node, { + actingUserId: userId, + ancestors: async () => ancestors(), + }); + +const posted = (region: Region, count = 1): Promise => + vi.waitFor( + () => expect(region.posts.length).toBeGreaterThanOrEqual(count), + { + timeout: 3_000, + interval: 10, + }, + ); + +const arrived = (region: Region, count = 1): Promise => + vi.waitFor(() => expect(region.sent.length).toBeGreaterThanOrEqual(count), { + timeout: 3_000, + interval: 10, + }); + +const quiet = (ms = 150): Promise => + new Promise((resolve) => setTimeout(resolve, ms)); + +/** Move the clock past a lease without waiting out its window. */ +const jump = (ms: number): void => { + vi.useFakeTimers({ toFake: ['Date'] }); + vi.setSystemTime(Date.now() + ms); +}; + +const deliveriesIn = (region: Region): ForwardDelivery[] => + region.posts.flatMap((post) => + post.batch.items.filter( + (item): item is ForwardDelivery => item.kind === 'delivery', + ), + ); + +beforeEach(() => { + seq++; + userId = 9000 + seq; + table = new Map(); + tableReads = 0; + tableWrites = 0; + regions = new Map(); + rows = new Map(); + workerOutcome = 'deferred'; + EventForwardService.LEAVE_DELAY_MIN_MS = 60; + EventForwardService.LEAVE_DELAY_MAX_MS = 60; +}); + +afterEach(async () => { + for (const region of regions.values()) + await region.forward.onServerPrepareShutdown(); + vi.useRealTimers(); +}); + +// -- Presence transitions --------------------------------------------- + +describe('what a connection writes', () => { + it('writes once for the first, and nothing for the reconnects after it', async () => { + const west = makeRegion('west', ['east']); + makeRegion('east', ['west']); + + for (let i = 0; i < 5; i++) await west.forward.noteConnect(actorFor()); + + expect(tableWrites).toBe(1); + expect(rowFor()?.regions).toEqual({ west: expect.any(Number) }); + }); + + it('writes nothing at all for a reload — a disconnect and a reconnect inside the window', async () => { + const west = makeRegion('west', ['east']); + makeRegion('east', ['west']); + + await west.forward.noteConnect(actorFor()); + const afterConnect = tableWrites; + + await west.forward.noteDisconnect(actorFor()); + await west.forward.noteConnect(actorFor()); + await quiet(150); + + expect(tableWrites).toBe(afterConnect); + expect(rowFor()?.regions).toEqual({ west: expect.any(Number) }); + }); + + it('takes the region out of the row once the last connection stays gone', async () => { + const west = makeRegion('west', ['east']); + makeRegion('east', ['west']); + + await west.forward.noteConnect(actorFor()); + await west.forward.noteDisconnect(actorFor()); + + await vi.waitFor(() => expect(rowFor()?.regions).toEqual({}), { + timeout: 2_000, + interval: 10, + }); + }); + + it('holds the row while any connection remains', async () => { + const west = makeRegion('west', ['east']); + makeRegion('east', ['west']); + + await west.forward.noteConnect(actorFor()); + await west.forward.noteConnect(actorFor()); + await west.forward.noteDisconnect(actorFor()); + await quiet(150); + + expect(rowFor()?.regions).toEqual({ west: expect.any(Number) }); + }); + + it('never writes on a timer — a session sits connected across many windows', async () => { + const west = makeRegion('west', ['east']); + makeRegion('east', ['west']); + await west.forward.noteConnect(actorFor()); + expect(tableWrites).toBe(1); + + vi.useFakeTimers(); + // Several times over any refresh interval anyone would have chosen. + await vi.advanceTimersByTimeAsync(30 * 60 * 1000); + vi.useRealTimers(); + + expect(tableWrites).toBe(1); + expect(rowFor()?.regions).toEqual({ west: expect.any(Number) }); + }); + + it('skips the leaving write under a drain, where every socket goes at once', async () => { + const west = makeRegion('west', ['east']); + makeRegion('east', ['west']); + await west.forward.noteConnect(actorFor()); + + await west.forward.onServerPrepareShutdown(); + await west.forward.noteDisconnect(actorFor()); + await quiet(150); + + expect(tableWrites).toBe(1); + expect(rowFor()?.regions).toEqual({ west: expect.any(Number) }); + }); + + it('counts each app separately, because each has its own room', async () => { + const west = makeRegion('west', ['east']); + makeRegion('east', ['west']); + + await west.forward.noteConnect(actorFor('app-a')); + await west.forward.noteConnect(actorFor('app-b')); + + expect(rowFor('app-a')?.regions).toEqual({ west: expect.any(Number) }); + expect(rowFor('app-b')?.regions).toEqual({ west: expect.any(Number) }); + }); + + it('moves with the connection itself, not with what it subscribes to', async () => { + const west = makeRegion('west', ['east']); + makeRegion('east', ['west']); + const handlers = new Map void>(); + const socket = { + id: 'socket-1', + on: () => undefined, + once: (event: string, handler: () => void) => + handlers.set(event, handler), + }; + + west.events.attachSocket(socket, actorFor()); + await vi.waitFor(() => + expect(rowFor()?.regions).toEqual({ west: expect.any(Number) }), + ); + + handlers.get('disconnect')?.(); + await vi.waitFor(() => expect(rowFor()?.regions).toEqual({}), { + timeout: 2_000, + interval: 10, + }); + }); + + it('takes no part in presence where nothing is configured to read it', async () => { + const alone = makeRegion('west', []); + + await alone.forward.noteConnect(actorFor()); + await alone.forward.noteDisconnect(actorFor()); + await quiet(150); + + expect(tableWrites).toBe(0); + expect(tableReads).toBe(0); + expect(table.size).toBe(0); + }); + + it('writes nothing where events are switched off', async () => { + const off = makeRegion('west', ['east'], { + events: { enabled: false }, + } as Partial); + + await off.forward.noteConnect(actorFor()); + + expect(tableWrites).toBe(0); + }); +}); + +// -- Forwarding ------------------------------------------------------- + +describe('a broadcast delivery', () => { + it('makes no cross-region traffic at all when nobody is connected anywhere', async () => { + const west = makeRegion('west', ['east']); + makeRegion('east', ['west']); + await register(west); + + await dispatch(west); + await quiet(200); + + expect(west.posts).toEqual([]); + }); + + it('goes to the one region holding the socket, and only that one', async () => { + const west = makeRegion('west', ['east', 'south']); + const east = makeRegion('east', ['west', 'south']); + makeRegion('south', ['west', 'east']); + east.rooms.add(String(userId)); + await east.forward.noteConnect(actorFor()); + await register(west); + + await dispatch(west); + await posted(west); + await arrived(east); + + expect(west.posts).toHaveLength(1); + expect(west.posts[0].peerId).toBe('east'); + expect(east.sent[0].event).toMatchObject({ op: 'write' }); + }); + + it('stays put for a subscription tied to one connection, which is here', async () => { + const west = makeRegion('west', ['east']); + const east = makeRegion('east', ['west']); + east.rooms.add(String(userId)); + await east.forward.noteConnect(actorFor()); + // A session row is addressed at the socket that made it, so another + // region's tabs are not among its subscribers. + await register(west, { socketId: 'socket-here', durable: undefined }); + + await dispatch(west); + await quiet(200); + + expect(west.posts).toEqual([]); + }); + + it('reaches every region a subscriber is connected in', async () => { + const west = makeRegion('west', ['east', 'south']); + const east = makeRegion('east', ['west', 'south']); + const south = makeRegion('south', ['west', 'east']); + east.rooms.add(String(userId)); + south.rooms.add(String(userId)); + await east.forward.noteConnect(actorFor()); + await south.forward.noteConnect(actorFor()); + await register(west); + + await dispatch(west); + await posted(west, 2); + + expect(west.posts.map((post) => post.peerId).sort()).toEqual([ + 'east', + 'south', + ]); + }); + + it('carries a window of events in one request, not one request each', async () => { + const west = makeRegion('west', ['east']); + const east = makeRegion('east', ['west']); + east.rooms.add(String(userId)); + await east.forward.noteConnect(actorFor()); + await register(west); + + // Distinct nodes, so nothing is coalesced away upstream. + for (let i = 0; i < 6; i++) + await dispatch( + west, + entry({ uid: `file-${i}`, path: `${anchorPath()}/n${i}.txt` }), + ); + await posted(west); + await quiet(200); + + expect(west.posts).toHaveLength(1); + expect(west.posts[0].batch.items).toHaveLength(6); + }); + + it('reads presence once and then answers from what it already knows', async () => { + const west = makeRegion('west', ['east']); + const east = makeRegion('east', ['west']); + east.rooms.add(String(userId)); + await east.forward.noteConnect(actorFor()); + await register(west); + + tableReads = 0; + for (let i = 0; i < 8; i++) + await dispatch( + west, + entry({ uid: `file-${i}`, path: `${anchorPath()}/n${i}.txt` }), + ); + await posted(west); + await quiet(200); + + expect(tableReads).toBe(1); + }); +}); + +// -- Lazy repair ------------------------------------------------------ + +describe('a row naming a region that holds nothing', () => { + it('is corrected once, on the region saying so itself', async () => { + const west = makeRegion('west', ['east']); + const east = makeRegion('east', ['west']); + await east.forward.noteConnect(actorFor()); + // The socket goes without east ever reporting it — a node that died. + east.rooms.clear(); + await east.presence.removeConnection(userId, PRESENCE_NO_APP); + await register(west); + + tableWrites = 0; + await dispatch(west); + await posted(west); + + await vi.waitFor(() => expect(rowFor()?.regions).toEqual({}), { + timeout: 2_000, + interval: 10, + }); + expect(tableWrites).toBe(1); + }); + + it('is left alone when the region simply did not answer', async () => { + const west = makeRegion('west', ['east']); + const east = makeRegion('east', ['west']); + await east.forward.noteConnect(actorFor()); + east.rooms.clear(); + await east.presence.removeConnection(userId, PRESENCE_NO_APP); + west.unreachable.add('east'); + await register(west); + + tableWrites = 0; + await dispatch(west); + await posted(west); + await quiet(200); + + // A timeout says nothing about whether the socket is there. + expect(tableWrites).toBe(0); + expect(rowFor()?.regions).toEqual({ east: expect.any(Number) }); + }); + + it('cannot be stormed by a busy subscription', async () => { + const west = makeRegion('west', ['east']); + const east = makeRegion('east', ['west']); + await east.forward.noteConnect(actorFor()); + east.rooms.clear(); + await east.presence.removeConnection(userId, PRESENCE_NO_APP); + await register(west); + + tableWrites = 0; + for (let i = 0; i < 10; i++) { + await dispatch( + west, + entry({ uid: `file-${i}`, path: `${anchorPath()}/n${i}.txt` }), + ); + await quiet(60); + } + await quiet(200); + + expect(tableWrites).toBe(1); + }); + + it('loses harmlessly to a connect that landed while it was deciding', async () => { + const west = makeRegion('west', ['east']); + const east = makeRegion('east', ['west']); + await east.forward.noteConnect(actorFor()); + east.rooms.clear(); + await east.presence.removeConnection(userId, PRESENCE_NO_APP); + await register(west); + + // West reads the row, and only then does east get a socket back. + await west.forward.regionsFor(userId, null); + east.rooms.add(String(userId)); + await east.forward.noteConnect(actorFor()); + const fresh = rowFor()?.version; + + await dispatch(west); + await posted(west); + await quiet(200); + + // The stale repair is refused, so the reconnected region is still there. + expect(rowFor()?.regions).toEqual({ east: expect.any(Number) }); + expect(rowFor()?.version).toBe(fresh); + }); + + it('is not repaired away while the region holding it is still inside its own disconnect window', async () => { + EventForwardService.LEAVE_DELAY_MIN_MS = 10_000; + EventForwardService.LEAVE_DELAY_MAX_MS = 10_000; + try { + const west = makeRegion('west', ['east']); + const east = makeRegion('east', ['west']); + east.rooms.add(String(userId)); + await east.forward.noteConnect(actorFor()); + await register(west); + + // The socket is gone and the disconnect is in flight, but still + // inside its window: east expects the pair back, and has not + // written itself out of the row for it. + east.rooms.delete(String(userId)); + await east.forward.noteDisconnect(actorFor()); + + tableWrites = 0; + await dispatch(west); + await posted(west); + await arrived(east); + + expect(tableWrites).toBe(0); + expect(rowFor()?.regions).toEqual({ east: expect.any(Number) }); + } finally { + EventForwardService.LEAVE_DELAY_MIN_MS = 60; + EventForwardService.LEAVE_DELAY_MAX_MS = 60; + } + }); +}); + +// -- `single` across regions ------------------------------------------ + +describe('a delivery owed to exactly one consumer', () => { + const single = (region: Region) => + register(region, { + delivery: 'single', + targets: ['socket', 'worker'] as SubscriptionTarget[], + handlerName: 'onWrite', + }); + + it('takes the emitting region`s own socket before any other', async () => { + const west = makeRegion('west', ['east']); + const east = makeRegion('east', ['west']); + west.rooms.add(String(userId)); + east.rooms.add(String(userId)); + await west.forward.noteConnect(actorFor()); + await east.forward.noteConnect(actorFor()); + await single(west); + + await dispatch(west); + await arrived(west); + + expect(west.posts).toEqual([]); + expect(west.sent[0]).toMatchObject({ ackRequired: true }); + }); + + it('hands off to the most recently connected region when it holds none', async () => { + const west = makeRegion('west', ['east', 'south']); + const east = makeRegion('east', ['west', 'south']); + const south = makeRegion('south', ['west', 'east']); + east.rooms.add(String(userId)); + south.rooms.add(String(userId)); + await east.forward.noteConnect(actorFor()); + await quiet(5); + await south.forward.noteConnect(actorFor()); + const row = await single(west); + + await dispatch(west); + await posted(west); + await arrived(south); + + expect(west.posts[0].peerId).toBe('south'); + const carried = deliveriesIn(west)[0]; + expect(carried).toMatchObject({ + subId: row.subId, + ackRequired: true, + origin: 'west', + }); + expect(typeof carried.ackId).toBe('string'); + expect(south.sent[0]).toMatchObject({ + ackRequired: true, + origin: 'west', + }); + }); + + it('keeps every trace of the delivery in the region that emitted it', async () => { + const west = makeRegion('west', ['east']); + const east = makeRegion('east', ['west']); + east.rooms.add(String(userId)); + await east.forward.noteConnect(actorFor()); + const row = await single(west); + + await dispatch(west); + await posted(west); + await arrived(east); + + const held = async (region: Region) => + (await region.redis.keys('*')).filter((key: string) => + key.includes(row.subId), + ); + expect((await held(west)).length).toBeGreaterThan(0); + expect(await held(east)).toEqual([]); + }); + + it('settles on an ack the client gave to whichever region it reached', async () => { + const west = makeRegion('west', ['east']); + const east = makeRegion('east', ['west']); + east.rooms.add(String(userId)); + await east.forward.noteConnect(actorFor()); + const row = await single(west); + + await dispatch(west); + await posted(west); + await arrived(east); + + const envelope = east.sent[0]; + await east.events.ackDelivery(actorFor(), { + subId: row.subId, + id: envelope.ackId, + origin: envelope.origin, + }); + // The relay rides the same addressed channel the delivery did. + await vi.waitFor( + async () => expect(await west.pending.depth(row.subId)).toBe(0), + { timeout: 3_000, interval: 20 }, + ); + expect(east.posts.at(-1)?.batch.items[0]).toMatchObject({ + kind: 'ack', + subId: row.subId, + }); + }); + + it('never settles a lease in the region the client happened to reach', async () => { + const west = makeRegion('west', ['east']); + const east = makeRegion('east', ['west']); + east.rooms.add(String(userId)); + await east.forward.noteConnect(actorFor()); + const row = await single(west); + + await dispatch(west); + await posted(west); + await arrived(east); + + await east.events.ackDelivery(actorFor(), { + subId: row.subId, + id: east.sent[0].ackId, + origin: 'west', + }); + + expect(await east.pending.depth(row.subId)).toBe(0); + }); + + it('gives the handler what two socket candidates could not take', async () => { + const west = makeRegion('west', ['east']); + const east = makeRegion('east', ['west']); + east.rooms.add(String(userId)); + await east.forward.noteConnect(actorFor()); + workerOutcome = 'settled'; + const row = await single(west); + + await dispatch(west); + await posted(west); + + // Neither candidate acked, so the leases lapse and the handler takes + // what the sockets would not. + for (let attempt = 0; attempt < 3; attempt++) { + jump(31_000); + await west.events.sweepPending(); + } + + expect(west.invoked.length).toBeGreaterThan(0); + expect(west.invoked[0].subId).toBe(row.subId); + }); +}); + +// -- Overflow --------------------------------------------------------- + +describe('a forward queue that cannot keep up', () => { + it('replaces what it shed with a marker, and says so', async () => { + EventForwardService.MAX_QUEUED = 2; + try { + const west = makeRegion('west', ['east']); + const east = makeRegion('east', ['west']); + east.rooms.add(String(userId)); + await east.forward.noteConnect(actorFor()); + const row = await register(west); + + for (let i = 0; i < 6; i++) + west.forward.handOff('east', { + holderUserId: userId, + appUid: null, + subId: row.subId, + event: { + id: `ev-${i}`, + subject: `fs:${anchorPath()}`, + op: 'write', + uid: `file-${i}`, + path: `${anchorPath()}/n${i}.txt`, + self: true, + seq: 0, + ts: Date.now(), + }, + }); + + await posted(west); + await arrived(east, 2); + + expect(west.alarms).toHaveBeenCalledWith( + 'events_forward_overflow', + expect.any(String), + expect.objectContaining({ peerId: 'east' }), + 'warning', + expect.anything(), + ); + // One marker for the one subscription that lost events, with the + // newest delivery still behind it: a shed never eats the queue. + expect( + east.sent.filter((envelope) => envelope.event.op === 'gap'), + ).toHaveLength(1); + expect( + east.sent.some((envelope) => envelope.event.id === 'ev-5'), + ).toBe(true); + } finally { + EventForwardService.MAX_QUEUED = 5_000; + } + }); +}); + +// -- Generation propagation ------------------------------------------- + +describe('presence moving in another region', () => { + it('invalidates what this one cached about it', async () => { + const west = makeRegion('west', ['east']); + const east = makeRegion('east', ['west']); + await register(west); + + expect(await west.forward.regionsFor(userId, null)).toEqual([]); + const reads = tableReads; + + east.rooms.add(String(userId)); + await east.forward.noteConnect(actorFor()); + + expect(await west.forward.regionsFor(userId, null)).toEqual(['east']); + expect(tableReads).toBeGreaterThan(reads); + }); +}); diff --git a/src/backend/services/events/metering.test.ts b/src/backend/services/events/metering.test.ts index fdbcdb2f6..87b00324b 100644 --- a/src/backend/services/events/metering.test.ts +++ b/src/backend/services/events/metering.test.ts @@ -185,6 +185,18 @@ beforeEach(() => { permission: { getCacheGeneration: async () => 1 }, } as never, { + eventForward: { + // A deployment with no peers has nowhere to forward to, which + // is what every test here is. + region: 'local', + isPeer: () => false, + noteConnect: async () => undefined, + noteDisconnect: async () => undefined, + candidateRegion: async () => null, + fanOut: async () => undefined, + handOff: () => undefined, + relayAck: () => undefined, + }, socket: { send: vi.fn(), has: () => false }, fs: { getAncestorChain: async () => [] }, acl: { diff --git a/src/backend/services/events/presenceCache.test.ts b/src/backend/services/events/presenceCache.test.ts new file mode 100644 index 000000000..857c2547a --- /dev/null +++ b/src/backend/services/events/presenceCache.test.ts @@ -0,0 +1,174 @@ +/* + * Copyright (C) 2024-present Puter Technologies Inc. + * + * This file is part of Puter. + * + * Puter is free software: you can redistribute it and/or modify + * it under the terms of the GNU Affero General Public License as published + * by the Free Software Foundation, either version 3 of the License, or + * (at your option) any later version. + * + * This program is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + * GNU Affero General Public License for more details. + * + * You should have received a copy of the GNU Affero General Public License + * along with this program. If not, see . + */ + +import { describe, expect, it } from 'vitest'; +import type { PresenceRow } from '../../stores/events/PresenceStore.js'; +import { PresenceCache, remoteRegions } from './presenceCache.js'; + +/** + * What the cache promises: an answer survives until presence moves, and one + * window buys one repair. Both are what keep reads and corrective writes + * proportional to session churn instead of to how busy a subscription is. + */ + +const row = (regions: Record, version = 1): PresenceRow => ({ + regions, + version, +}); + +describe('the presence cache', () => { + it('answers from memory until something bumps it', () => { + const cache = new PresenceCache(); + cache.write(7, 'app-a', cache.generationOf(7), row({ west: 10 })); + + expect(cache.read(7, 'app-a')).toEqual(row({ west: 10 })); + expect(cache.read(7, 'app-a')).toEqual(row({ west: 10 })); + + cache.bump(7); + expect(cache.read(7, 'app-a')).toBeNull(); + }); + + it('has no expiry — a settled row is never re-read on its own', () => { + const cache = new PresenceCache(); + cache.write(7, 'app-a', cache.generationOf(7), row({ west: 10 })); + + const later = Date.now() + 60 * 60 * 1000; + const realNow = Date.now; + Date.now = () => later; + try { + expect(cache.read(7, 'app-a')).not.toBeNull(); + } finally { + Date.now = realNow; + } + }); + + it('drops an answer computed against a generation that has since moved', () => { + const cache = new PresenceCache(); + const epoch = cache.generationOf(7); + cache.bump(7); + cache.write(7, 'app-a', epoch, row({ west: 10 })); + + expect(cache.read(7, 'app-a')).toBeNull(); + }); + + it('bumps one user without touching another', () => { + const cache = new PresenceCache(); + cache.write(7, 'app-a', cache.generationOf(7), row({ west: 10 })); + cache.write(8, 'app-a', cache.generationOf(8), row({ east: 11 })); + + cache.bump(7); + + expect(cache.read(7, 'app-a')).toBeNull(); + expect(cache.read(8, 'app-a')).toEqual(row({ east: 11 })); + }); + + it('ignores a numbered bump it has already applied', () => { + const cache = new PresenceCache(); + cache.bump(7, 5); + cache.write(7, 'app-a', cache.generationOf(7), row({ west: 10 })); + + cache.bump(7, 4); + expect(cache.read(7, 'app-a')).not.toBeNull(); + + cache.bump(7, 6); + expect(cache.read(7, 'app-a')).toBeNull(); + }); + + it('advances on an unnumbered bump whatever number it last saw', () => { + const cache = new PresenceCache(); + cache.bump(7, 9); + cache.write(7, 'app-a', cache.generationOf(7), row({ west: 10 })); + + // A peer's counter cannot be compared against this region's own. + cache.bump(7); + expect(cache.read(7, 'app-a')).toBeNull(); + }); + + it('gives one window one repair per region', () => { + const cache = new PresenceCache(); + cache.write(7, 'app-a', cache.generationOf(7), row({ east: 10 })); + + expect(cache.claimRepair(7, 'app-a', 'east')).toBe(true); + expect(cache.claimRepair(7, 'app-a', 'east')).toBe(false); + expect(cache.claimRepair(7, 'app-a', 'east')).toBe(false); + // A different region is its own claim. + expect(cache.claimRepair(7, 'app-a', 'south')).toBe(true); + }); + + it('lets the next window repair again', () => { + const cache = new PresenceCache(); + cache.write(7, 'app-a', cache.generationOf(7), row({ east: 10 })); + expect(cache.claimRepair(7, 'app-a', 'east')).toBe(true); + + cache.bump(7); + cache.write(7, 'app-a', cache.generationOf(7), row({ east: 10 })); + + expect(cache.claimRepair(7, 'app-a', 'east')).toBe(true); + }); + + it('refuses a repair for something it holds no row for', () => { + const cache = new PresenceCache(); + expect(cache.claimRepair(7, 'app-a', 'east')).toBe(false); + }); + + it('stops naming a region the moment one is repaired away', () => { + const cache = new PresenceCache(); + cache.write( + 7, + 'app-a', + cache.generationOf(7), + row({ east: 10, south: 11 }), + ); + + cache.forget(7, 'app-a', 'east'); + + expect(cache.read(7, 'app-a')?.regions).toEqual({ south: 11 }); + }); + + it('evicts least-recently-read once it is full', () => { + const cache = new PresenceCache(2); + cache.write(1, 'a', 0, row({ west: 1 })); + cache.write(2, 'a', 0, row({ west: 2 })); + cache.read(1, 'a'); + cache.write(3, 'a', 0, row({ west: 3 })); + + expect(cache.size).toBe(2); + expect(cache.read(2, 'a')).toBeNull(); + expect(cache.read(1, 'a')).not.toBeNull(); + }); +}); + +describe('candidate regions', () => { + it('leaves out the region asking, which never consults presence about itself', () => { + expect(remoteRegions(row({ west: 1, east: 2 }), 'west')).toEqual([ + 'east', + ]); + }); + + it('offers the most recently connected first', () => { + expect( + remoteRegions(row({ east: 10, south: 30, north: 20 }), 'west'), + ).toEqual(['south', 'north', 'east']); + }); + + it('is empty for a row naming nobody else', () => { + expect(remoteRegions(row({ west: 1 }), 'west')).toEqual([]); + expect(remoteRegions(row({}), 'west')).toEqual([]); + }); +}); diff --git a/src/backend/services/events/presenceCache.ts b/src/backend/services/events/presenceCache.ts new file mode 100644 index 000000000..d8278c547 --- /dev/null +++ b/src/backend/services/events/presenceCache.ts @@ -0,0 +1,178 @@ +/* + * Copyright (C) 2024-present Puter Technologies Inc. + * + * This file is part of Puter. + * + * Puter is free software: you can redistribute it and/or modify + * it under the terms of the GNU Affero General Public License as published + * by the Free Software Foundation, either version 3 of the License, or + * (at your option) any later version. + * + * This program is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + * GNU Affero General Public License for more details. + * + * You should have received a copy of the GNU Affero General Public License + * along with this program. If not, see . + */ + +import type { PresenceRow } from '../../stores/events/PresenceStore.js'; + +/** + * Where a user's sockets are, per region, held until presence actually moves. + * + * Keyed by a per-user generation and by nothing else: no expiry. A fixed + * lifetime would make presence _reads_ scale with matched event volume, which + * is the one thing this whole path exists to avoid — a busy subscription + * against a settled row must read the table once and never again. The + * generation is bumped on every transition and every repair and carried to peer + * regions over the broadcast channel, so staleness ends when presence changes + * rather than on a clock. + * + * The entry also carries which regions it has already asked to be repaired, + * which is what stops a busy stream from storming conditional writes at a row + * it believes is wrong. A bump clears the entry, so the next window may repair + * again. + * + * Bounded, and evicted least-recently-used: a `Map` iterates in insertion + * order, so re-inserting on read moves an entry to the young end. + */ + +interface CacheEntry { + userId: number; + epoch: number; + row: PresenceRow; + /** Regions this window has already spent its one repair write on. */ + repaired: Set; +} + +interface UserEntry { + epoch: number; + /** `null` until a local transition has supplied a real counter value. */ + redisGeneration: number | null; +} + +export const PRESENCE_CACHE_MAX_ENTRIES = 10_000; + +const entryKey = (userId: number, appUid: string): string => + `${userId}|${appUid}`; + +export class PresenceCache { + readonly #entries = new Map(); + readonly #users = new Map(); + readonly #maxEntries: number; + + constructor(maxEntries: number = PRESENCE_CACHE_MAX_ENTRIES) { + this.#maxEntries = Math.max(1, maxEntries); + } + + get size(): number { + return this.#entries.size; + } + + /** The epoch a read must capture before going to the table. */ + generationOf(userId: number): number { + return this.#users.get(userId)?.epoch ?? 0; + } + + /** The cached row, or `null` when this region has to go and look. */ + read(userId: number, appUid: string): PresenceRow | null { + const key = entryKey(userId, appUid); + const entry = this.#entries.get(key); + if (!entry) return null; + if (entry.epoch !== this.generationOf(userId)) { + this.#entries.delete(key); + return null; + } + this.#entries.delete(key); + this.#entries.set(key, entry); + return entry.row; + } + + /** + * Record a row against the epoch it was read under. A bump that landed + * mid-read leaves the epochs mismatched and the answer is dropped rather + * than cached stale. + */ + write( + userId: number, + appUid: string, + epoch: number, + row: PresenceRow, + ): void { + if (this.generationOf(userId) !== epoch) return; + this.#set(entryKey(userId, appUid), { + userId, + epoch, + row, + repaired: new Set(), + }); + } + + /** + * Invalidate everything cached for a user. With `generation`, this is a + * local transition reporting the counter's new value, applied only when it + * is ahead of the last one recorded so two racing bumps cannot land out of + * order. Without one — every signal from another region — the epoch + * advances unconditionally: the counter is region-local, so a peer's number + * cannot be compared against this region's own. + */ + bump(userId: number, generation?: number): void { + const user = this.#users.get(userId); + if ( + generation !== undefined && + generation <= (user?.redisGeneration ?? 0) + ) + return; + this.#users.set(userId, { + epoch: (user?.epoch ?? 0) + 1, + redisGeneration: generation ?? user?.redisGeneration ?? null, + }); + } + + /** + * Take this window's one repair for a (user, app, region), or refuse it + * because the window already spent it. + */ + claimRepair(userId: number, appUid: string, region: string): boolean { + const entry = this.#entries.get(entryKey(userId, appUid)); + if (!entry || entry.epoch !== this.generationOf(userId)) return false; + if (entry.repaired.has(region)) return false; + entry.repaired.add(region); + return true; + } + + /** Drop a region from the cached row, so the window stops forwarding to it. */ + forget(userId: number, appUid: string, region: string): void { + const entry = this.#entries.get(entryKey(userId, appUid)); + if (!entry) return; + const { [region]: _dropped, ...rest } = entry.row.regions; + entry.row = { ...entry.row, regions: rest }; + } + + clear(): void { + this.#entries.clear(); + this.#users.clear(); + } + + #set(key: string, entry: CacheEntry): void { + this.#entries.delete(key); + this.#entries.set(key, entry); + while (this.#entries.size > this.#maxEntries) { + const oldest = this.#entries.keys().next(); + if (oldest.done) break; + this.#entries.delete(oldest.value); + } + } +} + +/** + * Regions to try for one delivery, most recently connected first, with the + * emitting region — which never consults presence about itself — dropped. + */ +export const remoteRegions = (row: PresenceRow, self: string): string[] => + Object.entries(row.regions ?? {}) + .filter(([region]) => region !== self) + .sort(([, left], [, right]) => Number(right) - Number(left)) + .map(([region]) => region); diff --git a/src/backend/services/events/singleDelivery.test.ts b/src/backend/services/events/singleDelivery.test.ts index 80f8ad8cc..18dc14782 100644 --- a/src/backend/services/events/singleDelivery.test.ts +++ b/src/backend/services/events/singleDelivery.test.ts @@ -248,6 +248,18 @@ beforeEach(async () => { permission: { getCacheGeneration: async () => 1 }, } as never, { + eventForward: { + // A deployment with no peers has nowhere to forward to, which + // is what every test here is. + region: 'local', + isPeer: () => false, + noteConnect: async () => undefined, + noteDisconnect: async () => undefined, + candidateRegion: async () => null, + fanOut: async () => undefined, + handOff: () => undefined, + relayAck: () => undefined, + }, socket: { send: vi.fn(async (_spec, _key, data) => { sent.push(data as DeliveryEnvelope); diff --git a/src/backend/services/index.ts b/src/backend/services/index.ts index e769bce03..345404ef6 100644 --- a/src/backend/services/index.ts +++ b/src/backend/services/index.ts @@ -28,6 +28,7 @@ import { OIDCService } from './auth/OIDCService'; import { TokenService } from './auth/TokenService'; import { BroadcastService } from './broadcast/BroadcastService'; import { CacheReplicationService } from './cache/CacheReplicationService'; +import { EventForwardService } from './events/EventForwardService'; import { EventsService } from './events/EventsService'; import { AppFeedbackService } from './feedback/AppFeedbackService'; import { FSService } from './fs/FSService'; @@ -70,6 +71,7 @@ declare module './types' { suggestedApps: SuggestedAppsService; socket: SocketService; events: EventsService; + eventForward: EventForwardService; notification: NotificationService; appFeedback: AppFeedbackService; broadcast: BroadcastService; @@ -131,6 +133,9 @@ export const puterServices = { // AuthService.appUidFromOrigin). appFeedback: AppFeedbackService, broadcast: BroadcastService, + // Forwards through `broadcast` and puts deliveries down through `socket`, + // so it follows both; `events` reaches it at call time only. + eventForward: EventForwardService, // Independent — only needs the event client and redis. cacheReplication: CacheReplicationService, oidc: OIDCService, diff --git a/src/backend/stores/events/PresenceStore.integration.test.ts b/src/backend/stores/events/PresenceStore.integration.test.ts new file mode 100644 index 000000000..d8ab04005 --- /dev/null +++ b/src/backend/stores/events/PresenceStore.integration.test.ts @@ -0,0 +1,206 @@ +/* + * Copyright (C) 2024-present Puter Technologies Inc. + * + * This file is part of Puter. + * + * Puter is free software: you can redistribute it and/or modify + * it under the terms of the GNU Affero General Public License as published + * by the Free Software Foundation, either version 3 of the License, or + * (at your option) any later version. + * + * This program is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + * GNU Affero General Public License for more details. + * + * You should have received a copy of the GNU Affero General Public License + * along with this program. If not, see . + */ + +/** + * Presence rows against the real table. + * + * Two things are on the hook. The first is that these rows go in through the + * direct item path and nothing else: they are not a user's data, so nothing + * about them may be metered, listed, cached, invalidated, or announced as a + * key-value change — a subscription firing on a presence write would be a + * feedback loop. The second is that the version really is a compare-and-set: + * every write but the first depends on it, and that is the whole of what keeps + * a reconnect from being erased by a repair that read the row before it. + */ + +import { afterAll, beforeAll, beforeEach, describe, expect, it, vi } from 'vitest'; +import { setupTestServer } from '../../testUtil.js'; +import type { PuterServer } from '../../server.js'; +import type { IConfig } from '../../types.js'; +import { presenceItemKey } from './PresenceStore.js'; + +const BOOT_TIMEOUT_MS = 120_000; + +let server: PuterServer; +let emitted: string[]; +let seq = 0; + +const presence = () => server.stores.presence; +const userUuid = () => `presence-user-${seq}`; +const appUid = () => `presence-app-${seq}`; + +beforeAll(async () => { + server = await setupTestServer({ events: { enabled: true } } as IConfig); +}, BOOT_TIMEOUT_MS); + +beforeEach(() => { + seq++; + emitted = []; + vi.spyOn(server.clients.event, 'emit').mockImplementation((( + key: string, + ...rest: unknown[] + ) => { + emitted.push(key); + return ( + server.clients.event.constructor.prototype.emit as never + ) as never; + }) as never); +}); + +afterAll(async () => { + vi.restoreAllMocks(); + await server?.shutdown(); +}); + +describe('a presence row', () => { + it('is created by the first region to claim it, and read back whole', async () => { + const version = await presence().join(userUuid(), appUid(), 'west', 11); + + expect(version).toBe(1); + expect(await presence().read(userUuid(), appUid())).toEqual({ + regions: { west: 11 }, + version: 1, + }); + }); + + it('reads as empty when nobody has ever claimed it', async () => { + expect(await presence().read(userUuid(), appUid())).toEqual({ + regions: {}, + version: 0, + }); + }); + + it('keeps both regions when two claim it, and moves the version each time', async () => { + await presence().join(userUuid(), appUid(), 'west', 11); + await presence().join(userUuid(), appUid(), 'east', 22); + + expect(await presence().read(userUuid(), appUid())).toEqual({ + regions: { west: 11, east: 22 }, + version: 2, + }); + }); + + it('takes one region out and leaves the other where it was', async () => { + await presence().join(userUuid(), appUid(), 'west', 11); + await presence().join(userUuid(), appUid(), 'east', 22); + const row = await presence().read(userUuid(), appUid()); + + expect( + await presence().leave(userUuid(), appUid(), 'east', row.version), + ).toBe(true); + expect(await presence().read(userUuid(), appUid())).toEqual({ + regions: { west: 11 }, + version: 3, + }); + }); + + it('refuses a removal that read the row before it moved', async () => { + await presence().join(userUuid(), appUid(), 'west', 11); + const stale = await presence().read(userUuid(), appUid()); + // The socket comes back in that region while the removal is deciding. + await presence().join(userUuid(), appUid(), 'west', 33); + + expect( + await presence().leave(userUuid(), appUid(), 'west', stale.version), + ).toBe(false); + expect( + (await presence().read(userUuid(), appUid())).regions, + ).toEqual({ west: 33 }); + }); + + it('refuses a removal against a row that does not exist', async () => { + expect( + await presence().leave(userUuid(), appUid(), 'west', 0), + ).toBe(false); + }); + + it('says nothing on the event bus — not a mutation, not an invalidation', async () => { + await presence().join(userUuid(), appUid(), 'west', 11); + const row = await presence().read(userUuid(), appUid()); + await presence().leave(userUuid(), appUid(), 'west', row.version); + + expect(emitted).not.toContain('kv.mutated'); + expect(emitted).not.toContain('kv.flushed'); + expect(emitted).not.toContain('outer.kv.cacheInvalidated'); + }); + + it('cannot be reached through the key-value surface', async () => { + await presence().join(userUuid(), appUid(), 'west', 11); + const actor = { + user: { id: 1, uuid: userUuid(), username: 'presence' }, + app: { uid: appUid() }, + effectiveApp: { uid: appUid() }, + }; + + const read = await server.stores.kv.get( + { key: presenceItemKey(userUuid(), appUid()) }, + { actor: actor as never }, + ); + expect(read.res).toBeNull(); + + const listed = await server.stores.kv.list( + {}, + { actor: actor as never }, + ); + expect(JSON.stringify(listed.res)).not.toContain('pr#'); + }); +}); + +describe('this region`s connection count', () => { + it('crosses zero once, however many connections come and go', async () => { + const store = presence(); + expect(await store.addConnection(seq, appUid())).toBe(1); + expect(await store.addConnection(seq, appUid())).toBe(2); + expect(await store.addConnection(seq, appUid())).toBe(3); + + expect(await store.removeConnection(seq, appUid())).toBe(2); + expect(await store.removeConnection(seq, appUid())).toBe(1); + expect(await store.removeConnection(seq, appUid())).toBe(0); + }); + + it('keeps nothing behind once the last connection goes', async () => { + const store = presence(); + await store.addConnection(seq, appUid()); + await store.removeConnection(seq, appUid()); + + expect(await store.holdsConnection(seq, appUid())).toBe(false); + // A double reap must not drive it below zero and hide the next connect. + expect(await store.removeConnection(seq, appUid())).toBe(0); + expect(await store.addConnection(seq, appUid())).toBe(1); + }); + + it('counts each app of one user separately', async () => { + const store = presence(); + await store.addConnection(seq, appUid()); + + expect(await store.holdsConnection(seq, `${appUid()}-other`)).toBe( + false, + ); + expect(await store.holdsConnection(seq, appUid())).toBe(true); + }); +}); + +describe('the presence generation', () => { + it('moves on every transition, so nothing cached under the old one is used', async () => { + const first = await presence().bumpGeneration(seq); + const second = await presence().bumpGeneration(seq); + + expect(second).toBeGreaterThan(first); + }); +}); diff --git a/src/backend/stores/events/PresenceStore.ts b/src/backend/stores/events/PresenceStore.ts new file mode 100644 index 000000000..2f7cccb9f --- /dev/null +++ b/src/backend/stores/events/PresenceStore.ts @@ -0,0 +1,197 @@ +/* + * Copyright (C) 2024-present Puter Technologies Inc. + * + * This file is part of Puter. + * + * Puter is free software: you can redistribute it and/or modify + * it under the terms of the GNU Affero General Public License as published + * by the Free Software Foundation, either version 3 of the License, or + * (at your option) any later version. + * + * This program is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + * GNU Affero General Public License for more details. + * + * You should have received a copy of the GNU Affero General Public License + * along with this program. If not, see . + */ + +import { PuterStore } from '../types.js'; +import { KV_GLOBAL_APP_KEY } from '../systemKv/SystemKVStore.js'; + +/** + * Which regions hold a socket for a (user, app). + * + * The row itself is a reserved item in the replicated key-value table, written + * with the store's direct item path: presence is platform bookkeeping, not the + * user's data, so it is never metered, never listed, and never announced as a + * key-value change. + * + * Two things keep the cost of that table proportional to session churn rather + * than to connected population: + * + * - **Writes are transition-only.** The first socket a region holds for a pair + * writes; every reconnect after it writes nothing, answered by a region-local + * counter in Redis. There is no periodic refresh of any kind — a row is + * corrected when someone reads it and finds it wrong, not on a timer. + * - **Reads are keyed by a per-user generation**, bumped here on every transition + * and repair and broadcast to peer regions, so a caller's cache turns over + * when presence actually moves. + * + * A region that dies without disconnecting leaves itself in the row. That is + * the designed failure: the next forward to it comes back "no socket", and the + * emitting region removes it with a write conditional on `version`. + */ + +// -- Keys ------------------------------------------------------------- + +/** The app a socket with no app of its own is counted under. */ +export const PRESENCE_NO_APP = KV_GLOBAL_APP_KEY; + +/** Reserved-item key of one pair's row. */ +export const presenceItemKey = (userUuid: string, appUid: string): string => + `pr#${userUuid}#${appUid}`; + +/** Map field of the row that holds the regions. */ +const REGIONS_FIELD = 'regions'; + +const connectionsKey = (userId: number | string, appUid: string): string => + `ev:pc:{${userId}}:${appUid}`; + +const generationKey = (userId: number | string): string => `ev:pg:{${userId}}`; + +// -- Lifetimes -------------------------------------------------------- + +/** + * How long a region's connection count survives untouched. Only ever reached by + * a node that died holding sockets, and a count stuck high keeps its region in + * the row — which lazy repair is what corrects, so the backstop is generous. + */ +const CONNECTION_COUNT_TTL_SECONDS = 24 * 60 * 60; + +/** + * The generation outlives the sessions it orders: one that expired and + * restarted at zero would let a cached row look current again. + */ +const GENERATION_TTL_SECONDS = 24 * 60 * 60; + +// -- Row -------------------------------------------------------------- + +/** One pair's presence, as the table stores it. */ +export interface PresenceRow { + /** Region name to the moment its first socket for this pair connected. */ + regions: Record; + /** Optimistic-concurrency counter every write moves. */ + version: number; +} + +const readRow = (item: Record | null): PresenceRow => { + const regions = item?.[REGIONS_FIELD]; + const version = Number(item?.version ?? 0); + return { + regions: + regions && typeof regions === 'object' && !Array.isArray(regions) + ? (regions as Record) + : {}, + version: Number.isFinite(version) ? version : 0, + }; +}; + +export class PresenceStore extends PuterStore { + // -- The row ----------------------------------------------------- + + async read(userUuid: string, appUid: string): Promise { + const item = await this.stores.kv.getReservedItem< + Record + >(presenceItemKey(userUuid, appUid)); + return readRow(item); + } + + /** + * Record that this region now holds a socket for the pair. Unconditional: + * nothing else can teach a peer region that a socket exists, and the write + * names only its own region, so a concurrent connect elsewhere keeps both. + */ + async join( + userUuid: string, + appUid: string, + region: string, + connectedAt: number = Date.now(), + ): Promise { + return this.stores.kv.setReservedEntry( + presenceItemKey(userUuid, appUid), + REGIONS_FIELD, + region, + connectedAt, + ); + } + + /** + * Take a region out of the row, but only while it still carries the version + * that was read. False means a fresher connect won the race, which is + * exactly the outcome that must not be overwritten. + */ + async leave( + userUuid: string, + appUid: string, + region: string, + expectedVersion: number, + ): Promise { + return this.stores.kv.removeReservedEntry( + presenceItemKey(userUuid, appUid), + REGIONS_FIELD, + region, + expectedVersion, + ); + } + + // -- This region's connections ----------------------------------- + + /** + * Count one more connection for the pair in this region, and say whether it + * is the one that crossed zero — the only connect that owes a write. + * + * The existing concurrency slots cannot answer this: they expose no count, + * key on the user rather than the pair, and fail open, which is wrong in + * precisely the situation presence exists for. + */ + async addConnection(userId: number, appUid: string): Promise { + const key = connectionsKey(userId, appUid); + const count = await this.clients.redis.incr(key); + await this.clients.redis.expire(key, CONNECTION_COUNT_TTL_SECONDS); + return Number(count); + } + + /** Drop one connection. Zero is the count that owes the region's removal. */ + async removeConnection(userId: number, appUid: string): Promise { + const key = connectionsKey(userId, appUid); + const count = Number(await this.clients.redis.decr(key)); + // Gone at zero, so the keyspace stays proportional to connected pairs + // — and a count driven negative by a double-reap resets with it. + if (count <= 0) { + await this.clients.redis.del(key); + return 0; + } + await this.clients.redis.expire(key, CONNECTION_COUNT_TTL_SECONDS); + return count; + } + + /** Whether this region still holds any socket for the pair. */ + async holdsConnection(userId: number, appUid: string): Promise { + const raw = await this.clients.redis.get( + connectionsKey(userId, appUid), + ); + return raw !== null && Number(raw) > 0; + } + + // -- Generation -------------------------------------------------- + + /** Advance the user's presence generation. One key, so one command. */ + async bumpGeneration(userId: number): Promise { + const key = generationKey(userId); + const next = await this.clients.redis.incr(key); + await this.clients.redis.expire(key, GENERATION_TTL_SECONDS); + return typeof next === 'number' ? next : Number(next); + } +} diff --git a/src/backend/stores/index.ts b/src/backend/stores/index.ts index 44afc272a..237d65326 100644 --- a/src/backend/stores/index.ts +++ b/src/backend/stores/index.ts @@ -34,6 +34,7 @@ import { S3ObjectStore } from './fs/S3ObjectStore.js'; import { SessionStore } from './session/SessionStore.js'; import { ShareStore } from './share/ShareStore.js'; import { SubdomainStore } from './subdomain/SubdomainStore.js'; +import { PresenceStore } from './events/PresenceStore.js'; import { SystemKVStore } from './systemKv/SystemKVStore.js'; import { TeamStore } from './team/TeamStore.js'; import { UserBlockStore } from './userBlock/UserBlockStore.js'; @@ -70,6 +71,7 @@ declare module './types.js' { durableSubscription: DurableSubscriptionStore; eventHandler: EventHandlerStore; pendingDelivery: PendingDeliveryStore; + presence: PresenceStore; } } @@ -107,4 +109,6 @@ export const puterStores = { durableSubscription: DurableSubscriptionStore, // Table only, and reads the subscription table for its dependent counts. eventHandler: EventHandlerStore, + // Writes presence rows through `kv`'s reserved-item path, so it follows it. + presence: PresenceStore, } satisfies IPuterStoreRegistry; diff --git a/src/backend/stores/systemKv/SystemKVStore.ts b/src/backend/stores/systemKv/SystemKVStore.ts index eae2a0b91..a19ed0f9d 100644 --- a/src/backend/stores/systemKv/SystemKVStore.ts +++ b/src/backend/stores/systemKv/SystemKVStore.ts @@ -138,6 +138,17 @@ export interface RecursiveRecord { export const KV_GLOBAL_APP_KEY = 'os-global'; const SYSTEM_NAMESPACE = `v1:${SYSTEM_ACTOR_UUID}:${KV_GLOBAL_APP_KEY}`; const MAX_KEY_BYTES = 1024; + +/** Optimistic-concurrency counter every reserved-item write moves. */ +const RESERVED_VERSION_ATTR = 'version'; + +/** + * Whether a write was refused because the condition it carried no longer held. + * The compare-and-set answer, not a failure: the caller re-reads and decides. + */ +const isConditionRefused = (err: unknown): boolean => + (err as { name?: string })?.name === 'ConditionalCheckFailedException'; +const RESERVED_VERSION_BUMP = '#v = if_not_exists(#v, :zero) + :one'; const MAX_VALUE_BYTES = 399 * 1024; // A number anywhere inside a value is bounded too, to the IEEE-754 safe // integer range — past that it cannot round-trip, so it is clamped to the @@ -790,6 +801,122 @@ export class SystemKVStore extends PuterStore { ); } + // -- Reserved items ----------------------------------------------- + // + // Platform bookkeeping that happens to live in this table and is not + // anyone's key-value data. Reserved items sit in the system namespace under + // their own key prefix, so the driver can never address one, and they go + // straight to the table: no usage is returned because nothing is billed, no + // read cache is consulted or invalidated, and no mutation event is emitted + // — a reserved item is not a change to a namespace anyone can subscribe to. + + /** One reserved item, or `null`. Eventually consistent, which is enough. */ + async getReservedItem(key: string): Promise { + assertKey(key); + const response = await this.clients.dynamo.get(this.tableName, { + namespace: SYSTEM_NAMESPACE, + key, + }); + return (response.Item as T | undefined) ?? null; + } + + /** + * Set one entry of a reserved item's map field and advance the item's + * version. Unconditional and safe to race: the expression names only its + * own entry, so two writers adding different ones keep both. Returns the + * version the item now carries. + */ + async setReservedEntry( + key: string, + field: string, + entry: string, + value: unknown, + ): Promise { + assertKey(key); + assertSafeValueKeys({ [field]: { [entry]: value } }); + + // Two shapes because a nested path cannot be set on a map that is not + // there yet: the common one first, the seed only when it is missing. + const nested = () => + this.#bumpReserved( + key, + `SET #f.#e = :value, ${RESERVED_VERSION_BUMP}`, + { ':value': value }, + { '#f': field, '#e': entry }, + 'attribute_exists(#f)', + ); + + const applied = await nested(); + if (applied !== null) return applied; + + const seeded = await this.#bumpReserved( + key, + `SET #f = :seed, ${RESERVED_VERSION_BUMP}`, + { ':seed': { [entry]: value } }, + { '#f': field }, + 'attribute_not_exists(#f)', + ); + if (seeded !== null) return seeded; + + // Another writer seeded the field between the two attempts. + return (await nested()) ?? 0; + } + + /** + * Drop one entry of a reserved item's map field, only while the item still + * carries `expectedVersion`. False means it has moved on since it was read + * — the compare-and-set answer, not a failure. + */ + async removeReservedEntry( + key: string, + field: string, + entry: string, + expectedVersion: number, + ): Promise { + assertKey(key); + const applied = await this.#bumpReserved( + key, + `SET ${RESERVED_VERSION_BUMP} REMOVE #f.#e`, + {}, + { '#f': field, '#e': entry }, + '#v = :expected', + { ':expected': expectedVersion }, + ); + return applied !== null; + } + + /** + * One version-advancing update of a reserved item. `null` when the + * condition refused it, which is an answer rather than an error. + */ + async #bumpReserved( + key: string, + expression: string, + values: Record, + names: Record, + condition: string, + conditionValues: Record = {}, + ): Promise { + try { + const response = await this.clients.dynamo.update( + this.tableName, + { namespace: SYSTEM_NAMESPACE, key }, + expression, + { ...values, ...conditionValues, ':zero': 0, ':one': 1 }, + { ...names, '#v': RESERVED_VERSION_ATTR }, + { condition }, + ); + return Number( + (response.Attributes as Record | undefined)?.[ + RESERVED_VERSION_ATTR + ] ?? 0, + ); + } catch (err) { + if (isConditionRefused(err)) return null; + throw err; + } + } + // -- Public API --------------------------------------------------- /** diff --git a/src/docs/src/Events.md b/src/docs/src/Events.md index 0d6837b61..b03d80711 100644 --- a/src/docs/src/Events.md +++ b/src/docs/src/Events.md @@ -176,6 +176,12 @@ Pass `handler` as a **function** and it runs here too, whenever this client is t A persistent subscription can also stop without you unsubscribing: its handler was removed, its holder ran out of credit, the handler kept failing, or the share it was made under was withdrawn. It is then *suspended* rather than deleted, and [`list()`](/Events/list/) reports `suspendedAt` and `suspendedReason`. Everything but a withdrawn grant can resume. +### Where your client is connected does not matter + +Puter runs in several places, and a client connects to whichever one is nearest. Nothing about that is yours to think about: an event finds the connection wherever it is, `ack()` settles the delivery it belongs to whichever connection you called it on, and the shape of everything you receive is identical either way. + +The one consequence worth knowing is the one already stated: a `single` delivery is **at-least-once**. Undelivered events are held where the change happened, so a deployment going down loses only what it was still holding — the subscription itself, and everything already delivered, is unaffected. Handlers are asked to be idempotent for this reason, and `event.id` is the key to deduplicate on. + ## Limits Subscriptions per connection, persistent subscriptions per account, published handlers per app, subscribe calls per minute, and how much one event may fan out are all capped — see [Rate Limits and Quotas](/rate-limits-and-quotas/). Deliveries are coalesced over 250 ms per subject, so a multipart upload or a save loop arrives as one event rather than one per write. diff --git a/src/puter-js/src/modules/events/lib/channel.js b/src/puter-js/src/modules/events/lib/channel.js index e1d8be844..5a14d9535 100644 --- a/src/puter-js/src/modules/events/lib/channel.js +++ b/src/puter-js/src/modules/events/lib/channel.js @@ -394,7 +394,7 @@ export class EventChannel { * * @internal * @param {DurableRegistration} registration - * @param {{ event?: unknown, ackRequired?: boolean, ackId?: string }} envelope + * @param {{ event?: unknown, ackRequired?: boolean, ackId?: string, origin?: string }} envelope * @returns {void} */ runDurable (registration, envelope) { @@ -410,7 +410,11 @@ export class EventChannel { const ack = () => { if ( acked ) return Promise.resolve(); acked = true; - return this.ack(registration.subId, /** @type {string} */ (envelope.ackId)); + return this.ack( + registration.subId, + /** @type {string} */ (envelope.ackId), + envelope.origin, + ); }; const { puter } = this.module; settleHandler( @@ -433,11 +437,18 @@ export class EventChannel { * @internal * @param {string} subId * @param {string} ackId + * @param {string} [origin] Echoed back untouched: it names whichever + * deployment is holding the delivery, which need not be the one this + * connection reached. * @returns {Promise} */ - async ack (subId, ackId) { + async ack (subId, ackId, origin) { try { - await this.request(ACK_VERB, { subId, id: ackId }, DEFAULT_TIMEOUT_MS); + await this.request( + ACK_VERB, + { subId, id: ackId, ...(origin ? { origin } : {}) }, + DEFAULT_TIMEOUT_MS, + ); } catch (error) { console.warn('[puter.events] could not acknowledge a delivery', error); }