1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
|
// Next.js API route support: https://nextjs.org/docs/api-routes/introduction
import express from "express";
import http from "http";
import {
createPool,
DatabasePoolType,
DatabaseTransactionConnectionType,
sql,
} from "slonik";
import * as config from "./config";
import { ValueRecord } from "../shared-types";
import { Server as SocketIOServer, Socket } from "socket.io";
const app = express();
const httpServer = http.createServer(app);
const io = new SocketIOServer(httpServer, {
cors: {
origin: "*",
},
});
type Connection = DatabasePoolType | DatabaseTransactionConnectionType;
const _connection = createPool(config.postgresConnectionString);
const getCurrentValue = async (conn: Connection): Promise<ValueRecord["value"]> => {
return await conn.oneFirst<ValueRecord["value"]>(sql`
SELECT value
FROM state
`);
};
const incrementCurrentValue = async (conn: Connection): Promise<number> =>
await conn.transaction(async (trans) => {
const current = await getCurrentValue(trans);
const newValue = current + 1;
await trans.query(sql`
UPDATE state
SET value = ${newValue}
`);
return newValue;
});
io.on("connection", async (socket) => {
console.log(`a user connected: ${socket.id}`);
const currentValue = await getCurrentValue(_connection);
socket.emit("update", currentValue);
socket.on("increment", async () => {
const newValue = await incrementCurrentValue(_connection);
io.emit("update", newValue); // broadcast
});
});
httpServer.listen(config.port, () => {
console.log(`Listening on http://localhost:${config.port}`);
});
|