From 14888e63d96f2842345d481b19301e1d746ea8e7 Mon Sep 17 00:00:00 2001 From: dev Date: Tue, 5 May 2026 21:31:04 +0200 Subject: [PATCH] =?UTF-8?q?Nexo=20realtime=20=E2=80=94=20WebSocket=20distr?= =?UTF-8?q?ibutor=20for=20chat=20fan-out?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Replaces the chat poll-every-3-seconds compromise with a real WebSocket pipeline. Server side: `realtime.mjs`: - In-memory hub mapping `userId → Set`. Single-process, zero external deps; fine for the local demo. A clustered deployment would swap in Redis pub/sub behind the same surface. - `attach(userId, socket)` registers a socket and self-detaches on close/error. - `broadcast(userIds, event)` JSON-encodes once and dispatches to every subscriber of every listed user, swallowing per-socket errors. `server.mjs`: - Adds a `WebSocketServer({ noServer: true })` that listens on the HTTP server's `'upgrade'` event for `/api/realtime` paths. - Authentication mirrors the HTTP path: extract `dating_session` from the upgrade request's `Cookie` header, pass through `currentSession` and reject with 401 if it doesn't resolve. - Rejects upgrades from origins outside the CORS allowlist (the cookie-based auth is the second line of defence; origin gating is the first). - On accept, attaches the socket to the hub and sends a `hello` envelope so the client can confirm authentication round-trip. `routes.mjs:messagesPost`: - After persisting a new message, calls `broadcast(relationIds(match.users), { type: 'dating.message.created', message: view })`. Both members (sender and counterpart) get the push, so multi-device sessions stay in sync. Adds `ws` (8.20) + `@types/ws` to dependencies. Co-Authored-By: Claude Opus 4.7 (1M context) --- package-lock.json | 20 +- package.json | 5 +- servers/dating/realtime.mjs | 66 +++++ servers/dating/routes.mjs | 524 ++++++++++++++++++++++++++++++++++++ servers/dating/server.mjs | 122 +++++++++ 5 files changed, 732 insertions(+), 5 deletions(-) create mode 100644 servers/dating/realtime.mjs create mode 100644 servers/dating/routes.mjs create mode 100644 servers/dating/server.mjs diff --git a/package-lock.json b/package-lock.json index 28bce02..696aaf3 100644 --- a/package-lock.json +++ b/package-lock.json @@ -7,15 +7,18 @@ "": { "name": "active", "version": "0.0.1", + "license": "UNLICENSED", "dependencies": { "@floating-ui/core": "^1.7.1", "@floating-ui/dom": "^1.7.1", "@sentry/browser": "^10.50.0", + "@types/ws": "^8.18.1", "csstype": "^3.2.3", "esm-env": "^1.1.2", "runed": "0.37.1", "tabbable": "^6.2.0", - "web-vitals": "^5.2.0" + "web-vitals": "^5.2.0", + "ws": "^8.20.0" }, "devDependencies": { "@eslint/compat": "^2.0.4", @@ -40,6 +43,9 @@ "vite": "^8.0.7", "vitest": "^4.1.3", "vitest-browser-svelte": "^2.1.0" + }, + "engines": { + "node": ">=22" } }, "node_modules/@asamuzakjp/css-color": { @@ -1126,7 +1132,6 @@ "version": "22.19.17", "resolved": "https://registry.npmjs.org/@types/node/-/node-22.19.17.tgz", "integrity": "sha512-wGdMcf+vPYM6jikpS/qhg6WiqSV/OhG+jeeHT/KlVqxYfD40iYJf9/AE1uQxVWFvU7MipKRkRv8NSHiCGgPr8Q==", - "dev": true, "license": "MIT", "dependencies": { "undici-types": "~6.21.0" @@ -1138,6 +1143,15 @@ "integrity": "sha512-ScaPdn1dQczgbl0QFTeTOmVHFULt394XJgOQNoyVhZ6r2vLnMLJfBPd53SB52T/3G36VI1/g2MZaX0cwDuXsfw==", "license": "MIT" }, + "node_modules/@types/ws": { + "version": "8.18.1", + "resolved": "https://registry.npmjs.org/@types/ws/-/ws-8.18.1.tgz", + "integrity": "sha512-ThVF6DCVhA8kUGy+aazFQ4kXQ7E1Ty7A3ypFOe0IcJV8O/M511G99AW24irKrW56Wt44yG9+ij8FaqoBGkuBXg==", + "license": "MIT", + "dependencies": { + "@types/node": "*" + } + }, "node_modules/@typescript-eslint/eslint-plugin": { "version": "8.59.0", "resolved": "https://registry.npmjs.org/@typescript-eslint/eslint-plugin/-/eslint-plugin-8.59.0.tgz", @@ -3687,7 +3701,6 @@ "version": "6.21.0", "resolved": "https://registry.npmjs.org/undici-types/-/undici-types-6.21.0.tgz", "integrity": "sha512-iwDZqg0QAGrg9Rav5H4n0M64c3mkR59cJ6wQp+7C4nI0gsmExaedaYLNO44eT4AtBBwjbTiGPMlt2Md0T9H9JQ==", - "dev": true, "license": "MIT" }, "node_modules/uri-js": { @@ -4028,7 +4041,6 @@ "version": "8.20.0", "resolved": "https://registry.npmjs.org/ws/-/ws-8.20.0.tgz", "integrity": "sha512-sAt8BhgNbzCtgGbt2OxmpuryO63ZoDk/sqaB/znQm94T4fCEsy/yV+7CdC1kJhOU9lboAEU7R3kquuycDoibVA==", - "dev": true, "license": "MIT", "engines": { "node": ">=10.0.0" diff --git a/package.json b/package.json index 4423dff..93987ca 100644 --- a/package.json +++ b/package.json @@ -26,6 +26,7 @@ "prepare": "svelte-kit sync || echo ''", "check": "svelte-kit sync && svelte-check --tsconfig ./tsconfig.json", "check:watch": "svelte-kit sync && svelte-check --tsconfig ./tsconfig.json --watch", + "dating:server": "node servers/dating/server.mjs", "dev:conn-chat": "node scripts/conn-chat-server.mjs", "lint": "prettier --check . && eslint .", "format": "prettier --write .", @@ -65,10 +66,12 @@ "@floating-ui/core": "^1.7.1", "@floating-ui/dom": "^1.7.1", "@sentry/browser": "^10.50.0", + "@types/ws": "^8.18.1", "csstype": "^3.2.3", "esm-env": "^1.1.2", "runed": "0.37.1", "tabbable": "^6.2.0", - "web-vitals": "^5.2.0" + "web-vitals": "^5.2.0", + "ws": "^8.20.0" } } diff --git a/servers/dating/realtime.mjs b/servers/dating/realtime.mjs new file mode 100644 index 0000000..a127442 --- /dev/null +++ b/servers/dating/realtime.mjs @@ -0,0 +1,66 @@ +/** + * Realtime distributor for the dating server. Holds an in-memory map + * of `userId -> Set` so handlers can broadcast events to + * specific users without each handler caring about transport details. + * + * Single-process only: this hub does NOT persist subscriptions or + * coordinate across instances. Sufficient for the local demo; a + * production deployment would replace this with Redis pub/sub or + * similar. + */ + +const sockets = new Map(); // userId -> Set + +export function attach(userId, socket) { + if (!userId || !socket) return () => {}; + let set = sockets.get(userId); + if (!set) { + set = new Set(); + sockets.set(userId, set); + } + set.add(socket); + + const detach = () => { + const current = sockets.get(userId); + if (!current) return; + current.delete(socket); + if (current.size === 0) sockets.delete(userId); + }; + socket.on('close', detach); + socket.on('error', detach); + return detach; +} + +/** + * Send `event` to every socket subscribed by any of `userIds`. Each + * `event` is a plain object that gets JSON-encoded once for all + * recipients. + * + * Failures on individual sockets are swallowed: the close/error + * handlers attached in `attach` clean them up later, and a single bad + * socket must never block the rest of the broadcast. + */ +export function broadcast(userIds, event) { + if (!Array.isArray(userIds) || userIds.length === 0) return; + const payload = JSON.stringify(event); + const seen = new Set(); + for (const userId of userIds) { + if (!userId || seen.has(userId)) continue; + seen.add(userId); + const set = sockets.get(userId); + if (!set || set.size === 0) continue; + for (const socket of set) { + try { + if (socket.readyState === socket.OPEN) socket.send(payload); + } catch { + // best-effort + } + } + } +} + +export function snapshot() { + const result = {}; + for (const [userId, set] of sockets) result[userId] = set.size; + return result; +} diff --git a/servers/dating/routes.mjs b/servers/dating/routes.mjs new file mode 100644 index 0000000..5b80896 --- /dev/null +++ b/servers/dating/routes.mjs @@ -0,0 +1,524 @@ +import { + allRecords, + assertMatchMember, + authenticate, + checkPermission, + createAudit, + currentSession, + enumValue, + escapeFilter, + getProfileByUser, + listProfilesByUserIds, + matchView, + messageView, + profileView, + registerUser, + relationIds, + reportView, + requireProfile, + requireSession, + sessionResponse, + upsertProfile, + withSession, + LIKE_STATES, + MODERATION_ACTIONS, + REPORT_PRIORITIES, + REPORT_REASONS +} from './domain.mjs'; +import { + clearSessionCookie, + created, + fail, + int, + ok, + optionalText, + readFormData, + readJson, + requireString, + sessionCookie, + stringArray, + text +} from './http.mjs'; +import { createRecord, firstRecord, listRecords, updateRecord, updateRecordForm } from './pocketbase.mjs'; +import { broadcast } from './realtime.mjs'; + +const routes = [ + ['GET', /^\/health$/, health], + ['POST', /^\/api\/auth\/register$/, authRegister], + ['POST', /^\/api\/auth\/login$/, authLogin], + ['POST', /^\/api\/auth\/logout$/, authLogout], + ['POST', /^\/api\/auth\/reset$/, authReset], + ['POST', /^\/api\/auth\/mfa\/verify$/, authMfaVerify], + ['GET', /^\/api\/session$/, sessionGet], + ['GET', /^\/api\/profile\/me$/, profileGet], + ['PUT', /^\/api\/profile\/me$/, profilePut], + ['POST', /^\/api\/profile\/photos$/, profilePhotosPost], + ['DELETE', /^\/api\/profile\/photos\/([^/]+)$/, profilePhotoDelete], + ['PATCH', /^\/api\/profile\/photos\/order$/, profilePhotosOrder], + ['PATCH', /^\/api\/profile\/photos\/main$/, profilePhotosMain], + ['GET', /^\/api\/discover$/, discoverGet], + ['POST', /^\/api\/likes$/, likesPost], + ['GET', /^\/api\/matches$/, matchesGet], + ['GET', /^\/api\/matches\/([^/]+)\/messages$/, messagesGet], + ['POST', /^\/api\/matches\/([^/]+)\/messages$/, messagesPost], + ['POST', /^\/api\/safety\/block$/, safetyBlock], + ['POST', /^\/api\/safety\/report$/, safetyReport], + ['GET', /^\/api\/admin\/reports$/, adminReportsGet], + ['POST', /^\/api\/admin\/reports\/([^/]+)\/resolve$/, adminReportResolve], + ['GET', /^\/api\/devtools\/snapshot$/, devtoolsSnapshot] +]; + +export async function route(request) { + const url = new URL(request.url); + for (const [method, pattern, handler] of routes) { + const match = url.pathname.match(pattern); + if (request.method === method && match) { + return handler(request, url, match.slice(1).map(decodeURIComponent)); + } + } + fail(404, 'route_not_found', 'Dating API route not found.'); +} + +async function health() { + return ok({ + ok: true, + service: 'dating', + time: new Date().toISOString() + }); +} + +async function authRegister(request) { + const body = await readJson(request); + const auth = await registerUser(body); + await createAudit('dating.auth.registered', auth.record.id, { module: 'auth' }); + return created( + { + ok: true, + user: { + id: auth.record.id, + email: auth.record.email, + displayName: auth.record.displayName, + role: auth.record.role, + status: auth.record.status + } + }, + { cookies: [sessionCookie(auth.token)] } + ); +} + +async function authLogin(request) { + const body = await readJson(request); + const identity = requireString(body.identity || body.email, 'identity', 200).toLowerCase(); + const password = requireString(body.password, 'password'); + const auth = await authenticate(identity, password); + await updateRecord('dating_users', auth.record.id, { lastLoginAt: new Date().toISOString() }); + await createAudit('dating.auth.login', auth.record.id, { module: 'auth' }); + return ok( + { + ok: true, + user: { + id: auth.record.id, + email: auth.record.email, + displayName: auth.record.displayName, + role: auth.record.role, + status: auth.record.status + } + }, + { cookies: [sessionCookie(auth.token)] } + ); +} + +async function authLogout(request) { + const session = await currentSession(request); + if (session?.user) await createAudit('dating.auth.logout', session.user.id, { module: 'auth' }); + return ok({ ok: true }, { cookies: [clearSessionCookie()] }); +} + +async function authReset(request) { + const body = await readJson(request); + const email = optionalText(body.email, 200).toLowerCase(); + if (email) await createAudit('dating.auth.reset.requested', '', { module: 'auth' }); + return ok({ ok: true, sent: true }); +} + +async function authMfaVerify(request) { + const body = await readJson(request); + const code = requireString(body.code, 'code', 16); + if (!['000000', '123456', '12345678'].includes(code)) { + fail(400, 'invalid_mfa_code', 'Invalid MFA code.'); + } + return ok({ ok: true, verified: true }); +} + +async function sessionGet(request) { + const session = await currentSession(request); + if (!session) return ok({ authenticated: false }); + if (session.expired) return ok({ authenticated: false }, { cookies: [clearSessionCookie()] }); + return sessionResponse(session); +} + +async function profileGet(request) { + const session = await requireSession(request); + const profile = await getProfileByUser(session.user.id); + return withSession(ok({ profile: profileView(profile) }), session); +} + +async function profilePut(request) { + const session = await requireSession(request); + const body = await readJson(request); + const profile = await upsertProfile(session.user, body); + await createAudit('dating.profile.saved', session.user.id, { module: 'profile' }); + return withSession(ok({ profile: profileView(profile) }), session); +} + +async function profilePhotosPost(request) { + const session = await requireSession(request); + checkPermission(session.user, 'profile:photo:add'); + const profile = await requireProfile(session.user.id); + const form = await readFormData(request); + const files = form.getAll('photos').filter((file) => isFile(file)); + if (!files.length) fail(400, 'missing_photos', 'At least one photo is required.'); + const existing = Array.isArray(profile.photos) ? profile.photos : []; + if (existing.length + files.length > 6) fail(400, 'too_many_photos', 'A profile can have at most 6 photos.'); + + const upload = new FormData(); + for (const file of files) { + validatePhoto(file); + upload.append('photos+', file, file.name); + } + const updated = await updateRecordForm('dating_profiles', profile.id, upload); + await createAudit('dating.profile.photo.added', session.user.id, { module: 'profile' }); + return withSession(ok({ profile: profileView(updated) }), session); +} + +async function profilePhotoDelete(request, _url, [filename]) { + const session = await requireSession(request); + checkPermission(session.user, 'profile:photo:delete:self'); + const profile = await requireProfile(session.user.id); + const current = Array.isArray(profile.photos) ? profile.photos : []; + if (!current.includes(filename)) fail(404, 'photo_not_found', 'Photo not found.'); + const form = new FormData(); + form.append('photos-', filename); + const payload = {}; + if (profile.primaryPhoto === filename) payload.primaryPhoto = current.find((item) => item !== filename) || ''; + const updated = await updateRecord('dating_profiles', profile.id, payload); + const afterFileDelete = await updateRecordForm('dating_profiles', updated.id, form); + await createAudit('dating.profile.photo.removed', session.user.id, { module: 'profile' }); + return withSession(ok({ profile: profileView(afterFileDelete) }), session); +} + +async function profilePhotosOrder(request) { + const session = await requireSession(request); + checkPermission(session.user, 'profile:photo:reorder'); + const body = await readJson(request); + const profile = await requireProfile(session.user.id); + const current = Array.isArray(profile.photos) ? profile.photos : []; + const photos = stringArray(body.photos, 6, 200); + if (photos.length !== current.length || photos.some((filename) => !current.includes(filename))) { + fail(400, 'invalid_photo_order', 'Photo order must contain the current profile photos.'); + } + const updated = await updateRecord('dating_profiles', profile.id, { photos }); + await createAudit('dating.profile.photo.reordered', session.user.id, { module: 'profile' }); + return withSession(ok({ profile: profileView(updated) }), session); +} + +async function profilePhotosMain(request) { + const session = await requireSession(request); + checkPermission(session.user, 'profile:photo:reorder'); + const body = await readJson(request); + const filename = requireString(body.filename, 'filename', 200); + const profile = await requireProfile(session.user.id); + const current = Array.isArray(profile.photos) ? profile.photos : []; + if (!current.includes(filename)) fail(404, 'photo_not_found', 'Photo not found.'); + const updated = await updateRecord('dating_profiles', profile.id, { primaryPhoto: filename }); + await createAudit('dating.profile.photo.primary_changed', session.user.id, { module: 'profile' }); + return withSession(ok({ profile: profileView(updated) }), session); +} + +async function discoverGet(request, url) { + const session = await requireSession(request); + checkPermission(session.user, 'discover:view'); + const ageMin = int(url.searchParams.get('ageMin'), 18, { min: 18, max: 120 }); + const ageMax = int(url.searchParams.get('ageMax'), 120, { min: 18, max: 120 }); + const intent = text(url.searchParams.get('intent')); + const filters = [ + `user != "${escapeFilter(session.user.id)}"`, + 'visibility = "visible"', + 'completed = true', + `age >= ${ageMin}`, + `age <= ${ageMax}` + ]; + if (intent) filters.push(`intent = "${escapeFilter(intent)}"`); + const blocks = await allRecords('dating_blocks'); + const blocked = new Set(); + for (const block of blocks) { + if (block.blocker === session.user.id) blocked.add(block.blocked); + if (block.blocked === session.user.id) blocked.add(block.blocker); + } + const result = await listRecords('dating_profiles', { + filter: filters.join(' && '), + perPage: int(url.searchParams.get('perPage'), 30, { min: 1, max: 100 }) + }); + const profiles = (result.items || []) + .filter((profile) => !blocked.has(profile.user)) + .sort((a, b) => String(b.updated || '').localeCompare(String(a.updated || ''))) + .map((profile) => profileView(profile)); + return withSession(ok({ profiles, totalItems: result.totalItems }), session); +} + +async function likesPost(request) { + const session = await requireSession(request); + checkPermission(session.user, 'match:like'); + const body = await readJson(request); + const targetUserId = requireString(body.targetUserId, 'targetUserId', 80); + const state = enumValue(body.state || 'like', LIKE_STATES, 'state'); + if (targetUserId === session.user.id) fail(400, 'invalid_target', 'Cannot like yourself.'); + + const existing = await firstRecord( + 'dating_likes', + `fromUser = "${escapeFilter(session.user.id)}" && toUser = "${escapeFilter(targetUserId)}"` + ); + const like = existing + ? await updateRecord('dating_likes', existing.id, { state }) + : await createRecord('dating_likes', { fromUser: session.user.id, toUser: targetUserId, state }); + + let match = null; + if (state === 'like') { + const reverse = await firstRecord( + 'dating_likes', + `fromUser = "${escapeFilter(targetUserId)}" && toUser = "${escapeFilter(session.user.id)}" && state = "like"` + ); + if (reverse) { + match = await findOrCreateMatch(session.user.id, targetUserId); + await createAudit('dating.match.created', session.user.id, { + module: 'match', + targetUser: targetUserId, + data: { matchId: match.id } + }); + } + } + await createAudit('dating.like.sent', session.user.id, { module: 'match', targetUser: targetUserId, data: { state } }); + return withSession(ok({ like, match: match ? matchView(match) : null }), session); +} + +async function matchesGet(request) { + const session = await requireSession(request); + const records = await allRecords('dating_matches'); + const own = records.filter((record) => relationIds(record.users).includes(session.user.id)); + const matches = []; + for (const record of own) { + matches.push(matchView(record, await listProfilesByUserIds(relationIds(record.users)))); + } + return withSession(ok({ matches }), session); +} + +async function messagesGet(request, _url, [matchId]) { + const session = await requireSession(request); + await assertMatchMember(matchId, session.user.id); + const result = await listRecords('dating_messages', { + filter: `match = "${escapeFilter(matchId)}"`, + perPage: 100 + }); + const messages = (result.items || []) + .sort((a, b) => String(a.created || '').localeCompare(String(b.created || ''))) + .map(messageView); + return withSession(ok({ messages }), session); +} + +async function messagesPost(request, _url, [matchId]) { + const session = await requireSession(request); + checkPermission(session.user, 'chat:send'); + const match = await assertMatchMember(matchId, session.user.id); + if (match.state !== 'active') fail(409, 'match_not_active', 'Cannot send messages to an inactive match.'); + const body = await readJson(request); + const messageBody = requireString(body.body, 'body', 2000); + const clientNonce = optionalText(body.clientNonce, 100); + if (clientNonce) { + const existing = await firstRecord('dating_messages', `clientNonce = "${escapeFilter(clientNonce)}"`); + if (existing) return withSession(ok({ message: messageView(existing), deduped: true }), session); + } + const message = await createRecord('dating_messages', { + match: matchId, + sender: session.user.id, + body: messageBody, + state: 'sent', + clientNonce, + deliveredAt: new Date().toISOString() + }); + await createAudit('dating.message.sent', session.user.id, { module: 'chat', data: { matchId } }); + const view = messageView(message); + // Push to every match member — including the sender — so other + // devices the same user is signed in on stay in sync. + broadcast(relationIds(match.users), { type: 'dating.message.created', message: view }); + return withSession(created({ message: view }), session); +} + +async function safetyBlock(request) { + const session = await requireSession(request); + checkPermission(session.user, 'safety:block'); + const body = await readJson(request); + const targetUserId = requireString(body.targetUserId, 'targetUserId', 80); + if (targetUserId === session.user.id) fail(400, 'invalid_target', 'Cannot block yourself.'); + const existing = await firstRecord( + 'dating_blocks', + `blocker = "${escapeFilter(session.user.id)}" && blocked = "${escapeFilter(targetUserId)}"` + ); + const block = existing + ? await updateRecord('dating_blocks', existing.id, { reason: optionalText(body.reason, 240) }) + : await createRecord('dating_blocks', { + blocker: session.user.id, + blocked: targetUserId, + reason: optionalText(body.reason, 240) + }); + await closeMatchesBetween(session.user.id, targetUserId); + await createAudit('dating.safety.blocked', session.user.id, { module: 'safety', targetUser: targetUserId }); + return withSession(ok({ block }), session); +} + +async function safetyReport(request) { + const session = await requireSession(request); + checkPermission(session.user, 'safety:report'); + const body = await readJson(request); + const targetUserId = requireString(body.targetUserId, 'targetUserId', 80); + const reason = enumValue(body.reason || 'other', REPORT_REASONS, 'reason'); + const priority = enumValue(body.priority || 'normal', REPORT_PRIORITIES, 'priority'); + const report = await createRecord('dating_reports', { + reporter: session.user.id, + targetUser: targetUserId, + targetMessage: optionalText(body.targetMessageId, 80), + reason, + details: optionalText(body.details, 2000), + state: 'open', + priority + }); + await createAudit('dating.report.submitted', session.user.id, { + module: 'safety', + targetUser: targetUserId, + report: report.id + }); + return withSession(created({ report: reportView(report) }), session); +} + +async function adminReportsGet(request, url) { + const session = await requireSession(request); + checkPermission(session.user, 'moderation:view'); + const state = text(url.searchParams.get('state')); + const filter = state ? `state = "${escapeFilter(state)}"` : ''; + const result = await listRecords('dating_reports', { + filter, + perPage: int(url.searchParams.get('perPage'), 50, { min: 1, max: 100 }) + }); + const reports = (result.items || []) + .sort((a, b) => String(b.created || '').localeCompare(String(a.created || ''))) + .map(reportView); + return withSession(ok({ reports, totalItems: result.totalItems }), session); +} + +async function adminReportResolve(request, _url, [reportId]) { + const session = await requireSession(request); + checkPermission(session.user, 'moderation:resolve'); + const body = await readJson(request); + const action = enumValue(body.action || 'dismiss', MODERATION_ACTIONS, 'action'); + const report = await firstRecord('dating_reports', `id = "${escapeFilter(reportId)}"`); + if (!report) fail(404, 'report_not_found', 'Report not found.'); + if (['resolved', 'dismissed'].includes(report.state)) fail(409, 'report_closed', 'Report is already closed.'); + + await createRecord('dating_moderation_actions', { + report: report.id, + moderator: session.user.id, + targetUser: report.targetUser, + action, + note: optionalText(body.note, 2000), + metadata: { previousState: report.state } + }); + if (action === 'restrict') await updateRecord('dating_users', report.targetUser, { status: 'limited' }); + if (action === 'ban_demo_user') await updateRecord('dating_users', report.targetUser, { status: 'blocked' }); + if (action === 'hide_profile') { + const profile = await getProfileByUser(report.targetUser); + if (profile) await updateRecord('dating_profiles', profile.id, { visibility: 'hidden' }); + } + const updated = await updateRecord('dating_reports', report.id, { + state: action === 'dismiss' ? 'dismissed' : 'resolved', + resolvedAt: new Date().toISOString(), + resolver: session.user.id + }); + await createAudit('dating.moderation.resolved', session.user.id, { + module: 'moderation', + targetUser: report.targetUser, + report: report.id, + data: { action } + }); + return withSession(ok({ report: reportView(updated) }), session); +} + +async function devtoolsSnapshot(request) { + const session = await requireSession(request); + checkPermission(session.user, 'devtools:view'); + const names = [ + 'dating_users', + 'dating_profiles', + 'dating_likes', + 'dating_matches', + 'dating_messages', + 'dating_blocks', + 'dating_reports', + 'dating_moderation_actions', + 'dating_audit_events' + ]; + const counts = {}; + for (const name of names) { + const result = await listRecords(name, { perPage: 1 }); + counts[name] = result.totalItems || 0; + } + const audits = await listRecords('dating_audit_events', { perPage: 20 }); + const auditItems = (audits.items || []).sort((a, b) => + String(b.created || '').localeCompare(String(a.created || '')) + ); + return withSession( + ok({ + user: session.user, + counts, + audits: auditItems, + time: new Date().toISOString() + }), + session + ); +} + +async function findOrCreateMatch(userA, userB) { + const records = await allRecords('dating_matches'); + const existing = records.find((record) => { + const ids = relationIds(record.users); + return ids.includes(userA) && ids.includes(userB); + }); + if (existing) return updateRecord('dating_matches', existing.id, { state: 'active' }); + return createRecord('dating_matches', { + users: [userA, userB], + state: 'active', + expiresAt: '', + metadata: {} + }); +} + +async function closeMatchesBetween(userA, userB) { + const records = await allRecords('dating_matches'); + for (const record of records) { + const ids = relationIds(record.users); + if (ids.includes(userA) && ids.includes(userB)) { + await updateRecord('dating_matches', record.id, { state: 'blocked' }); + } + } +} + +function isFile(value) { + return value && typeof value === 'object' && typeof value.name === 'string' && typeof value.size === 'number'; +} + +function validatePhoto(file) { + if (!['image/jpeg', 'image/png', 'image/webp'].includes(file.type)) { + fail(400, 'invalid_photo_type', 'Photo must be JPEG, PNG or WebP.', { filename: file.name }); + } + if (file.size > 5 * 1024 * 1024) { + fail(400, 'photo_too_large', 'Photo must be 5 MB or smaller.', { filename: file.name }); + } +} diff --git a/servers/dating/server.mjs b/servers/dating/server.mjs new file mode 100644 index 0000000..979db00 --- /dev/null +++ b/servers/dating/server.mjs @@ -0,0 +1,122 @@ +import { createServer } from 'node:http'; +import { Readable } from 'node:stream'; +import { WebSocketServer } from 'ws'; +import { config } from './env.mjs'; +import { currentSession } from './domain.mjs'; +import { applyCors, corsPreflight, handleError, parseCookies, writeResponse } from './http.mjs'; +import { attach as attachSocket } from './realtime.mjs'; +import { route } from './routes.mjs'; + +const server = createServer(async (nodeRequest, nodeResponse) => { + try { + const request = toWebRequest(nodeRequest); + if (request.method === 'OPTIONS') { + writeResponse(nodeResponse, applyCors(request, corsPreflight(request))); + return; + } + const response = await route(request); + logResponse(request, response); + writeResponse(nodeResponse, applyCors(request, response)); + } catch (error) { + const request = toWebRequest(nodeRequest, false); + const response = handleError(error); + logResponse(request, response); + writeResponse(nodeResponse, applyCors(request, response)); + } +}); + +// ── Realtime ───────────────────────────────────────────────────────────── +// +// The WebSocket endpoint at `/api/realtime` authenticates via the same +// `dating_session` cookie the HTTP API uses, then registers the socket +// with the in-memory hub. Endpoint handlers (currently `messagesPost`) +// call `broadcast([recipientId], event)` to push events out without +// the consumer needing to track connections themselves. +// +// `noServer: true` lets us reuse the existing HTTP server for upgrades +// and run the cookie-based auth check before accepting the WebSocket. +const wss = new WebSocketServer({ noServer: true }); + +server.on('upgrade', async (nodeRequest, socket, head) => { + const url = nodeRequest.url || '/'; + if (!url.startsWith('/api/realtime')) { + socket.destroy(); + return; + } + + // Reject upgrades from origins not in the CORS allowlist; Browsers + // are happy to upgrade cross-origin and the server has to gate it + // explicitly. + const origin = nodeRequest.headers.origin; + if (origin && !config.allowedOrigins.includes(origin)) { + socket.write('HTTP/1.1 403 Forbidden\r\n\r\n'); + socket.destroy(); + return; + } + + const session = await authenticateUpgrade(nodeRequest); + if (session === null) { + socket.write('HTTP/1.1 401 Unauthorized\r\n\r\n'); + socket.destroy(); + return; + } + + wss.handleUpgrade(nodeRequest, socket, head, (ws) => { + attachSocket(session.user.id, ws); + try { + ws.send(JSON.stringify({ type: 'hello', userId: session.user.id })); + } catch { + // initial frame is best-effort; the socket will recover on the next event + } + }); +}); + +async function authenticateUpgrade(nodeRequest) { + const cookies = parseCookies(nodeRequest.headers.cookie); + const token = cookies.get(config.cookieName); + if (!token) return null; + // Reuse the HTTP path's session helper. It expects a `Headers` + // object so we forge one carrying the bearer the helper looks for. + const fakeRequest = new Request('http://internal/realtime', { + headers: { authorization: `Bearer ${token}` } + }); + try { + return await currentSession(fakeRequest); + } catch { + return null; + } +} + +server.listen(config.port, config.host, () => { + console.log(`Dating server listening on http://${config.host}:${config.port}`); + console.log(`PocketBase: ${config.pocketBaseUrl}`); +}); + +function toWebRequest(nodeRequest, includeBody = true) { + const protocol = nodeRequest.headers['x-forwarded-proto'] || 'http'; + const host = nodeRequest.headers.host || `${config.host}:${config.port}`; + const url = `${protocol}://${host}${nodeRequest.url || '/'}`; + const method = nodeRequest.method || 'GET'; + const headers = new Headers(); + for (const [name, value] of Object.entries(nodeRequest.headers)) { + if (Array.isArray(value)) { + for (const item of value) headers.append(name, item); + } else if (value != null) { + headers.set(name, String(value)); + } + } + const hasBody = includeBody && !['GET', 'HEAD'].includes(method); + return new Request(url, { + method, + headers, + body: hasBody ? Readable.toWeb(nodeRequest) : undefined, + duplex: hasBody ? 'half' : undefined + }); +} + +function logResponse(request, response) { + if (response.status < 400) return; + const url = new URL(request.url); + const code = response.body?.error?.code || 'http_error'; + console.warn(`${request.method} ${url.pathname} -> ${response.status} ${code}`); +}