room occupancy impl & race cond fix & local user roomID and psuedoID storage for cleanup and indexing
This commit is contained in:
@@ -1,7 +1,8 @@
|
|||||||
import { getRedis } from "./redis.js";
|
import { getRedis } from "../redis/redis.js";
|
||||||
import { prisma } from "../lib/prisma.js";
|
import { prisma } from "../lib/prisma.js";
|
||||||
import { RoomTransition, transitions } from "./transition.js";
|
import { RoomTransition, transitions } from "./transition.js";
|
||||||
import config from "../config/config.js";
|
import config from "../config/config.js";
|
||||||
|
import { type RedisUserData } from "../api/livedata/livedata.repository.js";
|
||||||
|
|
||||||
const adjectives = [
|
const adjectives = [
|
||||||
"Crazy",
|
"Crazy",
|
||||||
@@ -23,6 +24,31 @@ const nouns = [
|
|||||||
"Phoenix"
|
"Phoenix"
|
||||||
];
|
];
|
||||||
|
|
||||||
|
const running = new Set<string>();
|
||||||
|
const userRooms = new Map<string, string>();
|
||||||
|
const userPseudo = new Map<string, string>();
|
||||||
|
|
||||||
|
export async function updateRoomOccupancyListener(message: string, channel: string) {
|
||||||
|
console.log(`[Redis] Key ${message} expired.`);
|
||||||
|
|
||||||
|
const id = message.startsWith("user:") ? message.slice(5) : message;
|
||||||
|
const lastRoomID = userRooms.get(id);
|
||||||
|
const userPseudoID = userPseudo.get(id);
|
||||||
|
|
||||||
|
const redis = getRedis();
|
||||||
|
|
||||||
|
if (!lastRoomID) {
|
||||||
|
console.log("[!] Key not found in local mapι")
|
||||||
|
} else {
|
||||||
|
await redis.decr(`room:${lastRoomID}`);
|
||||||
|
}
|
||||||
|
|
||||||
|
await redis.del(`_user:${userPseudoID}`);
|
||||||
|
|
||||||
|
userRooms.delete(id);
|
||||||
|
userPseudo.delete(id);
|
||||||
|
}
|
||||||
|
|
||||||
function generateNickname() {
|
function generateNickname() {
|
||||||
const adjective =
|
const adjective =
|
||||||
adjectives[Math.floor(Math.random() * adjectives.length)];
|
adjectives[Math.floor(Math.random() * adjectives.length)];
|
||||||
@@ -46,6 +72,14 @@ export async function analyzeData(roomID: string, metrics: {
|
|||||||
tagID: string,
|
tagID: string,
|
||||||
rssi: number
|
rssi: number
|
||||||
}) {
|
}) {
|
||||||
|
if (running.has(metrics.tagID)) {
|
||||||
|
console.log(`[!] Skipping TAG_ID: ${metrics.tagID} already in process.`);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
running.add(metrics.tagID);
|
||||||
|
|
||||||
|
try {
|
||||||
console.log(`\n=======================================================\n` +
|
console.log(`\n=======================================================\n` +
|
||||||
`ROOM_ID: ${roomID}\nTAG_ID: ${metrics.tagID}\nRSSI: ${metrics.rssi}\n` +
|
`ROOM_ID: ${roomID}\nTAG_ID: ${metrics.tagID}\nRSSI: ${metrics.rssi}\n` +
|
||||||
`=======================================================\n`);
|
`=======================================================\n`);
|
||||||
@@ -80,12 +114,19 @@ export async function analyzeData(roomID: string, metrics: {
|
|||||||
timestamp: Date.now()
|
timestamp: Date.now()
|
||||||
}), { PX: config.core.userTTL });
|
}), { PX: config.core.userTTL });
|
||||||
|
|
||||||
|
await redis.set(`_user:${user.pseudoID}`, userID);
|
||||||
|
|
||||||
|
await redis.incr(`room:${roomID}`);
|
||||||
|
|
||||||
|
userRooms.set(userID, roomID);
|
||||||
|
userPseudo.set(userID, user.pseudoID);
|
||||||
|
|
||||||
console.log("[Redis] Stored new user in redis.");
|
console.log("[Redis] Stored new user in redis.");
|
||||||
|
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
let resObj = JSON.parse(res);
|
let resObj: RedisUserData = JSON.parse(res);
|
||||||
console.log(`[Redis] User exists in redis in ROOM_ID: ${resObj["room"]}`);
|
console.log(`[Redis] User exists in redis in ROOM_ID: ${resObj["room"]}`);
|
||||||
|
|
||||||
if (roomID == resObj["room"]) {
|
if (roomID == resObj["room"]) {
|
||||||
@@ -100,6 +141,8 @@ export async function analyzeData(roomID: string, metrics: {
|
|||||||
|
|
||||||
let transition = transitions.get(userID);
|
let transition = transitions.get(userID);
|
||||||
|
|
||||||
|
// update room occupancy
|
||||||
|
|
||||||
if (!transition) {
|
if (!transition) {
|
||||||
console.log("[!] Starting new transition");
|
console.log("[!] Starting new transition");
|
||||||
transition = new RoomTransition();
|
transition = new RoomTransition();
|
||||||
@@ -111,6 +154,11 @@ export async function analyzeData(roomID: string, metrics: {
|
|||||||
|
|
||||||
console.log(`[!] Transition done ${resObj["room"]} -> ${roomID}`);
|
console.log(`[!] Transition done ${resObj["room"]} -> ${roomID}`);
|
||||||
|
|
||||||
|
await redis.decr(`room:${resObj["room"]}`);
|
||||||
|
await redis.incr(`room:${roomID}`);
|
||||||
|
|
||||||
|
userRooms.set(userID, roomID);
|
||||||
|
|
||||||
resObj["rssi"] = metrics.rssi;
|
resObj["rssi"] = metrics.rssi;
|
||||||
resObj["room"] = roomID;
|
resObj["room"] = roomID;
|
||||||
resObj["timestamp"] = Date.now()
|
resObj["timestamp"] = Date.now()
|
||||||
@@ -119,4 +167,7 @@ export async function analyzeData(roomID: string, metrics: {
|
|||||||
await redis.set(`user:${userID}`, JSON.stringify(resObj), { PX: config.core.userTTL });
|
await redis.set(`user:${userID}`, JSON.stringify(resObj), { PX: config.core.userTTL });
|
||||||
} else
|
} else
|
||||||
console.log("[!] Transition declined");
|
console.log("[!] Transition declined");
|
||||||
|
} finally {
|
||||||
|
running.delete(metrics.tagID);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
Reference in New Issue
Block a user