From 9b5201332e919a4b793fff8100613618072499d0 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E2=80=9CNaeel=E2=80=9D?= Date: Wed, 22 Jul 2026 09:05:26 +0400 Subject: [PATCH] =?UTF-8?q?fix:=20resilient=20store=20=E2=80=94=20Redis/Mo?= =?UTF-8?q?ngo=20errors=20don't=20stop=20consumer?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- server.js | 45 +++++++++++++++++++++++++-------------------- 1 file changed, 25 insertions(+), 20 deletions(-) diff --git a/server.js b/server.js index 866c506..ef1150f 100644 --- a/server.js +++ b/server.js @@ -41,19 +41,24 @@ async function getMongo() { } async function store(event) { - const r = await getRedis(); - await r.hSet("iot:latest", event.sensor_id, JSON.stringify({ - value: event.value, unit: event.unit, location: event.location, timestamp: event.timestamp, - })); - await r.hIncrBy("iot:counters", event.device_type, 1); - await r.zAdd("iot:recent", [{ score: 0, value: JSON.stringify(event) }]); - await r.zRemRangeByRank("iot:recent", 0, -1001); - const coll = await getMongo(); - await coll.insertOne({ - event_id: event.event_id, sensor_id: event.sensor_id, - device_type: event.device_type, location: event.location, - value: event.value, unit: event.unit, timestamp: new Date(event.timestamp), - }); + try { + const r = await getRedis(); + await r.hSet("iot:latest", event.sensor_id, JSON.stringify({ + value: event.value, unit: event.unit, location: event.location, timestamp: event.timestamp, + })); + await r.hIncrBy("iot:counters", event.device_type, 1); + await r.zAdd("iot:recent", [{ score: 0, value: JSON.stringify(event) }]); + await r.zRemRangeByRank("iot:recent", 0, -1001); + } catch (e) { console.error("Redis error:", e.message); } + + try { + const coll = await getMongo(); + await coll.insertOne({ + event_id: event.event_id, sensor_id: event.sensor_id, + device_type: event.device_type, location: event.location, + value: event.value, unit: event.unit, timestamp: new Date(event.timestamp), + }); + } catch (e) { console.error("Mongo error:", e.message); } } async function start() { @@ -62,16 +67,16 @@ async function start() { await ch.assertQueue(RMQ_QUEUE, { durable: true }); ch.prefetch(10); await ch.consume(RMQ_QUEUE, async (msg) => { + if (!msg) return; try { - if (msg) { - const event = JSON.parse(msg.content.toString()); - await store(event); - consumed++; - ch.ack(msg); - if (consumed % 10 === 0) console.log(`Consumed ${consumed} events`); - } + const event = JSON.parse(msg.content.toString()); + await store(event); + consumed++; + ch.ack(msg); + if (consumed % 10 === 0) console.log(`Consumed ${consumed} events`); } catch (e) { errors++; + ch.nack(msg, false, true); console.error("Error:", e.message); } });