From bf809a0309f768dea68ed11798f5e1e515ed92cb Mon Sep 17 00:00:00 2001 From: Apher Date: Fri, 3 Jul 2026 06:09:19 +0000 Subject: [PATCH] Replace initializeStreams with improved ensureGroup method --- src/models/index.js | 47 ++++++++++----------------------------------- src/server.js | 6 ++++-- 2 files changed, 14 insertions(+), 39 deletions(-) diff --git a/src/models/index.js b/src/models/index.js index 65e3ca8..549b563 100644 --- a/src/models/index.js +++ b/src/models/index.js @@ -49,48 +49,21 @@ redis.client = Redis.createClient({ url: `redis://default@${redisConfig.REDIS_HOST}:${redisConfig.REDIS_PORT}` }); -redis.initializeStreams = async () => { - redis.client.on('error', (err) => console.error('Redis error:', err)); +redis.ensureGroup = async (stream, group) => { + await redis.client.connect().catch(console.error); try { - await redis.client.connect().catch(console.error); - - console.log("Initializing Redis Streams..."); + await redis.client.xGroupCreate(stream, group, "0", {MKSTREAM: true}); - // Create transcriptions streams consumer group - try { - await redis.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); - } - } - - const deviceIds = [config.DEVICE_ID]; - - for (const deviceId of deviceIds) { - const groupName = `${deviceId}_group`; - try { - await redis.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 redis.client.quit(); + console.log(`Created consumer group: ${group} (${stream})`); } catch (err) { - console.error("Error initializing streams:", err.message); + if (err.message.includes('BUSYGROUP')) { + console.log(`Consumer group already exists: ${group} (${stream})`); + } else { + console.error(err); + } } + await redis.client.quit(); } module.exports = { diff --git a/src/server.js b/src/server.js index 9890f7d..dad46a9 100644 --- a/src/server.js +++ b/src/server.js @@ -17,8 +17,10 @@ const Role = db.role; const authConfig = require("./config/auth.config"); const config = require("./config"); -const createApp = () => { - redis.initializeStreams(); +const createApp = async () => { + await redis.ensureGroup("speech:commands", `${config.DEVICE_ID}-cg`); + await redis.ensureGroup("speech:events", `${config.DEVICE_ID}-cg`); + const app = express(); // Security headers