Add an SQL DB abstraction layer
Because what self-respecting Enterprise(R) project doesn't have one of those?
This commit is contained in:
+29
-54
@@ -1,37 +1,34 @@
|
||||
import { PoolClient } from "pg";
|
||||
import sql from "sql-bricks";
|
||||
import sql from "sql-bricks-postgres";
|
||||
|
||||
import { CardRepository, Transaction } from "./repo";
|
||||
import { connect, generateId, generateExtId } from "../db";
|
||||
import { AimeId } from "../model";
|
||||
import { CardRepository, Repositories } from "./repo";
|
||||
import { AimeId, generateExtId } from "../model";
|
||||
import { Transaction, generateId } from "../sql";
|
||||
|
||||
class CardRepositoryImpl implements CardRepository {
|
||||
constructor(private readonly _conn: PoolClient) {}
|
||||
constructor(private readonly _txn: Transaction) {}
|
||||
|
||||
async lookup(luid: string, now: Date): Promise<AimeId | undefined> {
|
||||
const fetchSql = sql
|
||||
.select("c.id", "p.ext_id")
|
||||
.from("aime_card c")
|
||||
.join("aime_player p", { "c.player_id": "p.id" })
|
||||
.where("c.nfc_id", luid)
|
||||
.toParams();
|
||||
.where("c.nfc_id", luid);
|
||||
|
||||
const { rows } = await this._conn.query(fetchSql);
|
||||
const row = await this._txn.fetchRow(fetchSql);
|
||||
|
||||
if (rows.length === 0) {
|
||||
if (row === undefined) {
|
||||
return undefined;
|
||||
}
|
||||
|
||||
const id = rows[0].id;
|
||||
const extId = rows[0].ext_id;
|
||||
const id = row.id;
|
||||
const extId = row.ext_id;
|
||||
|
||||
const touchSql = sql
|
||||
.update("aime_card")
|
||||
.set({ access_time: now })
|
||||
.where("id", id)
|
||||
.toParams();
|
||||
.where("id", id);
|
||||
|
||||
await this._conn.query(touchSql);
|
||||
await this._txn.modify(touchSql);
|
||||
|
||||
return extId;
|
||||
}
|
||||
@@ -41,54 +38,32 @@ class CardRepositoryImpl implements CardRepository {
|
||||
const cardId = generateId();
|
||||
const aimeId = generateExtId() as AimeId;
|
||||
|
||||
const playerSql = sql
|
||||
.insert("aime_player", {
|
||||
id: playerId,
|
||||
ext_id: aimeId,
|
||||
register_time: now,
|
||||
})
|
||||
.toParams();
|
||||
const playerSql = sql.insert("aime_player", {
|
||||
id: playerId,
|
||||
ext_id: aimeId,
|
||||
register_time: now,
|
||||
});
|
||||
|
||||
await this._conn.query(playerSql);
|
||||
await this._txn.modify(playerSql);
|
||||
|
||||
const cardSql = sql
|
||||
.insert("aime_card", {
|
||||
id: cardId,
|
||||
player_id: playerId,
|
||||
nfc_id: luid,
|
||||
register_time: now,
|
||||
access_time: now,
|
||||
})
|
||||
.toParams();
|
||||
const cardSql = sql.insert("aime_card", {
|
||||
id: cardId,
|
||||
player_id: playerId,
|
||||
nfc_id: luid,
|
||||
register_time: now,
|
||||
access_time: now,
|
||||
});
|
||||
|
||||
await this._conn.query(cardSql);
|
||||
await this._txn.modify(cardSql);
|
||||
|
||||
return aimeId;
|
||||
}
|
||||
}
|
||||
|
||||
class TransactionImpl implements Transaction {
|
||||
constructor(private readonly _conn: PoolClient) {}
|
||||
export class SqlRepositories implements Repositories {
|
||||
constructor(private readonly _txn: Transaction) {}
|
||||
|
||||
cards(): CardRepository {
|
||||
return new CardRepositoryImpl(this._conn);
|
||||
}
|
||||
|
||||
async commit(): Promise<void> {
|
||||
await this._conn.query("commit");
|
||||
await this._conn.release();
|
||||
}
|
||||
|
||||
async rollback(): Promise<void> {
|
||||
await this._conn.query("rollback");
|
||||
await this._conn.release();
|
||||
return new CardRepositoryImpl(this._txn);
|
||||
}
|
||||
}
|
||||
|
||||
export async function beginDbSession(): Promise<Transaction> {
|
||||
const conn = await connect();
|
||||
|
||||
await conn.query("begin");
|
||||
|
||||
return new TransactionImpl(conn);
|
||||
}
|
||||
|
||||
@@ -3,6 +3,8 @@ import logger from "debug";
|
||||
import { Repositories } from "./repo";
|
||||
import * as Req from "./request";
|
||||
import * as Res from "./response";
|
||||
import { Transaction } from "../sql";
|
||||
import { SqlRepositories } from "./db";
|
||||
|
||||
const debug = logger("app:aimedb:ops");
|
||||
|
||||
@@ -101,10 +103,12 @@ function log(
|
||||
}
|
||||
|
||||
export async function dispatch(
|
||||
rep: Repositories,
|
||||
txn: Transaction,
|
||||
req: Req.AimeRequest,
|
||||
now: Date
|
||||
): Promise<Res.AimeResponse | undefined> {
|
||||
const rep = new SqlRepositories(txn);
|
||||
|
||||
switch (req.type) {
|
||||
case "hello":
|
||||
return hello(rep, req, now);
|
||||
|
||||
+21
-21
@@ -4,37 +4,37 @@ import { Socket } from "net";
|
||||
import { dispatch } from "./handler";
|
||||
import { AimeRequest } from "./request";
|
||||
import { setup } from "./pipeline";
|
||||
import { beginDbSession } from "./db";
|
||||
import { DataSource } from "../sql/api";
|
||||
|
||||
const debug = logger("app:aimedb:session");
|
||||
|
||||
export default async function aimedb(socket: Socket) {
|
||||
debug("Connection opened");
|
||||
export default function aimedb(db: DataSource) {
|
||||
return async function(socket: Socket) {
|
||||
debug("Connection opened");
|
||||
|
||||
const { input, output } = setup(socket);
|
||||
const txn = await beginDbSession();
|
||||
const { input, output } = setup(socket);
|
||||
|
||||
try {
|
||||
for await (const obj of input) {
|
||||
const now = new Date();
|
||||
const req = obj as AimeRequest;
|
||||
const res = await dispatch(txn, req, now);
|
||||
try {
|
||||
const now = new Date();
|
||||
const req = obj as AimeRequest;
|
||||
const res = await db.transaction(txn => dispatch(txn, req, now));
|
||||
|
||||
if (res === undefined) {
|
||||
debug("Closing connection");
|
||||
if (res === undefined) {
|
||||
debug("Closing connection");
|
||||
|
||||
break;
|
||||
}
|
||||
|
||||
output.write(res);
|
||||
} catch (e) {
|
||||
debug(`Connection error:\n${e.toString()}\n`);
|
||||
|
||||
break;
|
||||
}
|
||||
|
||||
output.write(res);
|
||||
}
|
||||
|
||||
await txn.commit();
|
||||
} catch (e) {
|
||||
debug(`Connection error:\n${e.toString()}\n`);
|
||||
await txn.rollback();
|
||||
}
|
||||
|
||||
debug("Connection closed");
|
||||
socket.end();
|
||||
debug("Connection closed");
|
||||
socket.end();
|
||||
};
|
||||
}
|
||||
|
||||
@@ -9,9 +9,3 @@ export interface CardRepository {
|
||||
export interface Repositories {
|
||||
cards(): CardRepository;
|
||||
}
|
||||
|
||||
export interface Transaction extends Repositories {
|
||||
commit(): Promise<void>;
|
||||
|
||||
rollback(): Promise<void>;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user