Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions backend/.env.example
Original file line number Diff line number Diff line change
Expand Up @@ -97,3 +97,7 @@ RESEND_API_KEY=your_resend_api_key
# Sender address for receipt emails (must be a verified domain in Resend)
# Defaults to: receipts@notifications.dripsnetwork.com
RECEIPT_FROM_EMAIL=receipts@yourdomain.com

# WebSocket relay limits (issue #1452)
WS_MAX_HTTP_BUFFER_SIZE=16384
WS_MAX_INVALID_EVENTS=10
78 changes: 78 additions & 0 deletions backend/docs/WEBSOCKET_RELAY_VALIDATION.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,78 @@
# WebSocket Relay: Payload Sanitization & Strict Validation

Issue: #1452

The Socket.IO relay (`src/lib/websocket-relay-server.js`) lets merchant dashboards
join `merchant:<id>` rooms and checkout pages join `checkout:<id>` rooms. Every
inbound event now passes through a guard (`src/lib/websocket-relay-validation.js`)
before any handler runs.

## Inbound pipeline

1. **Engine limit.** `maxHttpBufferSize` is 16 KiB (the Socket.IO default is 1 MB). Larger frames close the connection.
2. **Event allow-list.** Only `join:merchant`, `join:checkout` and `leave:checkout` are accepted. Any other event gets `UNKNOWN_EVENT`.
3. **Shape check.** The payload must be a plain JSON object. A missing payload, `null`, an array or a primitive gets `INVALID_PAYLOAD`.
4. **Byte limit.** The serialised payload must be 4 KiB or less (`PAYLOAD_TOO_LARGE`).
5. **Deep sanitization** (`sanitizeRelayPayload`):
- drops the `__proto__`, `constructor` and `prototype` keys at any depth
- NFC-normalises strings and strips C0/C1 control characters and bidi overrides (`\t \n \r` are kept)
- rejects non-finite numbers, BigInt, binary data and circular references
- bounds depth (8), keys per object (64), array length (256) and string length (2048)
6. **Strict schema** (zod `.strict()`). Unknown keys are rejected. IDs must be UUIDs and are trimmed and lower-cased, so room names are canonical and cannot be injected (for example `"<uuid>:admin"`).

## Failure handling

| Situation | Behaviour |
| --- | --- |
| Invalid event, client sent an ack callback | `ack({ ok: false, event, error: { code, message, issues? } })` |
| Invalid event, no ack callback | `socket.emit("relay:error", { ok: false, event, error })` |
| Every rejection | `logger.warn` records the socket id, event, error code and violation count. The raw payload is never logged or echoed. |
| `WS_MAX_INVALID_EVENTS` rejections on one socket | The socket is disconnected. |
| A handler throws | The error is logged, the client gets `INTERNAL_ERROR`, and the process keeps running. |

Error codes: `UNKNOWN_EVENT`, `INVALID_PAYLOAD`, `PAYLOAD_TOO_LARGE`,
`VALIDATION_FAILED`, `PAYLOAD_TOO_DEEP`, `TOO_MANY_KEYS`, `ARRAY_TOO_LONG`,
`STRING_TOO_LONG`, `INVALID_NUMBER`, `INVALID_DATE`, `UNSUPPORTED_TYPE`,
`CIRCULAR_REFERENCE`, `UNSERIALIZABLE`, `INTERNAL_ERROR`.

## Outbound sanitization

`sanitizeOutboundPayload()` is applied to `checkout:presence` emits and to all
payment events from the Horizon poller (`notifyPaymentEvent`). It uses the same
rules with looser size limits. Data that came from Horizon or from merchant
metadata cannot push prototype keys, control characters or non-JSON values to
dashboards. If an outbound payload cannot be sanitized, it is dropped and a
warning is logged; nothing crashes.

`sanitizeRelayMessage()` in `websocket-relay-security.js` now also deep-sanitizes
the value of each allowed field.

## Configuration

| Variable | Default | Purpose |
| --- | --- | --- |
| `WS_MAX_HTTP_BUFFER_SIZE` | `16384` | Maximum inbound frame size in bytes |
| `WS_MAX_INVALID_EVENTS` | `10` | Invalid events allowed per socket before it is disconnected |

## Adding a new inbound event

Add a strict zod schema to `INBOUND_EVENT_SCHEMAS`, then register the handler with
`relay.on(event, handler)`. `createRelayGuard().on()` refuses events that have no schema.

## Security notes

- **Fixed:** a `join:*` event with no payload used to throw a destructuring `TypeError` inside the listener.
- **Fixed:** any non-empty string was accepted as a room ID. Arbitrary room names could be created, and the relay's memory could grow without bound.
- **Fixed:** frames up to 1 MB were accepted for events that need fewer than 100 bytes.
- **Out of scope:** `join:merchant` is still unauthenticated. Anyone who knows a merchant UUID can listen to that merchant's room. The follow-up is to require a JWT at handshake (`verifyRelayToken`) and check that the token's merchant matches the room.

## Tests

