Skip to content
Merged
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
74 changes: 73 additions & 1 deletion server/src/__integration__/realtime-outbox.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@ import { createServer, type Server } from 'node:http'
import { after, before, beforeEach, test } from 'node:test'
import assert from 'node:assert/strict'
import { pool } from '../db/pool.js'
import { drainRealtimeOutbox } from '../realtime-outbox.js'
import { drainRealtimeOutbox, enqueueBroadcast } from '../realtime-outbox.js'
import type { BroadcastEvent } from '../redis.js'
import {
buildApiTestApp, ensureSchemaOnce, resetAllTables,
Expand Down Expand Up @@ -154,3 +154,75 @@ test('[integration] document and calendar creates replay without duplicate rows
)
assert.equal(events.rows[0]?.count, 2)
})

// ── a body cut mid-emoji must not roll the transaction back ─────────────────
//
// The quote path truncates by UTF-16 code unit (`qr[0].body.slice(0, 240)`), so
// a message with a non-BMP character straddling that index yields a lone
// surrogate. JSON.stringify emits it as the ASCII escape `\ud83d`, which reaches
// Postgres intact and is refused by the jsonb cast:
//
// ERROR: invalid input syntax for type json
// DETAIL: Unicode low surrogate must follow a high surrogate.
//
// Because the enqueue is inside the caller's transaction, that took the message
// with it. `messages.body` is TEXT and accepts the same bytes, so only the
// outbox half died — which is what made one emoji-bearing message permanently
// un-quotable rather than merely ugly.

/** Exactly what the quote path produces: 239 ASCII chars then an emoji, cut at 240. */
const CUT_MID_EMOJI = `${'a'.repeat(239)}\u{1F600}`.slice(0, 240)

test('[integration] a payload truncated mid-emoji still enqueues', async () => {
assert.equal(CUT_MID_EMOJI.charCodeAt(239), 0xd83d, 'fixture no longer ends in a lone surrogate')

const id = await enqueueBroadcast(pool, 'message:new', {
type: 'message.new',
conversationId: 'c-outbox',
quoted: { body: CUT_MID_EMOJI },
} as unknown as BroadcastEvent)

const { rows } = await pool.query<{ body: string }>(
`SELECT payload->'quoted'->>'body' AS body FROM realtime_outbox WHERE id = $1`,
[id],
)
assert.equal(rows.length, 1, 'the row was not written')
// The broken half is dropped, not the message: 239 readable characters survive.
assert.equal(rows[0].body, 'a'.repeat(239))
})

test('[integration] a lone surrogate anywhere in the tree is scrubbed, not just at the top', async () => {
// Payload fields grow over time; scrubbing lives at the single cast so a new
// nested field cannot reintroduce this.
const id = await enqueueBroadcast(pool, 'message:new', {
type: 'message.new',
conversationId: 'c-outbox',
quoted: { body: '\ud83d', authorName: 'tail \ud83d' },
tags: ['\ud83d', 'ok'],
} as unknown as BroadcastEvent)

const { rows } = await pool.query<{ payload: Record<string, unknown> }>(
`SELECT payload FROM realtime_outbox WHERE id = $1`, [id],
)
const payload = rows[0].payload as {
quoted: { body: string; authorName: string }; tags: string[]
}
assert.equal(payload.quoted.body, '')
assert.equal(payload.quoted.authorName, 'tail ')
assert.deepEqual(payload.tags, ['', 'ok'])
})

test('[integration] a well-formed emoji is untouched', async () => {
// The scrub must only remove UNPAIRED halves — a real emoji has to survive
// intact or every message carrying one arrives mangled.
const id = await enqueueBroadcast(pool, 'message:new', {
type: 'message.new',
conversationId: 'c-outbox',
quoted: { body: 'ship it \u{1F680} done' },
} as unknown as BroadcastEvent)
const { rows } = await pool.query<{ body: string }>(
`SELECT payload->'quoted'->>'body' AS body FROM realtime_outbox WHERE id = $1`,
[id],
)
assert.equal(rows[0].body, 'ship it \u{1F680} done')
})
28 changes: 27 additions & 1 deletion server/src/realtime-outbox.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import { randomUUID } from 'node:crypto'
import type { PoolClient, QueryResult, QueryResultRow } from 'pg'
import { pool } from './db/pool.js'
import { stripLoneSurrogates } from './agents/text-safety.js'
import { publish, type BroadcastEvent } from './redis.js'

/**
Expand Down Expand Up @@ -60,11 +61,36 @@ export async function enqueueBroadcast(
await db.query(
`INSERT INTO realtime_outbox (id, channel, payload)
VALUES ($1, $2, $3::jsonb)`,
[id, channel, JSON.stringify(payload)],
[id, channel, serializePayload(payload)],
)
return id
}

/** Serialize an outbox payload so Postgres will accept it as `jsonb`.
*
* A body truncated mid-emoji — `qr[0].body.slice(0, 240)` on the quote path —
* ends in a lone UTF-16 surrogate. JSON.stringify turns that into the literal
* ASCII escape `\ud83d`, which survives transport intact and is then rejected:
*
* ERROR: invalid input syntax for type json
* DETAIL: Unicode low surrogate must follow a high surrogate.
*
* That matters here and not before because the enqueue is INSIDE the caller's
* transaction. The same payload used to go to redis.publish AFTER COMMIT, where
* a lone surrogate was a cosmetic glyph; behind a `::jsonb` cast it rolls the
* message back and 500s, so a quote-reply to one emoji-bearing message fails
* forever. `messages.body` is TEXT and accepts the same bytes, which is why
* only the outbox half dies.
*
* Scrubbing here rather than at each call site is deliberate: every field of
* every future event goes through this one cast. The replacer visits every
* string in the tree, so nesting and arrays need no traversal of our own. */
function serializePayload(payload: unknown): string {
return JSON.stringify(payload, (_key, value) =>
typeof value === 'string' ? stripLoneSurrogates(value) : value,
)
}

/** Run a mutation and all of its realtime invalidations in one transaction. */
export async function withOutboxTransaction<T>(
run: (client: PoolClient) => Promise<T>,
Expand Down
Loading