From 2bd0547451dead484346c59983469741c2c8a443 Mon Sep 17 00:00:00 2001 From: Apher Date: Mon, 29 Jun 2026 00:51:59 +0000 Subject: [PATCH] Add streams initializer to redis model --- src/models/index.js | 44 ++++++++++++++++++++++++++++++++------------ src/server.js | 2 ++ 2 files changed, 34 insertions(+), 12 deletions(-) diff --git a/src/models/index.js b/src/models/index.js index 0dd6460..4ded5f4 100644 --- a/src/models/index.js +++ b/src/models/index.js @@ -1,3 +1,4 @@ +const config = require("../config"); const dbConfig = require("../config/db.config"); const redisConfig = require("../config/redis.config"); @@ -51,27 +52,46 @@ redis.client = Redis.createClient({ redis.client.on('error', (err) => console.error('Redis error:', err)); - - -const initializeStreams = async () => { +redis.initializeStreams = async () => { try { await redis.client.connect().catch(console.error); console.log("Initializing Redis Streams..."); // Create transcriptions streams consumer group - await client.xGroupCreate("communication", "mtf-devices", '0', {MKSTREAM: true}); + try { + await client.xGroupCreate("communication", "mtf-devices", '0', {MKSTREAM: true}); - console.log("Created consumer group: mtf-devices (communication)"); - } catch (err) { - if (err.message.includes('BUSYGROUP')) { - console.log("Consumer group already exists: mtf-devices"); - } else { - console.error(err); + console.log("Created consumer group: mtf-devices (communication)"); + } catch (err) { + if (err.message.includes('BUSYGROUP')) { + console.log("Consumer group already exists: mtf-devices"); + } else { + console.error(err); + } } - } - const deviceIds = [""] + const deviceIds = [config.DEVICE_ID]; + + for (const deviceId of deviceIds) { + const groupName = `${deviceId}_group`; + try { + await client.xGroupCreate('communication', groupName, '0', {MKSTREAM: true}); + console.log(`Created consumer group: ${groupName} (communication)`); + } catch (err) { + if (err.message.includes('BUSYGROUP')) { + console.log(`Consumer group already exists: ${groupName}`); + } else { + console.error(err); + } + } + } + + console.log("Streams initialized successfully"); + await client.quit(); + } catch (err) { + console.error("Error initializing streams:", err.message); + } } module.exports = { diff --git a/src/server.js b/src/server.js index 0ab3ddc..6f2bf79 100644 --- a/src/server.js +++ b/src/server.js @@ -9,6 +9,7 @@ const rateLimit = require("express-rate-limit"); const morgan = require("morgan"); const models = require('./models'); const db = models.db; +const redis = models.redis; const {sequelize, Sequelize} = db; const Role = db.role; @@ -16,6 +17,7 @@ const authConfig = require("./config/auth.config"); const config = require("./config"); const createApp = () => { + redis.initializeStreams(); const app = express(); // Security headers