```
npx vitest run src/lib/websocket-relay-validation.test.js \
src/lib/websocket-relay-server.test.js \
src/lib/websocket-relay-security.test.js
```

`websocket-relay-server.test.js` starts a real Socket.IO server and drives it with
a raw WebSocket client that speaks the Engine.IO v4 / Socket.IO v5 wire protocol.
67 changes: 8 additions & 59 deletions backend/src/app.js
Original file line number Diff line number Diff line change
@@ -1,7 +1,6 @@
import cors from "cors";
import helmet from "helmet";
import express from "express";
import { Server as SocketIOServer } from "socket.io";
import swaggerUi from "swagger-ui-express";
import { ZodError } from "zod";
import path from "node:path";
Expand Down Expand Up @@ -49,6 +48,7 @@ import { versionDeprecationMiddleware } from "./lib/version-deprecation.js";
import oracleRouter from "./routes/oracle.js";
import { getPaymentSessionValidatorHealth } from "./lib/payment-session-validator.js";
import { configureExchangeRateCoordination } from "./services/exchangeRateService.js";
import { createRelayServer } from "./lib/websocket-relay-server.js";

export async function createApp({ redisClient }) {
const app = express();
Expand All @@ -62,64 +62,13 @@ export async function createApp({ redisClient }) {
const __dirname = path.dirname(__filename);
const publicDir = path.join(__dirname, "..", "public");

// Create socket.io instance (attached to HTTP server in server.js)
const io = new SocketIOServer({
cors: {
origin: process.env.CORS_ALLOWED_ORIGINS
? process.env.CORS_ALLOWED_ORIGINS.split(",").map((o) => o.trim())
: ["http://localhost:3000"],
credentials: true,
},
});

const checkoutRoomName = (paymentId) => `checkout:${paymentId}`;
const emitCheckoutPresence = (paymentId) => {
const room = checkoutRoomName(paymentId);
const activeViewers = io.sockets.adapter.rooms.get(room)?.size ?? 0;

io.to(room).emit("checkout:presence", {
payment_id: paymentId,
active_viewers: activeViewers,
});
};

// Socket.io room management: clients join their merchant-specific room
io.on("connection", (socket) => {
const joinedCheckoutRooms = new Set();

socket.on("join:merchant", ({ merchant_id }) => {
if (typeof merchant_id === "string" && merchant_id.length > 0) {
socket.join(`merchant:${merchant_id}`);
}
});

socket.on("join:checkout", ({ payment_id }) => {
if (typeof payment_id !== "string" || payment_id.length === 0) {
return;
}

const room = checkoutRoomName(payment_id);
joinedCheckoutRooms.add(payment_id);
socket.join(room);
emitCheckoutPresence(payment_id);
});

socket.on("leave:checkout", ({ payment_id }) => {
if (typeof payment_id !== "string" || payment_id.length === 0) {
return;
}

joinedCheckoutRooms.delete(payment_id);
socket.leave(checkoutRoomName(payment_id));
emitCheckoutPresence(payment_id);
});

socket.on("disconnect", () => {
for (const paymentId of joinedCheckoutRooms) {
emitCheckoutPresence(paymentId);
}
joinedCheckoutRooms.clear();
});
// Create socket.io relay (attached to HTTP server in server.js). Inbound
// events are sanitized and strictly validated before use (issue #1452).
const io = createRelayServer({
corsOrigins: process.env.CORS_ALLOWED_ORIGINS
? process.env.CORS_ALLOWED_ORIGINS.split(",").map((o) => o.trim())
: ["http://localhost:3000"],
logger,
});

// Make DB pool and io accessible on every request
Expand Down
10 changes: 9 additions & 1 deletion backend/src/lib/horizon-poller.js
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,7 @@ import { sendReceiptEmail } from "./email.js";
import { renderReceiptEmail } from "./email-templates.js";
import { getPayloadForVersion } from "../webhooks/resolver.js";
import { streamManager } from "./stream-manager.js";
import { sanitizeOutboundPayload } from "./websocket-relay-validation.js";
import { connectRedisClient, invalidatePaymentCache } from "./redis.js";
import { logger } from "./logger.js";
import {
Expand Down Expand Up @@ -790,7 +791,14 @@ function sleep(ms) {
function notifyPaymentEvent(payment, { sseEvent, sseData, socketEvent, socketData }) {
streamManager.notify(payment.id, sseEvent, sseData);
if (_io && payment.merchant_id) {
_io.to(`merchant:${payment.merchant_id}`).emit(socketEvent, socketData);
let payload;
try {
payload = sanitizeOutboundPayload(socketData);
} catch (err) {
logger.warn({ err, paymentId: payment.id, socketEvent }, "Horizon poller: dropped unsafe socket payload");
return;
}
_io.to(`merchant:${payment.merchant_id}`).emit(socketEvent, payload);
}
}

Expand Down
10 changes: 7 additions & 3 deletions backend/src/lib/websocket-relay-security.js
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
*/

import jwt from "jsonwebtoken";
import { sanitizeRelayPayload } from "./websocket-relay-validation.js";

// ─── Allowed message fields ───────────────────────────────────────────────────

Expand Down Expand Up @@ -135,7 +136,8 @@ function verifyRelayToken(token, secret) {
*
* @param {any} msg - The parsed WebSocket message object
* @returns {{ sanitized: object, warnings: string[] }}
* @throws {Error} When `msg` is not a non-null object, or when required fields are missing
* @throws {Error} When `msg` is not a non-null object, when required fields are missing,
* or (RelayValidationError) when a field value breaks the relay payload limits
*/
function sanitizeRelayMessage(msg) {
if (msg === null || typeof msg !== "object" || Array.isArray(msg)) {
Expand All @@ -145,10 +147,12 @@ function sanitizeRelayMessage(msg) {
const warnings = [];
const sanitized = {};

// Copy only allowed fields
// Copy only allowed fields, deep-sanitizing each value so nested payloads
// cannot carry prototype-pollution keys or control characters (issue #1452)
for (const [key, value] of Object.entries(msg)) {
if (ALLOWED_MESSAGE_FIELDS.has(key)) {
sanitized[key] = value;
const clean = sanitizeRelayPayload(value);
if (clean !== undefined) sanitized[key] = clean;
} else {
warnings.push(`Unknown field stripped: '${key}'`);
}
Expand Down
108 changes: 108 additions & 0 deletions backend/src/lib/websocket-relay-server.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,108 @@
/**
* websocket-relay-server.js
*
* Builds the Socket.IO relay server used for merchant dashboards and checkout
* presence. Every inbound event goes through the relay guard, which sanitizes
* and strictly validates the payload before any room is joined (issue #1452).
*/

import { Server as SocketIOServer } from "socket.io";
import {
createRelayGuard,
sanitizeOutboundPayload,
DEFAULT_MAX_VIOLATIONS,
} from "./websocket-relay-validation.js";

/** Inbound relay events are tiny room joins; the socket.io default is 1 MB. */
export const DEFAULT_WS_MAX_HTTP_BUFFER_SIZE = 16 * 1024;

function parsePositiveInt(value, fallback) {
const parsed = Number.parseInt(value, 10);
return Number.isFinite(parsed) && parsed > 0 ? parsed : fallback;
}

export const checkoutRoomName = (paymentId) => `checkout:${paymentId}`;

/**
* Register the relay's connection handlers on an existing Socket.IO server.
*
* @param {import("socket.io").Server} io
* @param {object} [opts]
* @param {{ warn: Function, error: Function }} [opts.logger]
* @param {number} [opts.maxViolations] - Invalid events tolerated per socket
*/
export function attachRelayHandlers(io, opts = {}) {
const emitCheckoutPresence = (paymentId) => {
const room = checkoutRoomName(paymentId);
const activeViewers = io.sockets.adapter.rooms.get(room)?.size ?? 0;

io.to(room).emit(
"checkout:presence",
sanitizeOutboundPayload({
payment_id: paymentId,
active_viewers: activeViewers,
}),
);
};

io.on("connection", (socket) => {
const joinedCheckoutRooms = new Set();
const relay = createRelayGuard(socket, {
logger: opts.logger,
maxViolations: opts.maxViolations,
});

relay.on("join:merchant", ({ merchant_id }) => {
socket.join(`merchant:${merchant_id}`);
});

relay.on("join:checkout", ({ payment_id }) => {
joinedCheckoutRooms.add(payment_id);
socket.join(checkoutRoomName(payment_id));
emitCheckoutPresence(payment_id);
});

relay.on("leave:checkout", ({ payment_id }) => {
joinedCheckoutRooms.delete(payment_id);
socket.leave(checkoutRoomName(payment_id));
emitCheckoutPresence(payment_id);
});

socket.on("disconnect", () => {
for (const paymentId of joinedCheckoutRooms) {
emitCheckoutPresence(paymentId);
}
joinedCheckoutRooms.clear();
});
});

return io;
}

/**
* Create the relay Socket.IO server (attached to the HTTP server in server.js).
*
* Environment:
* WS_MAX_HTTP_BUFFER_SIZE - max inbound frame size in bytes (default 16 KiB)
* WS_MAX_INVALID_EVENTS - invalid events before a socket is disconnected (default 10)
*
* @param {object} opts
* @param {string[]} opts.corsOrigins
* @param {{ warn: Function, error: Function }} [opts.logger]
* @param {NodeJS.ProcessEnv} [opts.env]
* @returns {import("socket.io").Server}
*/
export function createRelayServer({ corsOrigins, logger, env = process.env }) {
const io = new SocketIOServer({
cors: { origin: corsOrigins, credentials: true },
maxHttpBufferSize: parsePositiveInt(
env.WS_MAX_HTTP_BUFFER_SIZE,
DEFAULT_WS_MAX_HTTP_BUFFER_SIZE,
),
});

return attachRelayHandlers(io, {
logger,
maxViolations: parsePositiveInt(env.WS_MAX_INVALID_EVENTS, DEFAULT_MAX_VIOLATIONS),
});
}
Loading
Loading