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
1 change: 1 addition & 0 deletions changelog.d/fixed-nps-relay-conditional-writes.md
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
- Bump `@netlify/blobs` to 10.x and make the NPS relay update its rate-limit, nonce and install records with conditional (ETag-guarded) writes, so concurrent requests are bounded by the same caps as sequential ones.
124 changes: 122 additions & 2 deletions netlify/nps-relay/relay.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ import { fileURLToPath } from 'node:url'
import { beforeEach, describe, expect, it } from 'vitest'
import {
ACK_PATH,
CAS_MAX_ATTEMPTS,
ENTRY_PREFIX,
HEADER_INSTALL_ID,
HEADER_NONCE,
Expand All @@ -15,6 +16,7 @@ import {
MAX_FEEDBACK_CHARS,
MAX_NONCES_PER_INSTALL,
MAX_PER_INSTALL_PER_WINDOW,
MAX_PER_IP_PER_WINDOW,
MAX_PULL_BATCH,
MAX_REGISTER_BODY_BYTES,
MAX_REGISTRATIONS_PER_DAY,
Expand All @@ -34,10 +36,12 @@ import {
signingInput,
validateRegistration,
verifyEd25519,
type ConditionalSetOptions,
type InstallRecord,
type RelayDeps,
type RelayEntry,
type RelayStore,
type WriteOutcome,
} from './relay'

const HERE = path.dirname(fileURLToPath(import.meta.url))
Expand All @@ -51,22 +55,43 @@ function sha256(value: string): string {
return createHash('sha256').update(value, 'utf8').digest('hex')
}

/** In-memory stand-in for the Netlify Blobs store. */
/**
* In-memory stand-in for the Netlify Blobs store, including its ETag-based
* conditional writes. Every awaited call yields to the event loop so that
* concurrent handler invocations interleave the way separate function
* instances would.
*/
class MemoryStore implements RelayStore {
data = new Map<string, unknown>()
etags = new Map<string, string>()
writes = 0
private version = 0

async get(key: string): Promise<unknown> {
await new Promise((r) => setImmediate(r))
return this.data.has(key) ? JSON.parse(JSON.stringify(this.data.get(key))) : null
}

async setJSON(key: string, value: unknown): Promise<void> {
async getWithMetadata(key: string): Promise<{ data: unknown; etag?: string } | null> {
const data = await this.get(key)
if (data === null) return null
return { data, etag: this.etags.get(key) }
}

async setJSON(key: string, value: unknown, options?: ConditionalSetOptions): Promise<WriteOutcome> {
await new Promise((r) => setImmediate(r))
if (options?.onlyIfNew && this.data.has(key)) return { modified: false }
if (options?.onlyIfMatch !== undefined && this.etags.get(key) !== options.onlyIfMatch) return { modified: false }
this.writes++
const etag = `"v${++this.version}"`
this.data.set(key, JSON.parse(JSON.stringify(value)))
this.etags.set(key, etag)
return { modified: true, etag }
}

async delete(key: string): Promise<void> {
this.data.delete(key)
this.etags.delete(key)
}

async list({ prefix }: { prefix: string }): Promise<{ blobs: Array<{ key: string }> }> {
Expand Down Expand Up @@ -584,6 +609,101 @@ describe('hub pull and ack', () => {
})
})

describe('concurrent requests', () => {
// Each test fires a burst of requests without awaiting between them, so
// every handler reads the same pre-burst records: the race an attacker
// creates by sending many requests at once to parallel function instances.
function statuses(responses: Response[]): Record<number, number> {
return responses.reduce<Record<number, number>>((acc, r) => ({ ...acc, [r.status]: (acc[r.status] ?? 0) + 1 }), {})
}

it('holds the per-IP submission cap under a concurrent burst', async () => {
await registerOk()
const burst = 8
const responses = await Promise.all(Array.from({ length: burst }, () => handleRelayRequest(submit(GOOD), deps())))
expect(statuses(responses)[201]).toBe(MAX_PER_IP_PER_WINDOW)
expect(statuses(responses)[429]).toBe(burst - MAX_PER_IP_PER_WINDOW)
expect(store.entries()).toHaveLength(MAX_PER_IP_PER_WINDOW)
})

it('holds the per-install submission cap under a concurrent burst from many IPs', async () => {
await registerOk()
const burst = MAX_PER_INSTALL_PER_WINDOW + 6
const responses = await Promise.all(
Array.from({ length: burst }, (_, i) => handleRelayRequest(submit(GOOD), deps({ clientIp: `198.51.100.${i}` }))),
)
// Losers of the compare-and-swap fail closed, so a lockstep burst may
// admit fewer than the cap; it must never admit more.
const accepted = statuses(responses)[201] ?? 0
expect(accepted).toBeGreaterThan(0)
expect(accepted).toBeLessThanOrEqual(MAX_PER_INSTALL_PER_WINDOW)
expect(statuses(responses)[429]).toBe(burst - accepted)
expect(store.entries()).toHaveLength(accepted)
})

it('stores a replayed request at most once even when the copies race', async () => {
await registerOk()
const nonce = 'b'.repeat(32)
const responses = await Promise.all(
Array.from({ length: 6 }, (_, i) => handleRelayRequest(submit(GOOD, { nonce }), deps({ clientIp: `198.51.100.${i}` }))),
)
expect(statuses(responses)[201]).toBe(1)
expect(statuses(responses)[401]).toBe(5)
expect(store.entries()).toHaveLength(1)
})

it('holds the per-IP and per-day registration caps under a concurrent burst', async () => {
const burst = MAX_REGISTRATIONS_PER_IP_PER_WINDOW + 5
const responses = await Promise.all(Array.from({ length: burst }, () => handleRelayRequest(register(newHive()), deps())))
expect(statuses(responses)[201]).toBe(MAX_REGISTRATIONS_PER_IP_PER_WINDOW)
expect([...store.data.keys()].filter((k) => k.startsWith(INSTALL_PREFIX))).toHaveLength(MAX_REGISTRATIONS_PER_IP_PER_WINDOW)

const dayKey = `${REGISTER_DAY_PREFIX}${new Date(T0).toISOString().slice(0, 10)}`
expect((store.data.get(dayKey) as { count: number }).count).toBe(MAX_REGISTRATIONS_PER_IP_PER_WINDOW)
store.data.set(dayKey, { count: MAX_REGISTRATIONS_PER_DAY - 1 })
const dayBurst = await Promise.all(
Array.from({ length: 5 }, (_, i) => handleRelayRequest(register(newHive()), deps({ clientIp: `192.0.2.${i}` }))),
)
expect(statuses(dayBurst)[201]).toBe(1)
expect((store.data.get(dayKey) as { count: number }).count).toBe(MAX_REGISTRATIONS_PER_DAY)
})

it('binds a raced install id to exactly one key', async () => {
const a = newHive()
const b = { ...newHive(), installId: a.installId }
const [ra, rb] = await Promise.all([
handleRelayRequest(register(a), deps({ clientIp: '198.51.100.1' })),
handleRelayRequest(register(b), deps({ clientIp: '198.51.100.2' })),
])
expect([ra.status, rb.status].sort()).toEqual([201, 409])
const bound = (store.data.get(`${INSTALL_PREFIX}${a.installId}`) as InstallRecord).public_key
expect(bound).toBe(ra.status === 201 ? a.publicKey : b.publicKey)
})

it('fails closed when a record stays contended past the retry budget', async () => {
await registerOk()
let attempts = 0
const contended: RelayStore = {
...store,
get: (k) => store.get(k),
getWithMetadata: (k) => store.getWithMetadata(k),
delete: (k) => store.delete(k),
list: (o) => store.list(o),
setJSON: async (k, v, o) => {
if (k.startsWith(NONCE_PREFIX)) {
attempts++
return { modified: false }
}
return store.setJSON(k, v, o)
},
}
const res = await handleRelayRequest(submit(GOOD), deps({ store: contended }))
expect(res.status).toBe(429)
expect(attempts).toBe(CAS_MAX_ATTEMPTS)
expect(store.entries()).toHaveLength(0)
})
})

describe('Netlify function wiring', () => {
it('routes exactly the relay paths and reads only the hub secret env var', () => {
const fn = readFileSync(path.join(HERE, '..', 'functions', 'nps.mts'), 'utf8')
Expand Down
137 changes: 107 additions & 30 deletions netlify/nps-relay/relay.ts
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,9 @@
* - Bodies are capped by the bytes actually READ, never by Content-Length,
* so a chunked or lying request cannot make the function buffer more
* (console#16666).
* - Every bound above is enforced with conditional (compare-and-swap)
* writes, so concurrent requests cannot all pass a check against the
* same stale record; under persistent contention the relay fails closed.
* - Score must be an integer 1-4; feedback at most 500 characters; hive_id
* must match the hub's name rule.
* - Entries live in Netlify Blobs as a rolling window.
Expand Down Expand Up @@ -148,6 +151,12 @@ export const SIGNATURE_MAX_SKEW_MS = 5 * 60 * 1000
export const NONCE_RETENTION_MS = 2 * SIGNATURE_MAX_SKEW_MS
/** Most live nonces kept per install; more signed requests in the window get 429. */
export const MAX_NONCES_PER_INSTALL = 100
/**
* Attempts made to update one record with a conditional write before giving
* up. Each retry re-reads the record, so a lost race is re-checked against
* the winner's write rather than the stale copy.
*/
export const CAS_MAX_ATTEMPTS = 4

/** Prefix of the signed input; a future format must change it. */
export const SIGNATURE_VERSION = 'hive-nps-relay-v1'
Expand Down Expand Up @@ -194,10 +203,22 @@ const CONTROL_CHARS = /[\u0000-\u0008\u000b-\u001f\u007f]/g

// ── Types ────────────────────────────────────────────────────────────

/** Conditions a write may carry: update only this version, or create only. */
export type ConditionalSetOptions = { onlyIfMatch?: string; onlyIfNew?: never } | { onlyIfNew?: boolean; onlyIfMatch?: never }

/** Result of a (possibly conditional) write. */
export interface WriteOutcome {
/** False when the condition was not met and nothing was written. */
modified: boolean
etag?: string
}

/** The subset of the Netlify Blobs store API the relay uses. */
export interface RelayStore {
get(key: string, options: { type: 'json' }): Promise<unknown>
setJSON(key: string, value: unknown): Promise<unknown>
/** Like get, also returning the ETag the record can be conditionally updated with. */
getWithMetadata(key: string, options: { type: 'json' }): Promise<{ data: unknown; etag?: string } | null>
setJSON(key: string, value: unknown, options?: ConditionalSetOptions): Promise<WriteOutcome>
delete(key: string): Promise<unknown>
list(options: { prefix: string }): Promise<{ blobs: Array<{ key: string }> }>
}
Expand Down Expand Up @@ -437,12 +458,53 @@ async function entryKeysOldestFirst(store: RelayStore): Promise<string[]> {
}

/** Timestamps inside the window, from a stored rate record. */
async function recentTimes(store: RelayStore, key: string, now: number): Promise<number[]> {
const rec = (await store.get(key, { type: 'json' })) as RateRecord | null
const times = rec && Array.isArray(rec.times) ? rec.times : []
function liveTimes(rec: unknown, now: number): number[] {
const times = rec && Array.isArray((rec as RateRecord).times) ? (rec as RateRecord).times : []
return times.filter((t) => typeof t === 'number' && now - t < RATE_WINDOW_MS)
}

/** One compare-and-swap step: the record to write, or the response to stop with. */
type CasStep = { next: unknown } | { reject: Response }

/**
* Updates one record atomically. `step` sees the current record and decides
* the next one (or rejects); the write only lands if the record is still the
* version that was read. A lost race re-reads and re-decides, so a check
* such as "under the limit" is always made against the record the write
* replaces. Returns null on success, else the response to send. Persistent
* contention fails closed (429) rather than letting a write through
* unchecked.
*/
async function compareAndSwap(store: RelayStore, key: string, step: (current: unknown) => CasStep): Promise<Response | null> {
for (let attempt = 0; attempt < CAS_MAX_ATTEMPTS; attempt++) {
const current = await store.getWithMetadata(key, { type: 'json' })
const decision = step(current ? current.data : null)
if ('reject' in decision) return decision.reject
const condition: ConditionalSetOptions = current ? { onlyIfMatch: current.etag } : { onlyIfNew: true }
const result = await store.setJSON(key, decision.next, condition)
if (result.modified) return null
}
return RATE_LIMITED()
}

/**
* Reserves one slot in a windowed rate record, or returns 429 when the
* record already holds `max` live timestamps. Check and reservation are one
* conditional write, so concurrent callers cannot all fit in the last slot.
*/
function reserveSlot(store: RelayStore, key: string, max: number, now: number): Promise<Response | null> {
return compareAndSwap(store, key, (current) => {
const times = liveTimes(current, now)
if (times.length >= max) return { reject: RATE_LIMITED() }
return { next: { times: [...times, now] } satisfies RateRecord }
})
}

/** Whether a windowed rate record is already at `max`, without reserving. */
async function atLimit(store: RelayStore, key: string, max: number, now: number): Promise<boolean> {
return liveTimes(await store.get(key, { type: 'json' }), now).length >= max
}

async function getInstall(store: RelayStore, installId: string): Promise<InstallRecord | null> {
const rec = (await store.get(`${INSTALL_PREFIX}${installId}`, { type: 'json' })) as InstallRecord | null
return rec && typeof rec.public_key === 'string' ? rec : null
Expand Down Expand Up @@ -535,39 +597,49 @@ async function handleRegister(req: Request, deps: RelayDeps, now: number): Promi
return json(409, { error: 'install_id is registered to a different key' })
}

const ipKey = `${REGISTER_IP_RATE_PREFIX}${ipHash(deps)}`
const ipTimes = await recentTimes(store, ipKey, now)
if (ipTimes.length >= MAX_REGISTRATIONS_PER_IP_PER_WINDOW) return RATE_LIMITED()
// Reserve the per-IP and per-day slots before creating the install, so a
// limit that is hit never leaves a registration behind.
const ipDenied = await reserveSlot(store, `${REGISTER_IP_RATE_PREFIX}${ipHash(deps)}`, MAX_REGISTRATIONS_PER_IP_PER_WINDOW, now)
if (ipDenied) return ipDenied
const dayKey = `${REGISTER_DAY_PREFIX}${new Date(now).toISOString().slice(0, ISO_DATE_LENGTH)}`
const day = (await store.get(dayKey, { type: 'json' })) as DayCounter | null
const dayCount = day && typeof day.count === 'number' ? day.count : 0
if (dayCount >= MAX_REGISTRATIONS_PER_DAY) return RATE_LIMITED()
const dayDenied = await compareAndSwap(store, dayKey, (current) => {
const count = current && typeof (current as DayCounter).count === 'number' ? (current as DayCounter).count : 0
if (count >= MAX_REGISTRATIONS_PER_DAY) return { reject: RATE_LIMITED() }
return { next: { count: count + 1 } satisfies DayCounter }
})
if (dayDenied) return dayDenied

const record: InstallRecord = {
public_key: reg.public_key,
...(reg.hive_version ? { hive_version: reg.hive_version } : {}),
registered_at: new Date(now).toISOString(),
}
await store.setJSON(`${INSTALL_PREFIX}${reg.install_id}`, record)
await store.setJSON(ipKey, { times: [...ipTimes, now] })
await store.setJSON(dayKey, { count: dayCount + 1 })
// Create-only: two registrations racing for the same install id cannot both
// win, so the id is bound to exactly one key — the first write's.
const created = await store.setJSON(`${INSTALL_PREFIX}${reg.install_id}`, record, { onlyIfNew: true })
if (!created.modified) {
const winner = await getInstall(store, reg.install_id)
if (winner && winner.public_key === reg.public_key) return json(200, { ok: true, registered: false })
return json(409, { error: 'install_id is registered to a different key' })
}
return json(201, { ok: true, registered: true })
}

/**
* Records a nonce for an install, rejecting a reuse. Returns null when the
* nonce is fresh, else the response to send.
* nonce is fresh, else the response to send. Check and record are one
* conditional write, so two copies of the same request cannot both pass.
*/
async function consumeNonce(store: RelayStore, installId: string, nonce: string, now: number): Promise<Response | null> {
const key = `${NONCE_PREFIX}${installId}`
const rec = (await store.get(key, { type: 'json' })) as NonceRecord | null
const live = (rec && Array.isArray(rec.seen) ? rec.seen : []).filter(
(s) => s && typeof s.n === 'string' && typeof s.t === 'number' && now - s.t < NONCE_RETENTION_MS,
)
if (live.some((s) => s.n === nonce)) return json(401, { error: 'replayed request' })
if (live.length >= MAX_NONCES_PER_INSTALL) return RATE_LIMITED()
await store.setJSON(key, { seen: [...live, { n: nonce, t: now }] })
return null
function consumeNonce(store: RelayStore, installId: string, nonce: string, now: number): Promise<Response | null> {
return compareAndSwap(store, `${NONCE_PREFIX}${installId}`, (current) => {
const rec = current as NonceRecord | null
const live = (rec && Array.isArray(rec.seen) ? rec.seen : []).filter(
(s) => s && typeof s.n === 'string' && typeof s.t === 'number' && now - s.t < NONCE_RETENTION_MS,
)
if (live.some((s) => s.n === nonce)) return { reject: json(401, { error: 'replayed request' }) }
if (live.length >= MAX_NONCES_PER_INSTALL) return { reject: RATE_LIMITED() }
return { next: { seen: [...live, { n: nonce, t: now }] } satisfies NonceRecord }
})
}

async function handleSubmit(req: Request, deps: RelayDeps, now: number): Promise<Response> {
Expand All @@ -585,18 +657,25 @@ async function handleSubmit(req: Request, deps: RelayDeps, now: number): Promise
const replay = await consumeNonce(store, signed.installId, signed.nonce, now)
if (replay) return replay

// Cheap pre-checks keep today's response order (429 before 400) without
// reserving anything; the reservations below are what actually bound.
const ipKey = `${IP_RATE_PREFIX}${ipHash(deps)}`
const ipTimes = await recentTimes(store, ipKey, now)
if (ipTimes.length >= MAX_PER_IP_PER_WINDOW) return RATE_LIMITED()
if (await atLimit(store, ipKey, MAX_PER_IP_PER_WINDOW, now)) return RATE_LIMITED()

const body = parseJson(bytes)
if (body === undefined) return json(400, { error: 'invalid JSON body' })
const parsed = validateSubmission(body)
if (!parsed.ok) return json(400, { error: parsed.error })

const installKey = `${INSTALL_RATE_PREFIX}${signed.installId}`
const installTimes = await recentTimes(store, installKey, now)
if (installTimes.length >= MAX_PER_INSTALL_PER_WINDOW) return RATE_LIMITED()
if (await atLimit(store, installKey, MAX_PER_INSTALL_PER_WINDOW, now)) return RATE_LIMITED()

// Reserve both slots atomically before storing the entry, so concurrent
// submissions cannot all squeeze into the same remaining slot.
const ipDenied = await reserveSlot(store, ipKey, MAX_PER_IP_PER_WINDOW, now)
if (ipDenied) return ipDenied
const installDenied = await reserveSlot(store, installKey, MAX_PER_INSTALL_PER_WINDOW, now)
if (installDenied) return installDenied

const entry: RelayEntry = {
id: newEntryId(now),
Expand All @@ -608,8 +687,6 @@ async function handleSubmit(req: Request, deps: RelayDeps, now: number): Promise
timestamp: new Date(now).toISOString(),
}
await store.setJSON(`${ENTRY_PREFIX}${entry.id}`, entry)
await store.setJSON(ipKey, { times: [...ipTimes, now] })
await store.setJSON(installKey, { times: [...installTimes, now] })

// Rolling window: drop the oldest entries beyond the cap.
const keys = await entryKeysOldestFirst(store)
Expand Down
Loading
Loading