import { mkdir, readFile, rename, writeFile } from 'node:fs/promises'; import { dirname, resolve } from 'node:path'; import { Pool } from 'pg'; const empty = () => ({ users: [], sessions: [], rooms: [] }); export class FileStore { constructor(path) { this.path = resolve(path); this.queue = Promise.resolve(); } async init() { await mkdir(dirname(this.path), { recursive: true }); try { this.data = JSON.parse(await readFile(this.path, 'utf8')); } catch (e) { if (e.code !== 'ENOENT') throw e; this.data = empty(); await this.persist(this.data); } } async persist(data) { const temp = `${this.path}.tmp`; await writeFile(temp, JSON.stringify(data), { mode: 0o600 }); await rename(temp, this.path); } async mutate(fn) { const operation = this.queue.then(async () => { const draft = structuredClone(this.data); const result = await fn(draft); await this.persist(draft); this.data = draft; return structuredClone(result); }); this.queue = operation.catch(() => {}); return operation; } async userByEmail(email) { return structuredClone( this.data.users.find((u) => u.email === email) ?? null, ); } async userById(id) { return structuredClone(this.data.users.find((u) => u.id === id) ?? null); } async createUser(user) { return this.mutate((d) => { if ( d.users.some( (u) => (user.email && u.email === user.email) || u.nick.toLowerCase() === user.nick.toLowerCase(), ) ) throw new Error('El correo o nick ya están registrados.'); d.users.push(user); return user; }); } async createSession(session) { return this.mutate((d) => { d.sessions = d.sessions.filter((s) => s.expiresAt > Date.now()); d.sessions.push(session); return session; }); } async updateUser(id, fn) { return this.mutate(async (d) => { const user = d.users.find((u) => u.id === id); if (!user) throw new Error('Cuenta no encontrada.'); await fn(user); if ( d.users.some( (u) => u.id !== id && ((user.email && u.email === user.email) || u.nick.toLowerCase() === user.nick.toLowerCase()), ) ) throw new Error('El correo o nick ya están registrados.'); return user; }); } async deleteUserSessions(id) { return this.mutate((d) => { d.sessions = d.sessions.filter((s) => s.userId !== id); }); } async session(hash) { return structuredClone( this.data.sessions.find( (s) => s.hash === hash && s.expiresAt > Date.now(), ) ?? null, ); } async deleteSession(hash) { return this.mutate((d) => { d.sessions = d.sessions.filter((s) => s.hash !== hash); return null; }); } async listRooms(userId) { return structuredClone( this.data.rooms.filter( (r) => r.visibility === 'public' || r.players.some((p) => p.userId === userId), ), ); } async room(id) { return structuredClone( this.data.rooms.find((r) => r.id === id || r.code === id.toUpperCase()) ?? null, ); } async activeRooms() { return structuredClone( this.data.rooms.filter((r) => r.status !== 'finished'), ); } async createRoom(room) { return this.mutate((d) => { d.rooms.push(room); return room; }); } async updateRoom(id, fn) { return this.mutate(async (d) => { const room = d.rooms.find((r) => r.id === id); if (!room) throw new Error('Sala no encontrada.'); await fn(room); return room; }); } async history(userId) { return structuredClone( this.data.rooms.filter((r) => r.players.some( (p) => p.userId === userId && (r.status === 'finished' || p.withdrawnAt), ), ), ); } async close() { await this.queue; } } export class PostgresStore { constructor(url) { this.pool = new Pool({ connectionString: url, max: 10 }); } async init() { await this.pool.query('SELECT id FROM rooms LIMIT 1'); } async userByEmail(email) { const r = await this.pool.query('SELECT data FROM users WHERE email=$1', [ email, ]); return r.rows[0]?.data ?? null; } async userById(id) { const r = await this.pool.query('SELECT data FROM users WHERE id=$1', [id]); return r.rows[0]?.data ?? null; } async createUser(user) { try { await this.pool.query( 'INSERT INTO users(id,email,nick,data) VALUES($1,$2,$3,$4)', [user.id, user.email, user.nick.toLowerCase(), user], ); return user; } catch (e) { if (e.code === '23505') throw new Error('El correo o nick ya están registrados.'); throw e; } } async createSession(session) { await this.pool.query( 'INSERT INTO sessions(hash,user_id,expires_at) VALUES($1,$2,$3)', [session.hash, session.userId, new Date(session.expiresAt)], ); return session; } async updateUser(id, fn) { const client = await this.pool.connect(); try { await client.query('BEGIN'); const result = await client.query( 'SELECT data FROM users WHERE id=$1 FOR UPDATE', [id], ); const user = result.rows[0]?.data; if (!user) throw new Error('Cuenta no encontrada.'); await fn(user); await client.query( 'UPDATE users SET email=$2,nick=$3,data=$4 WHERE id=$1', [id, user.email, user.nick.toLowerCase(), user], ); await client.query('COMMIT'); return user; } catch (e) { await client.query('ROLLBACK'); if (e.code === '23505') throw new Error('El correo o nick ya están registrados.'); throw e; } finally { client.release(); } } async deleteUserSessions(id) { await this.pool.query('DELETE FROM sessions WHERE user_id=$1', [id]); } async session(hash) { const r = await this.pool.query( 'SELECT hash,user_id AS "userId",expires_at FROM sessions WHERE hash=$1 AND expires_at>now()', [hash], ); return r.rows[0] ? { ...r.rows[0], expiresAt: +r.rows[0].expires_at } : null; } async deleteSession(hash) { await this.pool.query('DELETE FROM sessions WHERE hash=$1', [hash]); } async listRooms(userId) { const r = await this.pool.query( "SELECT data FROM rooms WHERE visibility='public' OR data->'players' @> $1::jsonb ORDER BY created_at DESC", [JSON.stringify([{ userId }])], ); return r.rows.map((r) => r.data); } async activeRooms() { const r = await this.pool.query( "SELECT data FROM rooms WHERE status IN ('waiting','playing')", ); return r.rows.map((r) => r.data); } async room(id) { const r = await this.pool.query( 'SELECT data FROM rooms WHERE id::text=$1 OR code=$2', [id, id.toUpperCase()], ); return r.rows[0]?.data ?? null; } async createRoom(room) { await this.pool.query( 'INSERT INTO rooms(id,code,game_id,game_version,visibility,status,created_at,data) VALUES($1,$2,$3,$4,$5,$6,$7,$8)', [ room.id, room.code, room.gameId, room.gameVersion, room.visibility, room.status, room.createdAt, room, ], ); return room; } async updateRoom(id, fn) { const client = await this.pool.connect(); try { await client.query('BEGIN'); const row = await client.query( 'SELECT data FROM rooms WHERE id=$1 FOR UPDATE', [id], ); if (!row.rows[0]) throw new Error('Sala no encontrada.'); const room = row.rows[0].data; const eventCount = room.events.length, messageCount = room.messages.length; await fn(room); await client.query('UPDATE rooms SET status=$2,data=$3 WHERE id=$1', [ id, room.status, room, ]); for (const event of room.events.slice(eventCount)) await client.query( 'INSERT INTO match_events(room_id,sequence,action_id,actor_id,data) VALUES($1,$2,$3,$4,$5)', [id, event.sequence, event.actionId, event.actorId, event], ); for (const msg of room.messages.slice(messageCount)) await client.query( 'INSERT INTO chat_messages(id,room_id,user_id,data) VALUES($1,$2,$3,$4) ON CONFLICT DO NOTHING', [msg.id, id, msg.userId, msg], ); if (room.result && !room.result.cancelled) for (const player of room.players) await client.query( 'INSERT INTO match_results(room_id,user_id,outcome,placement,reason) VALUES($1,$2,$3,$4,$5) ON CONFLICT(room_id,user_id) DO NOTHING', [ id, player.userId, player.outcome, player.placement, room.result.reason, ], ); await client.query('COMMIT'); return room; } catch (e) { await client.query('ROLLBACK'); throw e; } finally { client.release(); } } async history(userId) { const r = await this.pool.query( "SELECT data FROM rooms WHERE EXISTS (SELECT 1 FROM jsonb_array_elements(data->'players') p WHERE p->>'userId'=$1 AND (status='finished' OR p->>'withdrawnAt' IS NOT NULL)) ORDER BY created_at DESC", [userId], ); return r.rows.map((r) => r.data); } async close() { await this.pool.end(); } } export async function createStore(env = process.env) { if (env.NODE_ENV === 'production' && !env.DATABASE_URL) throw new Error('Producción requiere DATABASE_URL.'); const store = env.DATABASE_URL ? new PostgresStore(env.DATABASE_URL) : new FileStore(env.DATA_FILE || '.data/development.json'); await store.init(); return store; }