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