From eb8f497e9f8e5cb4250bb0534f4c8dee0607b9e7 Mon Sep 17 00:00:00 2001 From: Daniel Salazar Date: Thu, 3 Sep 2026 15:39:21 -0700 Subject: [PATCH] feat: presence and cross-region event forwarding (PUT-1679) (#3686) * feat: presence and cross-region event forwarding (PUT-1679) * fix: fan cache bumps to sibling nodes and stop the forward shed cascading (PUT-1679) `outer.events.generationBumped` and `outer.events.presenceBumped` rode `outer.*`, which the broadcast service only webhooks to peer regions; only `outer.pubsub.*` also fans over Redis to a region's other nodes. Both caches are per-process maps with no expiry, so a bump landing on one node left its siblings stale until that user's next transition. Renamed onto `outer.pubsub.events.*`; the listeners already accept the `from_outside` copy the Redis re-emit carries. `PeerForwardQueue.push` called `onOverflow` synchronously and the handler pushed markers straight back, each of which re-tripped the bound and shed the next item: one item over a 5000 bound recursed ~2200 deep, threw a RangeError, and turned ~2200 queued deliveries into gap markers. It also re-summed `bytes` over the whole queue per drop. The handler now returns its markers and the queue appends them past the bound check, sheds deliveries before markers, keeps one pending marker per (peer, subscription), and subtracts bytes per dropped item. --- src/backend/clients/dynamodb/DDBClient.ts | 4 + src/backend/clients/event/types.ts | 20 +- .../broadcast/BroadcastController.ts | 71 +- .../services/broadcast/BroadcastService.ts | 136 ++- .../services/events/EventForwardService.ts | 592 +++++++++++ .../services/events/EventsService.test.ts | 24 +- src/backend/services/events/EventsService.ts | 162 ++- .../events/durable.integration.test.ts | 4 +- .../services/events/forwardQueue.test.ts | 275 +++++ src/backend/services/events/forwardQueue.ts | 269 +++++ .../services/events/forwarding.test.ts | 999 ++++++++++++++++++ src/backend/services/events/metering.test.ts | 12 + .../services/events/presenceCache.test.ts | 174 +++ src/backend/services/events/presenceCache.ts | 178 ++++ .../services/events/singleDelivery.test.ts | 12 + src/backend/services/index.ts | 5 + .../events/PresenceStore.integration.test.ts | 206 ++++ src/backend/stores/events/PresenceStore.ts | 197 ++++ src/backend/stores/index.ts | 4 + src/backend/stores/systemKv/SystemKVStore.ts | 127 +++ src/docs/src/Events.md | 6 + .../src/modules/events/lib/channel.js | 19 +- 22 files changed, 3438 insertions(+), 58 deletions(-) create mode 100644 src/backend/services/events/EventForwardService.ts create mode 100644 src/backend/services/events/forwardQueue.test.ts create mode 100644 src/backend/services/events/forwardQueue.ts create mode 100644 src/backend/services/events/forwarding.test.ts create mode 100644 src/backend/services/events/presenceCache.test.ts create mode 100644 src/backend/services/events/presenceCache.ts create mode 100644 src/backend/stores/events/PresenceStore.integration.test.ts create mode 100644 src/backend/stores/events/PresenceStore.ts 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); }