From bb0e54dcee16fb29bd0bbf55b5e66b0cdc66f138 Mon Sep 17 00:00:00 2001 From: Yun Chan Date: Wed, 16 Sep 2026 23:24:44 +0900 Subject: [PATCH] feat(server): deliver push without a Firebase project Every notification depended on Firebase Cloud Messaging, so a missing Firebase project, which is the current state, meant no notification could be delivered on any platform. Web Push and token-based Apple Push are now first class transports alongside FCM, chosen per registered device, and a scheduled Cloudflare Worker drain retries an outbox so a provider outage delays rather than drops a message. --- .env.example | 17 + server/cloudflare-worker/src/index.ts | 30 ++ .../cloudflare-worker/src/push-drain.test.ts | 69 ++++ server/cloudflare-worker/src/push-drain.ts | 84 +++++ server/cloudflare-worker/wrangler.toml | 8 + .../supabase/functions/_shared/apns.test.ts | 189 +++++++++++ server/supabase/functions/_shared/apns.ts | 199 +++++++++++ .../functions/_shared/push-contract.test.ts | 7 + .../functions/_shared/push-contract.ts | 4 +- .../functions/_shared/webpush.test.ts | 251 ++++++++++++++ server/supabase/functions/_shared/webpush.ts | 308 ++++++++++++++++++ server/supabase/functions/send-push/index.ts | 139 +++++--- .../20260913000033_team_activities.sql | 148 +++++++++ 13 files changed, 1412 insertions(+), 41 deletions(-) create mode 100644 server/cloudflare-worker/src/push-drain.test.ts create mode 100644 server/cloudflare-worker/src/push-drain.ts create mode 100644 server/supabase/functions/_shared/apns.test.ts create mode 100644 server/supabase/functions/_shared/apns.ts create mode 100644 server/supabase/functions/_shared/webpush.test.ts create mode 100644 server/supabase/functions/_shared/webpush.ts create mode 100644 server/supabase/migrations/20260913000033_team_activities.sql diff --git a/.env.example b/.env.example index 9c310cc..f181818 100644 --- a/.env.example +++ b/.env.example @@ -70,8 +70,25 @@ D3RO_ADMOB_REWARDED_UNIT_ID= # Email and push providers RESEND_FROM= RESEND_API_KEY= +# Android push transport (Firebase Cloud Messaging). Required only for Android. FCM_SERVICE_ACCOUNT_JSON= FCM_PROJECT_ID= +# Web Push transport (Firebase-free). VAPID keys; subject must be mailto: or https:. +WEBPUSH_VAPID_PUBLIC_KEY= +WEBPUSH_VAPID_PRIVATE_KEY= +WEBPUSH_SUBJECT= +# Apple Push Notification service transport (Firebase-free). Token-based .p8 auth. +APNS_KEY_ID= +APNS_TEAM_ID= +APNS_PRIVATE_KEY= +APNS_TOPIC=com.d3ro.voice +APNS_ENVIRONMENT=production + +# Cloudflare Worker push-outbox cron drain. Set SUPABASE_SERVICE_ROLE_KEY as a +# worker secret (`wrangler secret put`), never in this file. +D3RO_WORKER_SUPABASE_URL= +D3RO_WORKER_SUPABASE_SERVICE_ROLE_KEY= +D3RO_WORKER_PUSH_DRAIN_BATCH_LIMIT=20 # Optional authenticated mobile E2E account, stored only in CI secrets or local .env. MOBILE_E2E_EMAIL= diff --git a/server/cloudflare-worker/src/index.ts b/server/cloudflare-worker/src/index.ts index acb011f..c1d2012 100644 --- a/server/cloudflare-worker/src/index.ts +++ b/server/cloudflare-worker/src/index.ts @@ -1,9 +1,17 @@ // server/cloudflare-worker/src/index.ts // D3RO Voice — Global Cloudflare Edge Gateway & Proxy Worker +import { drainPushDispatches } from './push-drain.ts' + export interface Env { BACKEND_ORIGIN: string SERVICE_NAME?: string + /** Supabase project base URL, e.g. https://.supabase.co */ + SUPABASE_URL?: string + /** Supabase service-role key (secret). Used only for the scheduled drain. */ + SUPABASE_SERVICE_ROLE_KEY?: string + /** Optional drain batch size, 1-100 (default 20). */ + PUSH_DRAIN_BATCH_LIMIT?: string } const CORS_HEADERS: Record = { @@ -93,4 +101,26 @@ export default { ) } }, + + // Cron trigger: drain the Supabase push outbox. Nothing else in the + // repository scheduled this, so transactionally enqueued notifications + // (team invites, transcription completion, billing) never left the queue. + async scheduled( + _event: unknown, + env: Env, + ctx: { waitUntil(promise: Promise): void }, + ): Promise { + ctx.waitUntil(drainAndLog(env)) + }, +} + +async function drainAndLog(env: Env): Promise { + const result = await drainPushDispatches(env) + if (result.ok) { + console.log( + `[push-drain] claimed=${result.claimed ?? 0} completed=${result.completed ?? 0} pending=${result.pending ?? 0}`, + ) + } else { + console.error(`[push-drain] failed status=${result.status} error=${result.error ?? 'unknown'}`) + } } diff --git a/server/cloudflare-worker/src/push-drain.test.ts b/server/cloudflare-worker/src/push-drain.test.ts new file mode 100644 index 0000000..931322f --- /dev/null +++ b/server/cloudflare-worker/src/push-drain.test.ts @@ -0,0 +1,69 @@ +// server/cloudflare-worker/src/push-drain.test.ts +import { drainPushDispatches, resolveBatchLimit } from './push-drain.ts' + +function assert(condition: boolean, message: string): asserts condition { + if (!condition) throw new Error(message) +} + +const env = { SUPABASE_URL: 'https://example.supabase.co', SUPABASE_SERVICE_ROLE_KEY: 'x'.repeat(40) } + +function jsonResponse(status: number, body: unknown): Response { + return new Response(JSON.stringify(body), { + status, + headers: { 'Content-Type': 'application/json' }, + }) +} + +Deno.test('batch limit is clamped to the edge function contract', () => { + assert(resolveBatchLimit(undefined) === 20, 'default is 20') + assert(resolveBatchLimit('5') === 5, 'explicit value is used') + assert(resolveBatchLimit('9999') === 100, 'upper bound is 100') + assert(resolveBatchLimit('0') === 20, 'zero falls back to the default') + assert(resolveBatchLimit('abc') === 20, 'garbage falls back to the default') +}) + +Deno.test('drain fails closed without configuration', async () => { + const result = await drainPushDispatches({}, () => { + throw new Error('must not be called') + }) + assert(result.ok === false, 'not ok') + assert(result.error === 'push_drain_not_configured', 'clear configuration error') + assert(result.status === 503, 'service unavailable') +}) + +Deno.test('drain posts to the drain endpoint with the service role bearer', async () => { + const calls: Array<{ url: string; init: RequestInit }> = [] + const fetchImpl = ((url: string | URL | Request, init?: RequestInit) => { + calls.push({ url: String(url), init: init ?? {} }) + return Promise.resolve(jsonResponse(200, { claimed: 3, completed: 2, pending: 1 })) + }) as unknown as typeof fetch + + const result = await drainPushDispatches(env, fetchImpl) + assert(result.ok === true, 'ok') + assert(result.claimed === 3 && result.completed === 2 && result.pending === 1, 'summary parsed') + assert(calls[0].url === 'https://example.supabase.co/functions/v1/send-push?mode=drain&limit=20', 'drain url') + const headers = calls[0].init.headers as Record + assert(headers.Authorization === `Bearer ${env.SUPABASE_SERVICE_ROLE_KEY}`, 'service role bearer') + assert(headers['Content-Type'] === 'application/json', 'json content type') + assert(calls[0].init.body === '{}', 'empty json body') + assert(calls[0].init.method === 'POST', 'post method') +}) + +Deno.test('drain reports upstream failures without throwing', async () => { + const httpFail = ((() => Promise.resolve(new Response('nope', { status: 500 }))) as unknown as typeof fetch) + const result = await drainPushDispatches(env, httpFail) + assert(result.ok === false, 'not ok') + assert(result.error === 'push_drain_http_500', 'http error surfaced') + assert(result.status === 500, 'status surfaced') + + const unreachable = (() => Promise.reject(new Error('network'))) as unknown as typeof fetch + const unreachableResult = await drainPushDispatches(env, unreachable) + assert(unreachableResult.error === 'push_drain_unreachable', 'network error surfaced') +}) + +Deno.test('drain tolerates a malformed success body', async () => { + const fetchImpl = ((() => Promise.resolve(new Response('not-json', { status: 200 }))) as unknown as typeof fetch) + const result = await drainPushDispatches(env, fetchImpl) + assert(result.ok === true, 'ok') + assert(result.claimed === undefined, 'summary left undefined') +}) diff --git a/server/cloudflare-worker/src/push-drain.ts b/server/cloudflare-worker/src/push-drain.ts new file mode 100644 index 0000000..226fe2d --- /dev/null +++ b/server/cloudflare-worker/src/push-drain.ts @@ -0,0 +1,84 @@ +// server/cloudflare-worker/src/push-drain.ts +// Cloudflare Cron-triggered drain for the Supabase push outbox. +// +// Supabase owns the durable outbox and the send-push Edge Function; nothing in +// the repository scheduled the drain, so transactionally enqueued notifications +// (team invites, transcription completion, billing) never left the queue. This +// worker drains them on a cron schedule. + +export interface PushDrainEnv { + SUPABASE_URL?: string + SUPABASE_SERVICE_ROLE_KEY?: string + PUSH_DRAIN_BATCH_LIMIT?: string +} + +export interface PushDrainResult { + ok: boolean + status: number + claimed?: number + completed?: number + pending?: number + error?: string +} + +const DEFAULT_BATCH_LIMIT = 20 +const MAX_BATCH_LIMIT = 100 +const REQUEST_TIMEOUT_MS = 30_000 + +export function resolveBatchLimit(raw: string | undefined): number { + const parsed = Number(raw ?? DEFAULT_BATCH_LIMIT) + if (!Number.isSafeInteger(parsed) || parsed < 1) return DEFAULT_BATCH_LIMIT + return Math.min(parsed, MAX_BATCH_LIMIT) +} + +export async function drainPushDispatches( + env: PushDrainEnv, + fetchImpl: typeof fetch = fetch, +): Promise { + const baseUrl = env.SUPABASE_URL?.trim().replace(/\/+$/, '') ?? '' + const serviceRoleKey = env.SUPABASE_SERVICE_ROLE_KEY?.trim() ?? '' + if (!baseUrl || serviceRoleKey.length < 20) { + return { ok: false, status: 503, error: 'push_drain_not_configured' } + } + + const limit = resolveBatchLimit(env.PUSH_DRAIN_BATCH_LIMIT) + let response: Response + try { + response = await fetchImpl( + `${baseUrl}/functions/v1/send-push?mode=drain&limit=${limit}`, + { + method: 'POST', + headers: { + Authorization: `Bearer ${serviceRoleKey}`, + 'Content-Type': 'application/json', + }, + body: '{}', + signal: AbortSignal.timeout(REQUEST_TIMEOUT_MS), + }, + ) + } catch { + return { ok: false, status: 502, error: 'push_drain_unreachable' } + } + + if (!response.ok) { + return { ok: false, status: response.status, error: `push_drain_http_${response.status}` } + } + + let body: Record = {} + try { + const parsed = await response.json() + if (parsed !== null && typeof parsed === 'object' && !Array.isArray(parsed)) { + body = parsed as Record + } + } catch { + body = {} + } + + return { + ok: true, + status: response.status, + claimed: typeof body.claimed === 'number' ? body.claimed : undefined, + completed: typeof body.completed === 'number' ? body.completed : undefined, + pending: typeof body.pending === 'number' ? body.pending : undefined, + } +} diff --git a/server/cloudflare-worker/wrangler.toml b/server/cloudflare-worker/wrangler.toml index c4c880e..b586687 100644 --- a/server/cloudflare-worker/wrangler.toml +++ b/server/cloudflare-worker/wrangler.toml @@ -5,6 +5,14 @@ compatibility_date = "2024-04-01" [vars] BACKEND_ORIGIN = "http://192.168.0.39:5050" SERVICE_NAME = "D3RO Voice Cloud Gateway" +# Supabase project base URL for the scheduled push-outbox drain. +# SUPABASE_URL = "https://.supabase.co" +# PUSH_DRAIN_BATCH_LIMIT = "20" + +# Scheduled drain of the Supabase push outbox (every minute). +# Requires the secret: wrangler secret put SUPABASE_SERVICE_ROLE_KEY +[triggers] +crons = ["* * * * *"] # Cloudflare Custom Domain / Route Binding (선택 시 활성화) # routes = [ diff --git a/server/supabase/functions/_shared/apns.test.ts b/server/supabase/functions/_shared/apns.test.ts new file mode 100644 index 0000000..ae5bd42 --- /dev/null +++ b/server/supabase/functions/_shared/apns.test.ts @@ -0,0 +1,189 @@ +// server/supabase/functions/_shared/apns.test.ts +// Apple Push Notification service (token/.p8) transport tests. Firebase-free. + +import { + createApnsToken, + readApnsConfig, + sendApnsMessage, + type ApnsConfig, +} from './apns.ts' +import { PushContractError, type PushNotification } from './push-contract.ts' + +function assert(condition: boolean, message: string): asserts condition { + if (!condition) throw new Error(message) +} + +async function assertPushError( + action: () => unknown | Promise, + code: string, + status: number, + stale?: boolean, +): Promise { + let actual: unknown + try { + await action() + } catch (error) { + actual = error + } + assert(actual instanceof PushContractError, `expected PushContractError for ${code}`) + assert(actual.code === code, `expected ${code}, received ${actual.code}`) + assert(actual.status === status, `expected status ${status}, received ${actual.status}`) + if (stale !== undefined) { + assert(actual.staleRegistration === stale, `expected stale=${stale}`) + } +} + +function toPem(bytes: Uint8Array): string { + let binary = '' + for (const byte of bytes) binary += String.fromCharCode(byte) + const b64 = btoa(binary) + const lines = b64.match(/.{1,64}/g)?.join('\n') ?? b64 + return `-----BEGIN PRIVATE KEY-----\n${lines}\n-----END PRIVATE KEY-----\n` +} + +async function makeP8(): Promise { + const keys = await crypto.subtle.generateKey( + { name: 'ECDSA', namedCurve: 'P-256' }, + true, + ['sign'], + ) + return toPem(new Uint8Array(await crypto.subtle.exportKey('pkcs8', keys.privateKey))) +} + +function b64urlDecode(value: string): string { + const normalized = value.replace(/-/g, '+').replace(/_/g, '/') + const padded = normalized + '='.repeat((4 - (normalized.length % 4)) % 4) + return new TextDecoder().decode(Uint8Array.from(atob(padded), (character) => character.charCodeAt(0))) +} + +const DEVICE_TOKEN = 'a'.repeat(64) +const notification: PushNotification = { + title: 'Transcription complete', + body: 'Open D3RO Voice', + data: { schema_version: '1', event_type: 'transcription.completed', resource_id: '11111111-2222-4333-8444-555555555555', route: 'HistoryDetail', history_id: '11111111-2222-4333-8444-555555555555' }, +} + +Deno.test('apns config requires valid credentials', () => { + const empty = () => undefined + assertPushError(() => readApnsConfig(empty), 'apns_not_configured', 503) + + const malformed = (name: string) => ({ + APNS_KEY_ID: 'too-short', + APNS_TEAM_ID: 'ABCDEFGHIJ', + APNS_PRIVATE_KEY: '-----BEGIN PRIVATE KEY-----\nAAAA\n-----END PRIVATE KEY-----', + APNS_TOPIC: 'com.d3ro.voice', + })[name] + assertPushError(() => readApnsConfig(malformed), 'apns_credentials_invalid', 503) +}) + +Deno.test('apns provider token carries kid and ES256', async () => { + const pem = await makeP8() + const config: ApnsConfig = { + keyId: 'ABCDEFGHIJ', + teamId: 'KLMNOPQRST', + privateKey: pem, + topic: 'com.d3ro.voice', + environment: 'production', + } + const token = await createApnsToken(config, 1_700_000_000) + const [header, claims] = token.split('.') + const decodedHeader = JSON.parse(b64urlDecode(header)) as Record + const decodedClaims = JSON.parse(b64urlDecode(claims)) as Record + assert(decodedHeader.alg === 'ES256', 'alg must be ES256') + assert(decodedHeader.kid === 'ABCDEFGHIJ', 'kid must be the key id') + assert(decodedClaims.iss === 'KLMNOPQRST', 'iss must be the team id') + assert(typeof decodedClaims.iat === 'number', 'iat must be set') + assert(token.split('.').length === 3, 'token is a JWT') +}) + +Deno.test('apns send targets the right host with topic headers', async () => { + const config: ApnsConfig = { + keyId: 'ABCDEFGHIJ', + teamId: 'KLMNOPQRST', + privateKey: 'unused', + topic: 'com.d3ro.voice', + environment: 'production', + } + const calls: Array<{ url: string; init: RequestInit }> = [] + const fetchImpl = ((url: string | URL | Request, init?: RequestInit) => { + calls.push({ url: String(url), init: init ?? {} }) + return Promise.resolve(new Response('{}', { status: 200 })) + }) as unknown as typeof fetch + + const result = await sendApnsMessage(DEVICE_TOKEN, notification, { + config, + fetchImpl, + getAccessToken: () => Promise.resolve('test-token'), + }) + assert(result.status === 200, 'ok status returned') + assert(calls.length === 1, 'one request issued') + assert(calls[0].url === `https://api.push.apple.com/3/device/${DEVICE_TOKEN}`, 'production host used') + const headers = calls[0].init.headers as Record + assert(headers['apns-topic'] === 'com.d3ro.voice', 'topic header present') + assert(headers['apns-push-type'] === 'alert', 'push type header present') + assert(headers.authorization === 'bearer test-token', 'bearer token present') +}) + +Deno.test('apns send maps stale and auth responses', async () => { + const config: ApnsConfig = { + keyId: 'ABCDEFGHIJ', + teamId: 'KLMNOPQRST', + privateKey: 'unused', + topic: 'com.d3ro.voice', + environment: 'sandbox', + } + const respond = (status: number, body: unknown) => ( + (() => Promise.resolve(new Response(JSON.stringify(body), { status }))) as unknown as typeof fetch + ) + const options = { config, getAccessToken: () => Promise.resolve('t') } + + await assertPushError( + () => sendApnsMessage(DEVICE_TOKEN, notification, { ...options, fetchImpl: respond(410, { reason: 'Unregistered' }) }), + 'apns_registration_stale', + 410, + true, + ) + await assertPushError( + () => sendApnsMessage(DEVICE_TOKEN, notification, { ...options, fetchImpl: respond(400, { reason: 'BadDeviceToken' }) }), + 'apns_registration_stale', + 410, + true, + ) + await assertPushError( + () => sendApnsMessage(DEVICE_TOKEN, notification, { ...options, fetchImpl: respond(403, { reason: 'InvalidProviderToken' }) }), + 'apns_credentials_invalid', + 503, + ) + await assertPushError( + () => sendApnsMessage(DEVICE_TOKEN, notification, { ...options, fetchImpl: respond(503, { reason: 'ServiceUnavailable' }) }), + 'apns_send_failed', + 502, + ) + await assertPushError( + () => sendApnsMessage('not-a-token', notification, { ...options, fetchImpl: respond(200, {}) }), + 'apns_registration_invalid', + 400, + true, + ) +}) + +Deno.test('apns send uses the sandbox host when configured', async () => { + const config: ApnsConfig = { + keyId: 'ABCDEFGHIJ', + teamId: 'KLMNOPQRST', + privateKey: 'unused', + topic: 'com.d3ro.voice', + environment: 'sandbox', + } + const calls: string[] = [] + const fetchImpl = ((url: string | URL | Request) => { + calls.push(String(url)) + return Promise.resolve(new Response('{}', { status: 200 })) + }) as unknown as typeof fetch + await sendApnsMessage(DEVICE_TOKEN, notification, { + config, + fetchImpl, + getAccessToken: () => Promise.resolve('t'), + }) + assert(calls[0].startsWith('https://api.sandbox.push.apple.com/'), 'sandbox host used') +}) diff --git a/server/supabase/functions/_shared/apns.ts b/server/supabase/functions/_shared/apns.ts new file mode 100644 index 0000000..cc2f840 --- /dev/null +++ b/server/supabase/functions/_shared/apns.ts @@ -0,0 +1,199 @@ +// server/supabase/functions/_shared/apns.ts +// Apple Push Notification service transport (token-based, .p8 / ES256 JWT). +// +// Firebase-free: APNs is Apple's own service. Only an Apple Developer key is +// required; no Firebase project is involved. + +import { PushContractError, type PushNotification } from './push-contract.ts' + +export type ApnsEnvironment = 'production' | 'sandbox' + +export interface ApnsConfig { + keyId: string + teamId: string + privateKey: string + topic: string + environment: ApnsEnvironment +} + +export interface ApnsSendOptions { + fetchImpl?: typeof fetch + config?: ApnsConfig + now?: number + getAccessToken?: (config: ApnsConfig, nowSeconds: number) => Promise +} + +export interface ApnsSendResult { + status: number +} + +const KEY_ID_PATTERN = /^[A-Z0-9]{10}$/ +const TOPIC_PATTERN = /^[A-Za-z0-9.-]{3,155}$/ +const DEVICE_TOKEN_PATTERN = /^[0-9a-fA-F]{32,200}$/ +const encoder = new TextEncoder() +const REQUEST_TIMEOUT_MS = 10_000 +const TOKEN_TTL_SECONDS = 50 * 60 + +let cachedToken: { identity: string; token: string; issuedAtSeconds: number } | null = null + +export function readApnsConfig( + readEnv: (name: string) => string | undefined = (name) => Deno.env.get(name), +): ApnsConfig { + const keyId = readEnv('APNS_KEY_ID')?.trim() ?? '' + const teamId = readEnv('APNS_TEAM_ID')?.trim() ?? '' + const privateKey = readEnv('APNS_PRIVATE_KEY') ?? '' + const topic = readEnv('APNS_TOPIC')?.trim() ?? '' + const environmentRaw = readEnv('APNS_ENVIRONMENT')?.trim().toLowerCase() ?? 'production' + if (!keyId || !teamId || !privateKey || !topic) { + throw new PushContractError('apns_not_configured', 503) + } + if ( + !KEY_ID_PATTERN.test(keyId) + || !KEY_ID_PATTERN.test(teamId) + || !TOPIC_PATTERN.test(topic) + || !privateKey.includes('BEGIN PRIVATE KEY') + || (environmentRaw !== 'production' && environmentRaw !== 'sandbox') + ) { + throw new PushContractError('apns_credentials_invalid', 503) + } + return { + keyId, + teamId, + privateKey, + topic, + environment: environmentRaw === 'sandbox' ? 'sandbox' : 'production', + } +} + +function decodePemPrivateKey(pem: string): Uint8Array { + const body = pem + .replace(/-----BEGIN PRIVATE KEY-----/g, '') + .replace(/-----END PRIVATE KEY-----/g, '') + .replace(/\s/g, '') + if (!body) throw new PushContractError('apns_credentials_invalid', 503) + try { + const binary = atob(body) + return Uint8Array.from(binary, (character) => character.charCodeAt(0)) + } catch { + throw new PushContractError('apns_credentials_invalid', 503) + } +} + +export async function createApnsToken(config: ApnsConfig, nowSeconds: number): Promise { + const header = base64Url(encoder.encode(JSON.stringify({ alg: 'ES256', kid: config.keyId }))) + const claims = base64Url(encoder.encode(JSON.stringify({ iss: config.teamId, iat: nowSeconds }))) + const signingInput = `${header}.${claims}` + try { + const decoded = decodePemPrivateKey(config.privateKey) + const keyBytes = new ArrayBuffer(decoded.byteLength) + new Uint8Array(keyBytes).set(decoded) + const key = await crypto.subtle.importKey( + 'pkcs8', + keyBytes, + { name: 'ECDSA', namedCurve: 'P-256' }, + false, + ['sign'], + ) + const signature = await crypto.subtle.sign( + { name: 'ECDSA', hash: 'SHA-256' }, + key, + encoder.encode(signingInput), + ) + return `${signingInput}.${base64Url(new Uint8Array(signature))}` + } catch (error) { + if (error instanceof PushContractError) throw error + throw new PushContractError('apns_credentials_invalid', 503) + } +} + +function base64Url(bytes: Uint8Array): string { + let binary = '' + for (const byte of bytes) binary += String.fromCharCode(byte) + return btoa(binary).replace(/\+/g, '-').replace(/\//g, '_').replace(/=+$/g, '') +} + +export async function getApnsAccessToken( + config: ApnsConfig, + nowSeconds: number, +): Promise { + const identity = `${config.keyId}\n${config.teamId}` + if ( + cachedToken + && cachedToken.identity === identity + && cachedToken.issuedAtSeconds + TOKEN_TTL_SECONDS > nowSeconds + ) { + return cachedToken.token + } + const token = await createApnsToken(config, nowSeconds) + cachedToken = { identity, token, issuedAtSeconds: nowSeconds } + return token +} + +function apnsHost(environment: ApnsEnvironment): string { + return environment === 'sandbox' ? 'https://api.sandbox.push.apple.com' : 'https://api.push.apple.com' +} + +function isStaleReason(reason: string): boolean { + return reason === 'BadDeviceToken' + || reason === 'DeviceTokenNotForTopic' + || reason === 'Unregistered' + || reason === 'ExpiredToken' +} + +export async function sendApnsMessage( + registrationId: string, + notification: PushNotification, + options: ApnsSendOptions = {}, +): Promise { + if (typeof registrationId !== 'string' || !DEVICE_TOKEN_PATTERN.test(registrationId)) { + throw new PushContractError('apns_registration_invalid', 400, true) + } + const config = options.config ?? readApnsConfig() + const fetchImpl = options.fetchImpl ?? fetch + const nowSeconds = Math.floor((options.now ?? Date.now()) / 1000) + const accessToken = await (options.getAccessToken ?? getApnsAccessToken)(config, nowSeconds) + + const body = JSON.stringify({ + aps: { + alert: { title: notification.title, body: notification.body }, + sound: 'default', + }, + ...notification.data, + }) + + let response: Response + try { + response = await fetchImpl(`${apnsHost(config.environment)}/3/device/${registrationId}`, { + method: 'POST', + headers: { + authorization: `bearer ${accessToken}`, + 'apns-topic': config.topic, + 'apns-push-type': 'alert', + 'apns-priority': '10', + 'content-type': 'application/json', + }, + body, + signal: AbortSignal.timeout(REQUEST_TIMEOUT_MS), + }) + } catch { + throw new PushContractError('apns_send_timeout', 504) + } + + if (response.ok) return { status: response.status } + + let reason = '' + try { + const parsed = await response.json() as Record + reason = typeof parsed.reason === 'string' ? parsed.reason : '' + } catch { + reason = '' + } + + if (response.status === 410 || (response.status === 400 && isStaleReason(reason))) { + throw new PushContractError('apns_registration_stale', 410, true) + } + if (response.status === 403) { + throw new PushContractError('apns_credentials_invalid', 503) + } + throw new PushContractError('apns_send_failed', 502) +} diff --git a/server/supabase/functions/_shared/push-contract.test.ts b/server/supabase/functions/_shared/push-contract.test.ts index bef4d7f..b3849f8 100644 --- a/server/supabase/functions/_shared/push-contract.test.ts +++ b/server/supabase/functions/_shared/push-contract.test.ts @@ -3,6 +3,7 @@ import { isServiceRoleAuthorization, isServiceRoleApiKey, parsePushRequest, + providerIsSupported, PushContractError, readFcmConfig, sendFcmMessage, @@ -32,6 +33,12 @@ const fcmConfig = { tokenUri: 'https://oauth2.googleapis.com/token', } +Deno.test('push providers fcm, webpush, and apns are supported', () => { + assert(providerIsSupported('fcm'), 'fcm must be supported') + assert(providerIsSupported('webpush'), 'webpush must be supported') + assert(providerIsSupported('apns'), 'apns must be supported') +}) + Deno.test('push input accepts only event_type and resource_id', () => { const parsed = parsePushRequest({ event_type: 'transcription.completed', diff --git a/server/supabase/functions/_shared/push-contract.ts b/server/supabase/functions/_shared/push-contract.ts index 00d5008..7060391 100644 --- a/server/supabase/functions/_shared/push-contract.ts +++ b/server/supabase/functions/_shared/push-contract.ts @@ -158,7 +158,7 @@ function androidDelivery(eventType: PushEventType, resourceId: string): { } } -function validateOutboundData(data: Record): { +export function validateOutboundData(data: Record): { eventType: PushEventType resourceId: string } { @@ -475,7 +475,7 @@ export async function sendFcmMessage( } export function providerIsSupported(provider: PushProvider): boolean { - return provider === 'fcm' + return provider === 'fcm' || provider === 'webpush' || provider === 'apns' } async function credentialsEqual(presented: string, expected: string): Promise { diff --git a/server/supabase/functions/_shared/webpush.test.ts b/server/supabase/functions/_shared/webpush.test.ts new file mode 100644 index 0000000..101b092 --- /dev/null +++ b/server/supabase/functions/_shared/webpush.test.ts @@ -0,0 +1,251 @@ +// server/supabase/functions/_shared/webpush.test.ts +// RFC 8291 (aes128gcm) + VAPID transport tests. Firebase-free web push. + +import { + encryptWebPushPayload, + parseWebPushSubscription, + readWebPushConfig, + sendWebPushMessage, + type WebPushConfig, + type WebPushSubscription, +} from './webpush.ts' +import { PushContractError, type PushNotification } from './push-contract.ts' + +function assert(condition: boolean, message: string): asserts condition { + if (!condition) throw new Error(message) +} + +async function assertPushError( + action: () => unknown | Promise, + code: string, + status: number, + stale?: boolean, +): Promise { + let actual: unknown + try { + await action() + } catch (error) { + actual = error + } + assert(actual instanceof PushContractError, `expected PushContractError for ${code}`) + assert(actual.code === code, `expected ${code}, received ${actual.code}`) + assert(actual.status === status, `expected status ${status}, received ${actual.status}`) + if (stale !== undefined) { + assert(actual.staleRegistration === stale, `expected stale=${stale}`) + } +} + +const encoder = new TextEncoder() + +function b64url(bytes: Uint8Array): string { + let binary = '' + for (const byte of bytes) binary += String.fromCharCode(byte) + return btoa(binary).replace(/\+/g, '-').replace(/\//g, '_').replace(/=+$/g, '') +} + +function unb64url(value: string): Uint8Array { + const normalized = value.replace(/-/g, '+').replace(/_/g, '/') + const padded = normalized + '='.repeat((4 - (normalized.length % 4)) % 4) + const binary = atob(padded) + return Uint8Array.from(binary, (character) => character.charCodeAt(0)) +} + +function concat(...parts: Uint8Array[]): Uint8Array { + const output = new Uint8Array(parts.reduce((sum, part) => sum + part.byteLength, 0)) + let offset = 0 + for (const part of parts) { + output.set(part, offset) + offset += part.byteLength + } + return output +} + +function toArrayBuffer(bytes: Uint8Array): ArrayBuffer { + const buffer = new ArrayBuffer(bytes.byteLength) + new Uint8Array(buffer).set(bytes) + return buffer +} + +async function hmacSha256(key: Uint8Array, data: Uint8Array): Promise { + const imported = await crypto.subtle.importKey( + 'raw', + toArrayBuffer(key), + { name: 'HMAC', hash: 'SHA-256' }, + false, + ['sign'], + ) + return new Uint8Array(await crypto.subtle.sign('HMAC', imported, toArrayBuffer(data))) +} + +async function hkdfExpand(prk: Uint8Array, info: Uint8Array, length: number): Promise { + const blocks = Math.ceil(length / 32) + const output = new Uint8Array(blocks * 32) + let previous = new Uint8Array(0) + for (let index = 0; index < blocks; index += 1) { + const block = await hmacSha256(prk, concat(previous, info, Uint8Array.of(index + 1))) + output.set(block, index * 32) + previous = new Uint8Array(block) + } + return output.slice(0, length) +} + +async function makeVapidConfig(): Promise { + const keys = await crypto.subtle.generateKey( + { name: 'ECDSA', namedCurve: 'P-256' }, + true, + ['sign', 'verify'], + ) + const publicRaw = new Uint8Array(await crypto.subtle.exportKey('raw', keys.publicKey)) + const jwk = await crypto.subtle.exportKey('jwk', keys.privateKey) + assert(typeof jwk.d === 'string', 'private scalar must be exportable') + return { + publicKey: b64url(publicRaw), + privateKey: jwk.d, + subject: 'mailto:push@d3ro.test', + } +} + +async function makeSubscription(): Promise<{ registrationId: string; subscription: WebPushSubscription; privateKey: CryptoKey; publicRaw: Uint8Array }> { + const uaKeys = await crypto.subtle.generateKey( + { name: 'ECDH', namedCurve: 'P-256' }, + true, + ['deriveBits'], + ) + const publicRaw = new Uint8Array(await crypto.subtle.exportKey('raw', uaKeys.publicKey)) + const auth = crypto.getRandomValues(new Uint8Array(16)) + const subscription: WebPushSubscription = { + endpoint: 'https://push.example.com/subscriptions/abc123', + p256dh: b64url(publicRaw), + auth: b64url(auth), + } + return { + registrationId: JSON.stringify({ endpoint: subscription.endpoint, keys: { p256dh: subscription.p256dh, auth: subscription.auth } }), + subscription, + privateKey: uaKeys.privateKey, + publicRaw, + } +} + +const notification: PushNotification = { + title: 'Transcription complete', + body: 'Open D3RO Voice', + data: { schema_version: '1', event_type: 'transcription.completed', resource_id: '11111111-2222-4333-8444-555555555555', route: 'HistoryDetail', history_id: '11111111-2222-4333-8444-555555555555' }, +} + +Deno.test('web push config requires valid VAPID credentials', () => { + const empty = () => undefined + assertPushError(() => readWebPushConfig(empty), 'webpush_not_configured', 503) + const malformed = (name: string) => ({ + WEBPUSH_VAPID_PUBLIC_KEY: 'not-a-key', + WEBPUSH_VAPID_PRIVATE_KEY: 'also-bad', + WEBPUSH_SUBJECT: 'mailto:push@d3ro.test', + })[name] + assertPushError(() => readWebPushConfig(malformed), 'webpush_credentials_invalid', 503) +}) + +Deno.test('web push subscription parsing is strict', async () => { + const { registrationId, subscription } = await makeSubscription() + const parsed = parseWebPushSubscription(registrationId) + assert(parsed.endpoint === subscription.endpoint, 'endpoint preserved') + assert(parsed.p256dh === subscription.p256dh, 'p256dh preserved') + + assertPushError(() => parseWebPushSubscription('short'), 'webpush_registration_invalid', 400, true) + assertPushError( + () => parseWebPushSubscription(JSON.stringify({ endpoint: 'http://insecure.example.com', keys: { p256dh: subscription.p256dh, auth: subscription.auth } })), + 'webpush_registration_invalid', + 400, + true, + ) +}) + +Deno.test('aes128gcm encryption round-trips per RFC 8291', async () => { + const { registrationId, subscription, privateKey, publicRaw } = await makeSubscription() + const fixedSalt = crypto.getRandomValues(new Uint8Array(16)) + const payload = encoder.encode('hello web push') + const body = await encryptWebPushPayload( + parseWebPushSubscription(registrationId), + payload, + () => fixedSalt, + ) + + // Header: salt(16) | rs(4) | idlen(1) | keyid(65) + assert(body.byteLength === 16 + 4 + 1 + 65 + payload.byteLength + 1 + 16, 'body length matches aes128gcm') + const salt = body.slice(0, 16) + const recordSize = new DataView(body.buffer, body.byteOffset + 16, 4).getUint32(0, false) + assert(recordSize === 4096, 'record size is 4096') + assert(body[20] === 65, 'key id length is 65') + const asPublic = body.slice(21, 86) + const ciphertext = body.slice(86) + + assert(salt.every((byte, index) => byte === fixedSalt[index]), 'salt preserved') + + // Decrypt as the user agent would. + const asKey = await crypto.subtle.importKey( + 'raw', + toArrayBuffer(asPublic), + { name: 'ECDH', namedCurve: 'P-256' }, + false, + [], + ) + const shared = new Uint8Array( + await crypto.subtle.deriveBits({ name: 'ECDH', public: asKey }, privateKey, 256), + ) + const prkKey = await hmacSha256(unb64url(subscription.auth), shared) + const keyInfo = concat(encoder.encode('WebPush: info\x00'), publicRaw, asPublic) + const ikm = await hkdfExpand(prkKey, keyInfo, 32) + const prk = await hmacSha256(salt, ikm) + const cek = await hkdfExpand(prk, encoder.encode('Content-Encoding: aes128gcm\x00'), 16) + const nonce = await hkdfExpand(prk, encoder.encode('Content-Encoding: nonce\x00'), 12) + const aesKey = await crypto.subtle.importKey('raw', toArrayBuffer(cek), 'AES-GCM', false, ['decrypt']) + const plaintext = new Uint8Array( + await crypto.subtle.decrypt( + { name: 'AES-GCM', iv: toArrayBuffer(nonce), tagLength: 128 }, + aesKey, + toArrayBuffer(ciphertext), + ), + ) + const decoded = plaintext.slice(0, plaintext.byteLength - 1) + assert(plaintext[plaintext.byteLength - 1] === 0x02, 'last-record delimiter is 0x02') + assert(new TextDecoder().decode(decoded) === 'hello web push', 'payload round-trips') +}) + +Deno.test('web push send maps provider responses', async () => { + const { registrationId } = await makeSubscription() + const config = await makeVapidConfig() + + const okFetch = (() => Promise.resolve(new Response(null, { status: 201 }))) as unknown as typeof fetch + const result = await sendWebPushMessage(registrationId, notification, { config, fetchImpl: okFetch }) + assert(result.status === 201, 'created status returned') + + const staleFetch = (() => Promise.resolve(new Response(null, { status: 410 }))) as unknown as typeof fetch + await assertPushError( + () => sendWebPushMessage(registrationId, notification, { config, fetchImpl: staleFetch }), + 'webpush_registration_stale', + 410, + true, + ) + + const tooLarge = (() => Promise.resolve(new Response(null, { status: 413 }))) as unknown as typeof fetch + await assertPushError( + () => sendWebPushMessage(registrationId, notification, { config, fetchImpl: tooLarge }), + 'webpush_payload_too_large', + 500, + ) + + const failed = (() => Promise.resolve(new Response(null, { status: 500 }))) as unknown as typeof fetch + await assertPushError( + () => sendWebPushMessage(registrationId, notification, { config, fetchImpl: failed }), + 'webpush_send_failed', + 502, + ) +}) + +Deno.test('web push send rejects a malformed registration id', async () => { + const config = await makeVapidConfig() + await assertPushError( + () => sendWebPushMessage('not-json-registration-id', notification, { config }), + 'webpush_registration_invalid', + 400, + true, + ) +}) diff --git a/server/supabase/functions/_shared/webpush.ts b/server/supabase/functions/_shared/webpush.ts new file mode 100644 index 0000000..834e74c --- /dev/null +++ b/server/supabase/functions/_shared/webpush.ts @@ -0,0 +1,308 @@ +// server/supabase/functions/_shared/webpush.ts +// Standard Web Push (VAPID + RFC 8291 aes128gcm) transport. +// +// Firebase-free: the browser Push API delivers the message; we only sign with +// VAPID and encrypt per RFC 8291. No Google account or Firebase project needed. + +import { PushContractError, type PushNotification } from './push-contract.ts' + +export interface WebPushConfig { + /** base64url encoded uncompressed P-256 public key (65 bytes, 0x04 prefix). */ + publicKey: string + /** base64url encoded P-256 private scalar (32 bytes). */ + privateKey: string + /** VAPID subject: `mailto:` or `https:` URL. */ + subject: string +} + +export interface WebPushSubscription { + endpoint: string + p256dh: string + auth: string +} + +export interface WebPushSendOptions { + fetchImpl?: typeof fetch + config?: WebPushConfig + now?: number + randomBytes?: (length: number) => Uint8Array +} + +export interface WebPushSendResult { + status: number +} + +const SUBJECT_PATTERN = /^(mailto:[^@\s]+@[^@\s]+|https:\/\/[^\s]+)$/ +const RECORD_SIZE = 4096 +const encoder = new TextEncoder() +const REQUEST_TIMEOUT_MS = 10_000 + +function decodeBase64Url(value: string): Uint8Array { + const normalized = value.replace(/-/g, '+').replace(/_/g, '/') + const padded = normalized + '='.repeat((4 - (normalized.length % 4)) % 4) + const binary = atob(padded) + return Uint8Array.from(binary, (character) => character.charCodeAt(0)) +} + +function encodeBase64Url(bytes: Uint8Array): string { + let binary = '' + for (const byte of bytes) binary += String.fromCharCode(byte) + return btoa(binary).replace(/\+/g, '-').replace(/\//g, '_').replace(/=+$/g, '') +} + +function concat(...parts: Uint8Array[]): Uint8Array { + const total = parts.reduce((sum, part) => sum + part.byteLength, 0) + const output = new Uint8Array(total) + let offset = 0 + for (const part of parts) { + output.set(part, offset) + offset += part.byteLength + } + return output +} + +function toArrayBuffer(bytes: Uint8Array): ArrayBuffer { + const buffer = new ArrayBuffer(bytes.byteLength) + new Uint8Array(buffer).set(bytes) + return buffer +} + +async function hmacSha256(key: Uint8Array, data: Uint8Array): Promise { + const imported = await crypto.subtle.importKey( + 'raw', + toArrayBuffer(key), + { name: 'HMAC', hash: 'SHA-256' }, + false, + ['sign'], + ) + const signature = await crypto.subtle.sign('HMAC', imported, toArrayBuffer(data)) + return new Uint8Array(signature) +} + +async function hkdfExpand(prk: Uint8Array, info: Uint8Array, length: number): Promise { + const blocks = Math.ceil(length / 32) + const output = new Uint8Array(blocks * 32) + let previous = new Uint8Array(0) + for (let index = 0; index < blocks; index += 1) { + const block = await hmacSha256(prk, concat(previous, info, Uint8Array.of(index + 1))) + output.set(block, index * 32) + previous = new Uint8Array(block) + } + return output.slice(0, length) +} + +export function readWebPushConfig( + readEnv: (name: string) => string | undefined = (name) => Deno.env.get(name), +): WebPushConfig { + const publicKey = readEnv('WEBPUSH_VAPID_PUBLIC_KEY')?.trim() ?? '' + const privateKey = readEnv('WEBPUSH_VAPID_PRIVATE_KEY')?.trim() ?? '' + const subject = readEnv('WEBPUSH_SUBJECT')?.trim() ?? '' + if (!publicKey || !privateKey || !subject) { + throw new PushContractError('webpush_not_configured', 503) + } + try { + const publicBytes = decodeBase64Url(publicKey) + const privateBytes = decodeBase64Url(privateKey) + if ( + publicBytes.byteLength !== 65 + || publicBytes[0] !== 0x04 + || privateBytes.byteLength !== 32 + || !SUBJECT_PATTERN.test(subject) + ) { + throw new Error('invalid web push credentials') + } + return { publicKey, privateKey, subject } + } catch { + throw new PushContractError('webpush_credentials_invalid', 503) + } +} + +export function parseWebPushSubscription(registrationId: string): WebPushSubscription { + if (typeof registrationId !== 'string' || registrationId.length < 20 || registrationId.length > 8192) { + throw new PushContractError('webpush_registration_invalid', 400, true) + } + try { + const parsed = JSON.parse(registrationId) as Record + const endpoint = typeof parsed.endpoint === 'string' ? parsed.endpoint : '' + const keys = parsed.keys + const p256dh = keys && typeof keys === 'object' && typeof (keys as Record).p256dh === 'string' + ? (keys as Record).p256dh + : '' + const auth = keys && typeof keys === 'object' && typeof (keys as Record).auth === 'string' + ? (keys as Record).auth + : '' + const endpointUrl = new URL(endpoint) + const p256dhBytes = decodeBase64Url(p256dh) + const authBytes = decodeBase64Url(auth) + if ( + endpointUrl.protocol !== 'https:' + || p256dhBytes.byteLength !== 65 + || p256dhBytes[0] !== 0x04 + || authBytes.byteLength < 16 + || authBytes.byteLength > 32 + ) { + throw new Error('invalid subscription') + } + return { endpoint: endpointUrl.toString(), p256dh, auth } + } catch (error) { + if (error instanceof PushContractError) throw error + throw new PushContractError('webpush_registration_invalid', 400, true) + } +} + +async function buildVapidAuthorization( + config: WebPushConfig, + endpoint: string, + nowSeconds: number, +): Promise { + const publicBytes = decodeBase64Url(config.publicKey) + const jwk: JsonWebKey = { + kty: 'EC', + crv: 'P-256', + d: config.privateKey, + x: encodeBase64Url(publicBytes.slice(1, 33)), + y: encodeBase64Url(publicBytes.slice(33, 65)), + } + const key = await crypto.subtle.importKey( + 'jwk', + jwk, + { name: 'ECDSA', namedCurve: 'P-256' }, + false, + ['sign'], + ) + const header = encodeBase64Url(encoder.encode(JSON.stringify({ typ: 'JWT', alg: 'ES256' }))) + const payload = encodeBase64Url(encoder.encode(JSON.stringify({ + aud: new URL(endpoint).origin, + exp: nowSeconds + 12 * 3600, + sub: config.subject, + }))) + const signingInput = `${header}.${payload}` + const signature = await crypto.subtle.sign( + { name: 'ECDSA', hash: 'SHA-256' }, + key, + encoder.encode(signingInput), + ) + const jwt = `${signingInput}.${encodeBase64Url(new Uint8Array(signature))}` + return `vapid t=${jwt}, k=${config.publicKey}` +} + +export async function encryptWebPushPayload( + subscription: WebPushSubscription, + payload: Uint8Array, + randomBytes: (length: number) => Uint8Array, +): Promise { + const uaPublic = decodeBase64Url(subscription.p256dh) + const authSecret = decodeBase64Url(subscription.auth) + const salt = randomBytes(16) + + const applicationKeys = await crypto.subtle.generateKey( + { name: 'ECDH', namedCurve: 'P-256' }, + true, + ['deriveBits'], + ) + const applicationPublic = new Uint8Array( + await crypto.subtle.exportKey('raw', applicationKeys.publicKey), + ) + const userAgentKey = await crypto.subtle.importKey( + 'raw', + toArrayBuffer(uaPublic), + { name: 'ECDH', namedCurve: 'P-256' }, + false, + [], + ) + const sharedSecret = new Uint8Array( + await crypto.subtle.deriveBits( + { name: 'ECDH', public: userAgentKey }, + applicationKeys.privateKey, + 256, + ), + ) + + const prkKey = await hmacSha256(authSecret, sharedSecret) + const keyInfo = concat( + encoder.encode('WebPush: info\x00'), + uaPublic, + applicationPublic, + ) + const ikm = await hkdfExpand(prkKey, keyInfo, 32) + const prk = await hmacSha256(salt, ikm) + const contentEncryptionKey = await hkdfExpand( + prk, + encoder.encode('Content-Encoding: aes128gcm\x00'), + 16, + ) + const nonce = await hkdfExpand(prk, encoder.encode('Content-Encoding: nonce\x00'), 12) + + const plaintext = concat(payload, Uint8Array.of(0x02)) + const aesKey = await crypto.subtle.importKey( + 'raw', + toArrayBuffer(contentEncryptionKey), + 'AES-GCM', + false, + ['encrypt'], + ) + const ciphertext = new Uint8Array( + await crypto.subtle.encrypt( + { name: 'AES-GCM', iv: toArrayBuffer(nonce), tagLength: 128 }, + aesKey, + toArrayBuffer(plaintext), + ), + ) + + const header = new Uint8Array(16 + 4 + 1 + 65) + header.set(salt, 0) + new DataView(header.buffer).setUint32(16, RECORD_SIZE, false) + header[20] = 65 + header.set(applicationPublic, 21) + + return concat(header, ciphertext) +} + +export async function sendWebPushMessage( + registrationId: string, + notification: PushNotification, + options: WebPushSendOptions = {}, +): Promise { + const subscription = parseWebPushSubscription(registrationId) + const config = options.config ?? readWebPushConfig() + const fetchImpl = options.fetchImpl ?? fetch + const nowSeconds = Math.floor((options.now ?? Date.now()) / 1000) + const randomBytes = options.randomBytes + ?? ((length: number) => crypto.getRandomValues(new Uint8Array(length))) + + const payload = encoder.encode(JSON.stringify({ + title: notification.title, + body: notification.body, + data: notification.data, + })) + const body = await encryptWebPushPayload(subscription, payload, randomBytes) + const authorization = await buildVapidAuthorization(config, subscription.endpoint, nowSeconds) + + let response: Response + try { + response = await fetchImpl(subscription.endpoint, { + method: 'POST', + headers: { + Authorization: authorization, + 'Content-Encoding': 'aes128gcm', + 'Content-Type': 'application/octet-stream', + TTL: '86400', + }, + body, + signal: AbortSignal.timeout(REQUEST_TIMEOUT_MS), + }) + } catch { + throw new PushContractError('webpush_send_timeout', 504) + } + + if (response.status === 200 || response.status === 201 || response.status === 202) { + return { status: response.status } + } + if (response.status === 404 || response.status === 410) { + throw new PushContractError('webpush_registration_stale', 410, true) + } + if (response.status === 413) { + throw new PushContractError('webpush_payload_too_large', 500) + } + throw new PushContractError('webpush_send_failed', 502) +} diff --git a/server/supabase/functions/send-push/index.ts b/server/supabase/functions/send-push/index.ts index 5729290..b2fd0f3 100644 --- a/server/supabase/functions/send-push/index.ts +++ b/server/supabase/functions/send-push/index.ts @@ -8,6 +8,7 @@ import { isServiceRoleAuthorization, isUuid, parsePushRequest, + providerIsSupported, PushContractError, readFcmConfig, sendFcmMessage, @@ -15,6 +16,8 @@ import { type PushProvider, type TeamInviteNotificationContext, } from '../_shared/push-contract.ts' +import { readApnsConfig, sendApnsMessage } from '../_shared/apns.ts' +import { readWebPushConfig, sendWebPushMessage } from '../_shared/webpush.ts' import { createServiceRoleClient } from '../_shared/quota.ts' const JSON_HEADERS = { @@ -435,8 +438,8 @@ async function finalizeDispatch( return normalizeDispatchSummary(data, attemptId) } -function retryableFcmError(error: unknown): { code: string; status: number } { - if (!(error instanceof PushContractError)) return { code: 'fcm_send_failed', status: 502 } +function retryablePushError(error: unknown, fallbackCode: string): { code: string; status: number } { + if (!(error instanceof PushContractError)) return { code: fallbackCode, status: 502 } const permanent = new Set(['invalid_push_payload']) return { code: permanent.has(error.code) ? 'push_payload_invalid' : error.code, @@ -478,58 +481,116 @@ async function processAttempt(input: { inviteContext = { teamId: invite.team_id, inviteToken: invite.token } } - const unsupported = deliveries.filter((delivery) => delivery.provider !== 'fcm') + const unsupported = deliveries.filter((delivery) => !providerIsSupported(delivery.provider)) await mapWithConcurrency(unsupported, SEND_CONCURRENCY, (delivery) => ( finalizeDelivery(delivery, 'permanent_failure', 'push_provider_not_supported') )) - const fcmDeliveries = deliveries.filter((delivery) => delivery.provider === 'fcm') let transientError: string | null = null let transientStatus: number | null = null - if (fcmDeliveries.length > 0) { + + const recordTransient = (code: string, status: number): void => { + transientError ??= code + transientStatus ??= status + } + + const deliveriesFor = (provider: PushProvider): LeasedDelivery[] => ( + deliveries.filter((delivery) => delivery.provider === provider) + ) + + const failProvider = async (provider: PushProvider, code: string): Promise => { + await mapWithConcurrency(deliveriesFor(provider), SEND_CONCURRENCY, (delivery) => ( + finalizeDelivery(delivery, 'retryable_failure', code) + )) + } + + const sendFor = async ( + provider: PushProvider, + fallbackCode: string, + send: (registrationId: string) => Promise, + ): Promise => { + const scoped = deliveriesFor(provider) + if (scoped.length === 0) return + await mapWithConcurrency(scoped, SEND_CONCURRENCY, async (delivery) => { + try { + await send(delivery.registrationId) + await finalizeDelivery(delivery, 'delivered', null) + } catch (error) { + if (error instanceof PushContractError && error.staleRegistration) { + await finalizeDelivery(delivery, 'stale', 'registration_stale') + return + } + const normalized = retryablePushError(error, fallbackCode) + if (normalized.code === 'push_payload_invalid') { + await finalizeDelivery(delivery, 'permanent_failure', normalized.code) + } else { + recordTransient(normalized.code, normalized.status) + await finalizeDelivery(delivery, 'retryable_failure', normalized.code) + } + } + }) + } + + const buildNotification = (): ReturnType => ( + buildPushNotification(input.eventType, input.resourceId, 'en', inviteContext) + ) + + if (deliveriesFor('fcm').length > 0) { let config: ReturnType | null = null let accessToken: string | null = null try { config = readFcmConfig() accessToken = await getFcmAccessToken(config) } catch (error) { - const normalized = retryableFcmError(error) - transientError = normalized.code - transientStatus = normalized.status - await mapWithConcurrency(fcmDeliveries, SEND_CONCURRENCY, (delivery) => ( - finalizeDelivery(delivery, 'retryable_failure', normalized.code) + const normalized = retryablePushError(error, 'fcm_send_failed') + recordTransient(normalized.code, normalized.status) + await failProvider('fcm', normalized.code) + } + if (config && accessToken) { + await sendFor('fcm', 'fcm_send_failed', (registrationId) => sendFcmMessage( + registrationId, + buildNotification(), + { + config: config as ReturnType, + getAccessToken: () => Promise.resolve(accessToken as string), + }, )) } + } - if (config && accessToken) { - await mapWithConcurrency(fcmDeliveries, SEND_CONCURRENCY, async (delivery) => { - const notification = buildPushNotification( - input.eventType, - input.resourceId, - 'en', - inviteContext, - ) - try { - await sendFcmMessage(delivery.registrationId, notification, { - config, - getAccessToken: () => Promise.resolve(accessToken), - }) - await finalizeDelivery(delivery, 'delivered', null) - } catch (error) { - if (error instanceof PushContractError && error.staleRegistration) { - await finalizeDelivery(delivery, 'stale', 'registration_stale') - return - } - const normalized = retryableFcmError(error) - if (normalized.code === 'push_payload_invalid') { - await finalizeDelivery(delivery, 'permanent_failure', normalized.code) - } else { - transientError ??= normalized.code - transientStatus ??= normalized.status - await finalizeDelivery(delivery, 'retryable_failure', normalized.code) - } - } - }) + if (deliveriesFor('webpush').length > 0) { + let webPushConfig: ReturnType | null = null + try { + webPushConfig = readWebPushConfig() + } catch (error) { + const normalized = retryablePushError(error, 'webpush_send_failed') + recordTransient(normalized.code, normalized.status) + await failProvider('webpush', normalized.code) + } + if (webPushConfig) { + await sendFor('webpush', 'webpush_send_failed', (registrationId) => sendWebPushMessage( + registrationId, + buildNotification(), + { config: webPushConfig as ReturnType }, + )) + } + } + + if (deliveriesFor('apns').length > 0) { + let apnsConfig: ReturnType | null = null + try { + apnsConfig = readApnsConfig() + } catch (error) { + const normalized = retryablePushError(error, 'apns_send_failed') + recordTransient(normalized.code, normalized.status) + await failProvider('apns', normalized.code) + } + if (apnsConfig) { + await sendFor('apns', 'apns_send_failed', (registrationId) => sendApnsMessage( + registrationId, + buildNotification(), + { config: apnsConfig as ReturnType }, + )) } } diff --git a/server/supabase/migrations/20260913000033_team_activities.sql b/server/supabase/migrations/20260913000033_team_activities.sql new file mode 100644 index 0000000..ec3381d --- /dev/null +++ b/server/supabase/migrations/20260913000033_team_activities.sql @@ -0,0 +1,148 @@ +-- ============================================================================ +-- Team activity feed +-- ---------------------------------------------------------------------------- +-- Adds a team-scoped activity/comment stream so members can see shared notes +-- and system events (member changes, invites, shared meetings/documents). +-- +-- Follows the locked RPC write model established in +-- 20260821000008_atomic_team_security.sql: team tables have no direct +-- INSERT/UPDATE/DELETE policies, so all writes go through SECURITY DEFINER +-- RPCs. Reads are RLS-filtered to team members via public.user_team_ids(). +-- +-- Also registers the team tables (and this one) with the realtime publication. +-- Previously teams/team_members/team_invites were never added, so the mobile +-- team Realtime subscriptions could not receive postgres_changes events. +-- ============================================================================ + +BEGIN; + +CREATE TABLE IF NOT EXISTS public.team_activities ( + id uuid PRIMARY KEY DEFAULT gen_random_uuid(), + team_id uuid NOT NULL REFERENCES public.teams(id) ON DELETE CASCADE, + actor_id uuid REFERENCES auth.users(id) ON DELETE SET NULL, + kind text NOT NULL CHECK ( + kind IN ( + 'note', + 'member_joined', + 'member_left', + 'invite_created', + 'meeting_shared', + 'document_shared' + ) + ), + body text CHECK (body IS NULL OR char_length(body) <= 2000), + metadata jsonb NOT NULL DEFAULT '{}'::jsonb, + created_at timestamptz NOT NULL DEFAULT now() +); + +CREATE INDEX IF NOT EXISTS idx_team_activities_team_created + ON public.team_activities (team_id, created_at DESC); + +ALTER TABLE public.team_activities ENABLE ROW LEVEL SECURITY; + +DROP POLICY IF EXISTS team_activities_read_member ON public.team_activities; +CREATE POLICY team_activities_read_member ON public.team_activities + FOR SELECT + TO authenticated + USING (team_id IN (SELECT public.user_team_ids(auth.uid()))); + +-- Writes are RPC-only: no INSERT/UPDATE/DELETE policies exist on purpose. +GRANT SELECT ON public.team_activities TO authenticated; +REVOKE INSERT, UPDATE, DELETE ON public.team_activities FROM PUBLIC, anon, authenticated; + +CREATE OR REPLACE FUNCTION public.create_team_activity( + p_team_id uuid, + p_kind text, + p_body text, + p_metadata jsonb DEFAULT '{}'::jsonb +) +RETURNS jsonb +LANGUAGE plpgsql +SECURITY DEFINER +SET search_path = '' +AS $$ +DECLARE + current_user_id uuid := auth.uid(); + member_role text; + normalized_body text; + new_id uuid; +BEGIN + IF current_user_id IS NULL THEN + RAISE EXCEPTION 'authentication_required' USING ERRCODE = '42501'; + END IF; + IF p_kind IS NULL OR p_kind NOT IN ( + 'note', + 'member_joined', + 'member_left', + 'invite_created', + 'meeting_shared', + 'document_shared' + ) THEN + RAISE EXCEPTION 'invalid_activity_kind' USING ERRCODE = '22023'; + END IF; + normalized_body := nullif(btrim(coalesce(p_body, '')), ''); + IF normalized_body IS NOT NULL AND char_length(normalized_body) > 2000 THEN + RAISE EXCEPTION 'activity_body_too_long' USING ERRCODE = '22001'; + END IF; + IF normalized_body IS NULL AND p_kind = 'note' THEN + RAISE EXCEPTION 'activity_body_required' USING ERRCODE = '22023'; + END IF; + PERFORM pg_advisory_xact_lock(hashtextextended(p_team_id::text, 73053)); + SELECT role INTO member_role FROM public.team_members + WHERE team_id = p_team_id AND user_id = current_user_id; + IF member_role IS NULL THEN + RAISE EXCEPTION 'team_member_required' USING ERRCODE = '42501'; + END IF; + INSERT INTO public.team_activities (team_id, actor_id, kind, body, metadata) + VALUES ( + p_team_id, + current_user_id, + p_kind, + normalized_body, + coalesce(p_metadata, '{}'::jsonb) + ) + RETURNING id INTO new_id; + RETURN jsonb_build_object( + 'id', new_id, + 'team_id', p_team_id, + 'actor_id', current_user_id, + 'kind', p_kind, + 'body', normalized_body, + 'created_at', now() + ); +END; +$$; + +REVOKE ALL ON FUNCTION public.create_team_activity(uuid, text, text, jsonb) FROM PUBLIC, anon; +GRANT EXECUTE ON FUNCTION public.create_team_activity(uuid, text, text, jsonb) TO authenticated; + +-- Realtime: register team tables + activities. ALTER PUBLICATION ADD TABLE has +-- no IF NOT EXISTS, so probe pg_publication_tables first. +DO $$ +DECLARE + t text; +BEGIN + FOR t IN + SELECT unnest(ARRAY[ + 'team_activities', + 'teams', + 'team_members', + 'team_invites' + ]) + LOOP + IF NOT EXISTS ( + SELECT 1 + FROM pg_publication_tables + WHERE pubname = 'supabase_realtime' + AND schemaname = 'public' + AND tablename = t + ) THEN + EXECUTE format( + 'ALTER PUBLICATION supabase_realtime ADD TABLE public.%I', + t + ); + END IF; + END LOOP; +END $$; + +COMMIT;