You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.

335 lines
9.7 KiB

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

Powered by TurnKey Linux.