Nexo realtime — WebSocket distributor for chat fan-out

Replaces the chat poll-every-3-seconds compromise with a real
WebSocket pipeline. Server side:

`realtime.mjs`:
- In-memory hub mapping `userId → Set<WebSocket>`. 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) <noreply@anthropic.com>
master
dev 5 months ago
parent a713e0312d
commit 14888e63d9

20
package-lock.json generated

@ -7,15 +7,18 @@
"": { "": {
"name": "active", "name": "active",
"version": "0.0.1", "version": "0.0.1",
"license": "UNLICENSED",
"dependencies": { "dependencies": {
"@floating-ui/core": "^1.7.1", "@floating-ui/core": "^1.7.1",
"@floating-ui/dom": "^1.7.1", "@floating-ui/dom": "^1.7.1",
"@sentry/browser": "^10.50.0", "@sentry/browser": "^10.50.0",
"@types/ws": "^8.18.1",
"csstype": "^3.2.3", "csstype": "^3.2.3",
"esm-env": "^1.1.2", "esm-env": "^1.1.2",
"runed": "0.37.1", "runed": "0.37.1",
"tabbable": "^6.2.0", "tabbable": "^6.2.0",
"web-vitals": "^5.2.0" "web-vitals": "^5.2.0",
"ws": "^8.20.0"
}, },
"devDependencies": { "devDependencies": {
"@eslint/compat": "^2.0.4", "@eslint/compat": "^2.0.4",
@ -40,6 +43,9 @@
"vite": "^8.0.7", "vite": "^8.0.7",
"vitest": "^4.1.3", "vitest": "^4.1.3",
"vitest-browser-svelte": "^2.1.0" "vitest-browser-svelte": "^2.1.0"
},
"engines": {
"node": ">=22"
} }
}, },
"node_modules/@asamuzakjp/css-color": { "node_modules/@asamuzakjp/css-color": {
@ -1126,7 +1132,6 @@
"version": "22.19.17", "version": "22.19.17",
"resolved": "https://registry.npmjs.org/@types/node/-/node-22.19.17.tgz", "resolved": "https://registry.npmjs.org/@types/node/-/node-22.19.17.tgz",
"integrity": "sha512-wGdMcf+vPYM6jikpS/qhg6WiqSV/OhG+jeeHT/KlVqxYfD40iYJf9/AE1uQxVWFvU7MipKRkRv8NSHiCGgPr8Q==", "integrity": "sha512-wGdMcf+vPYM6jikpS/qhg6WiqSV/OhG+jeeHT/KlVqxYfD40iYJf9/AE1uQxVWFvU7MipKRkRv8NSHiCGgPr8Q==",
"dev": true,
"license": "MIT", "license": "MIT",
"dependencies": { "dependencies": {
"undici-types": "~6.21.0" "undici-types": "~6.21.0"
@ -1138,6 +1143,15 @@
"integrity": "sha512-ScaPdn1dQczgbl0QFTeTOmVHFULt394XJgOQNoyVhZ6r2vLnMLJfBPd53SB52T/3G36VI1/g2MZaX0cwDuXsfw==", "integrity": "sha512-ScaPdn1dQczgbl0QFTeTOmVHFULt394XJgOQNoyVhZ6r2vLnMLJfBPd53SB52T/3G36VI1/g2MZaX0cwDuXsfw==",
"license": "MIT" "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": { "node_modules/@typescript-eslint/eslint-plugin": {
"version": "8.59.0", "version": "8.59.0",
"resolved": "https://registry.npmjs.org/@typescript-eslint/eslint-plugin/-/eslint-plugin-8.59.0.tgz", "resolved": "https://registry.npmjs.org/@typescript-eslint/eslint-plugin/-/eslint-plugin-8.59.0.tgz",
@ -3687,7 +3701,6 @@
"version": "6.21.0", "version": "6.21.0",
"resolved": "https://registry.npmjs.org/undici-types/-/undici-types-6.21.0.tgz", "resolved": "https://registry.npmjs.org/undici-types/-/undici-types-6.21.0.tgz",
"integrity": "sha512-iwDZqg0QAGrg9Rav5H4n0M64c3mkR59cJ6wQp+7C4nI0gsmExaedaYLNO44eT4AtBBwjbTiGPMlt2Md0T9H9JQ==", "integrity": "sha512-iwDZqg0QAGrg9Rav5H4n0M64c3mkR59cJ6wQp+7C4nI0gsmExaedaYLNO44eT4AtBBwjbTiGPMlt2Md0T9H9JQ==",
"dev": true,
"license": "MIT" "license": "MIT"
}, },
"node_modules/uri-js": { "node_modules/uri-js": {
@ -4028,7 +4041,6 @@
"version": "8.20.0", "version": "8.20.0",
"resolved": "https://registry.npmjs.org/ws/-/ws-8.20.0.tgz", "resolved": "https://registry.npmjs.org/ws/-/ws-8.20.0.tgz",
"integrity": "sha512-sAt8BhgNbzCtgGbt2OxmpuryO63ZoDk/sqaB/znQm94T4fCEsy/yV+7CdC1kJhOU9lboAEU7R3kquuycDoibVA==", "integrity": "sha512-sAt8BhgNbzCtgGbt2OxmpuryO63ZoDk/sqaB/znQm94T4fCEsy/yV+7CdC1kJhOU9lboAEU7R3kquuycDoibVA==",
"dev": true,
"license": "MIT", "license": "MIT",
"engines": { "engines": {
"node": ">=10.0.0" "node": ">=10.0.0"

@ -26,6 +26,7 @@
"prepare": "svelte-kit sync || echo ''", "prepare": "svelte-kit sync || echo ''",
"check": "svelte-kit sync && svelte-check --tsconfig ./tsconfig.json", "check": "svelte-kit sync && svelte-check --tsconfig ./tsconfig.json",
"check:watch": "svelte-kit sync && svelte-check --tsconfig ./tsconfig.json --watch", "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", "dev:conn-chat": "node scripts/conn-chat-server.mjs",
"lint": "prettier --check . && eslint .", "lint": "prettier --check . && eslint .",
"format": "prettier --write .", "format": "prettier --write .",
@ -65,10 +66,12 @@
"@floating-ui/core": "^1.7.1", "@floating-ui/core": "^1.7.1",
"@floating-ui/dom": "^1.7.1", "@floating-ui/dom": "^1.7.1",
"@sentry/browser": "^10.50.0", "@sentry/browser": "^10.50.0",
"@types/ws": "^8.18.1",
"csstype": "^3.2.3", "csstype": "^3.2.3",
"esm-env": "^1.1.2", "esm-env": "^1.1.2",
"runed": "0.37.1", "runed": "0.37.1",
"tabbable": "^6.2.0", "tabbable": "^6.2.0",
"web-vitals": "^5.2.0" "web-vitals": "^5.2.0",
"ws": "^8.20.0"
} }
} }

@ -0,0 +1,66 @@
/**
* Realtime distributor for the dating server. Holds an in-memory map
* of `userId -> Set<WebSocket>` 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<WebSocket>
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;
}

@ -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 });
}
}

@ -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}`);
}
Loading…
Cancel
Save

Powered by TurnKey Linux.