Deploy: 2026-07-15 04:57:29
This commit is contained in:
@@ -0,0 +1,68 @@
|
|||||||
|
# تعريف نظام "نبيه" (Nabeh SaaS Platform) 🚀
|
||||||
|
|
||||||
|
## نظرة عامة (Overview)
|
||||||
|
**نبيه (Nabeh)** هو منصة برمجيات كخدمة (SaaS - Software as a Service) متكاملة ومصممة للشركات والمتاجر الإلكترونية وشركات التوصيل الذكي (مثل تطبيق "سيرو" Siro). تهدف المنصة إلى أتمتة خدمة العملاء وإدارة التفاعلات عبر تطبيق **WhatsApp** باستخدام تقنيات الذكاء الاصطناعي المتقدمة (AI) والربط البرمجي السلس.
|
||||||
|
|
||||||
|
يعتمد المشروع على نموذج تجاري قائم على الاشتراكات (Subscription-based)، حيث يوفر للشركات باقات مختلفة بناءً على احتياجاتهم من الرسائل، روبوتات الدردشة، وتقنية التعرف البصري (OCR).
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## الميزات الأساسية والقدرات (Core Features)
|
||||||
|
|
||||||
|
### 1. نظام متعدد المستأجرين (Multi-Tenant Architecture)
|
||||||
|
- **إدارة الشركات (Companies):** يمكن للمنصة استضافة عدد غير محدود من الشركات والمتاجر بشكل مستقل تماماً.
|
||||||
|
- **باقات الاشتراك (Subscription Plans):** نظام فوترة وباقات (Starter, Growth, Professional) يحدّد لكل شركة:
|
||||||
|
- عدد أرقام الواتساب المسموحة (Sessions).
|
||||||
|
- الحد الأقصى للرسائل النصية والصوتية (Limits).
|
||||||
|
- الحد الأقصى لعمليات الذكاء الاصطناعي وتحليل الصور (OCR Limits).
|
||||||
|
|
||||||
|
### 2. محرك واتساب المركزي (WhatsApp Gateway)
|
||||||
|
- **الربط المباشر (Baileys Node):** ربط أرقام الواتساب عبر مسح رمز الاستجابة السريعة (QR Code) دون الحاجة للمرور بخوادم طرف ثالث، مما يضمن الخصوصية والسرعة.
|
||||||
|
- **إدارة الجلسات (Sessions Management):** يمكن للشركة الواحدة تشغيل عدة أرقام واتساب في نفس الوقت (مثل: رقم للمبيعات، ورقم للدعم الفني).
|
||||||
|
- **الحملات الإعلانية (Broadcast Campaigns):** إمكانية جدولة وإرسال رسائل جماعية للعملاء مع دعم القوالب الديناميكية (Templates).
|
||||||
|
|
||||||
|
### 3. وكيل الذكاء الاصطناعي المتقدم (AI Agent & Chatbot)
|
||||||
|
- **تكامل Google Gemini:** الاعتماد على نماذج لغوية متطورة (LLM) للرد بذكاء على استفسارات العملاء بدلاً من الردود المبرمجة مسبقاً (Keywords).
|
||||||
|
- **تخصيص الهوية (Custom Prompting):** كل شركة قادرة على تزويد الروبوت بتعليمات مخصصة (الاسم، اللهجة، ساعات العمل، سياسات الإرجاع)، ليعمل وكأنه موظف خدمة عملاء بشري حقيقي.
|
||||||
|
- **نموذج "الخدمة الجاهزة" (Plug & Play):** بأسلوب SaaS التجاري، يتم إخفاء تعقيدات المفاتيح البرمجية (API Keys) عن العميل النهائي، حيث تتكفل المنصة باستهلاك الـ API وإدارة الحصة (Quota) بناءً على اشتراك التاجر.
|
||||||
|
- **دعم الرسائل الصوتية:** قدرة الذكاء الاصطناعي على استلام الملاحظات الصوتية وتحليلها، وكذلك الرد برسائل صوتية بفضل تقنيات (Text-to-Speech).
|
||||||
|
|
||||||
|
### 4. تسجيل الكباتن والتعرف البصري (Driver Onboarding & OCR)
|
||||||
|
- تم تخصيص جزء من النظام لخدمة تطبيقات التوصيل الذكي (مثل تطبيق "سيرو" في سوريا والأردن ومصر).
|
||||||
|
- **المعالجة الآلية للوثائق (Automated OCR):** يستطيع النظام استلام صور وثائق السائقين (الهوية الوطنية، رخصة القيادة، رخصة المركبة، لا حكم عليه) عبر الواتساب.
|
||||||
|
- **استخراج البيانات وتدقيقها:** يقوم الذكاء الاصطناعي بقراءة وتفريغ محتويات الصور إلى بيانات مهيكلة (JSON) مثل الرقم الوطني، نوع السيارة، ورقم اللوحة.
|
||||||
|
- **التكامل مع السيرفر الرئيسي (Siro API):** بعد نجاح التدقيق، يرسل "نبيه" كافة البيانات والصور إلى السيرفر الرئيسي لتطبيق سيرو لتفعيل السائق تلقائياً.
|
||||||
|
|
||||||
|
### 5. التجارة الإلكترونية (E-Commerce Integrations)
|
||||||
|
- **ربط المتاجر (WooCommerce & Salla):** إمكانية ربط متجر العميل بشكل مباشر مع المنصة عبر الـ Webhooks.
|
||||||
|
- **تتبع الطلبات الآلي:** يستطيع العميل الاستعلام عن حالة طلبه عبر الواتساب، ليقوم الذكاء الاصطناعي بجلب حالة الشحنة من المتجر وإرسالها بطريقة سلسة.
|
||||||
|
- **إشعارات تلقائية:** إرسال تحديثات أوتوماتيكية للعملاء على الواتساب عند تغيير حالة الطلب (تم التنفيذ، جاري التوصيل، الخ...).
|
||||||
|
|
||||||
|
### 6. لوحة تحكم الإدارة (Super Admin Dashboard)
|
||||||
|
- واجهة خاصة بمالك المنصة للتحكم بكافة العملاء والمتاجر.
|
||||||
|
- مراقبة الإحصائيات الحية (عدد الشركات، الرسائل المرسلة، استهلاك الـ API).
|
||||||
|
- إدارة اشتراكات الشركات وتجديد باقاتهم.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## البنية التقنية (Tech Stack)
|
||||||
|
1. **الواجهة الأمامية (Frontend):**
|
||||||
|
- تم بناء واجهة سهلة وخفيفة الاستخدام بالاعتماد على HTML5/CSS3 المخصصة بالكامل.
|
||||||
|
- يعتمد على `Alpine.js` كإطار تفاعلي لتقليل الحمل وزيادة سرعة الأداء والاستجابة.
|
||||||
|
2. **الواجهة الخلفية (Backend):**
|
||||||
|
- تم تطويره باستخدام لغة `PHP` (هندسة MVC حديثة بدون إطار عمل ثقيل).
|
||||||
|
- تصميم موجّه للخدمات الدقيقة (Micro-Services architecture).
|
||||||
|
3. **بوابة الواتساب (WhatsApp Gateway):**
|
||||||
|
- مبني بواسطة `Node.js` مع مكتبة `Baileys` للتعامل مع بروتوكولات الواتساب مباشرة.
|
||||||
|
4. **قاعدة البيانات (Database):**
|
||||||
|
- خادم `MySQL` بتصميم Multi-tenant يفصل بيانات كل شركة بدقة، مع الاعتماد على تشفير البيانات الحساسة (مثل أرقام الهواتف) باستخدام `AES-256-GCM`.
|
||||||
|
5. **الذكاء الاصطناعي (AI & OCR):**
|
||||||
|
- الاعتماد على `Google Gemini API` لفهم اللغات وتحليل الصور والملفات الصوتية.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## نموذج العمل الموصى به (Business Model Proposal)
|
||||||
|
يمكن تسويق النظام كنظام SaaS باشتراكات شهرية، بحيث:
|
||||||
|
- **الباقة التأسيسية:** تستهدف أصحاب الأعمال الصغيرة للردود الآلية النصية البسيطة.
|
||||||
|
- **باقة النمو:** تستهدف المتاجر الإلكترونية للرد على استفسارات العملاء الصوتية وربط المتاجر وإدارة الحملات.
|
||||||
|
- **الباقة الاحترافية:** للشركات الكبيرة (كشركات التوصيل) والتي تحتاج لعمليات OCR ضخمة وربط مخصص للموظفين والكباتن.
|
||||||
@@ -6,6 +6,7 @@ const NodeCache = require('node-cache');
|
|||||||
const axios = require('axios');
|
const axios = require('axios');
|
||||||
const fs = require('fs');
|
const fs = require('fs');
|
||||||
const path = require('path');
|
const path = require('path');
|
||||||
|
const throttleManager = require('./throttle-manager');
|
||||||
|
|
||||||
class InMemoryStore {
|
class InMemoryStore {
|
||||||
constructor() {
|
constructor() {
|
||||||
@@ -573,44 +574,40 @@ function convertToOggOpus(base64Audio) {
|
|||||||
/**
|
/**
|
||||||
* Send a message using an active session
|
* Send a message using an active session
|
||||||
*/
|
*/
|
||||||
async function sendMessage(session_key, phone, message, mediaUrl = null, audioBase64 = null, mimetype = null, imageBase64 = null) {
|
/**
|
||||||
const sock = sessions.get(session_key);
|
* Resolve the correct JID for a phone number (handles LID privacy scheme)
|
||||||
if (!sock) {
|
*/
|
||||||
throw new Error(`Session ${session_key} is not active or connected`);
|
async function resolveJid(sock, phone) {
|
||||||
}
|
|
||||||
|
|
||||||
// Use the LID JID if we have a mapping for this phone number.
|
|
||||||
// This is CRITICAL: Signal E2EE sessions are bound to the LID JID,
|
|
||||||
// so replying to the phone JID creates a separate session that
|
|
||||||
// conflicts with the LID session, causing "Waiting for this message".
|
|
||||||
let jid;
|
|
||||||
if (phone.includes('@')) {
|
if (phone.includes('@')) {
|
||||||
jid = phone;
|
return phone;
|
||||||
} else if (phoneToLid.has(phone)) {
|
|
||||||
jid = phoneToLid.get(phone);
|
|
||||||
console.log(`[LID] Routing reply to ${phone} via LID: ${jid}`);
|
|
||||||
} else {
|
|
||||||
// Query WhatsApp to resolve the correct JID/LID for this phone number.
|
|
||||||
// This is CRITICAL for newer WhatsApp versions (especially iOS) to prevent
|
|
||||||
// E2EE session conflicts and the "Waiting for this message" decryption delay.
|
|
||||||
try {
|
|
||||||
console.log(`[JID Resolve] Resolving JID for ${phone}...`);
|
|
||||||
const [result] = await sock.onWhatsApp(phone);
|
|
||||||
if (result && result.exists) {
|
|
||||||
jid = result.jid;
|
|
||||||
console.log(`[JID Resolve] Successfully resolved JID for ${phone}: ${jid}`);
|
|
||||||
// Cache it for subsequent messages
|
|
||||||
if (jid.endsWith('@lid')) {
|
|
||||||
phoneToLid.set(phone, jid);
|
|
||||||
}
|
|
||||||
} else {
|
|
||||||
jid = `${phone}@s.whatsapp.net`;
|
|
||||||
}
|
|
||||||
} catch (err) {
|
|
||||||
console.error(`[JID Resolve] Failed to query onWhatsApp for ${phone}:`, err.message);
|
|
||||||
jid = `${phone}@s.whatsapp.net`;
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
if (phoneToLid.has(phone)) {
|
||||||
|
const lid = phoneToLid.get(phone);
|
||||||
|
console.log(`[LID] Routing reply to ${phone} via LID: ${lid}`);
|
||||||
|
return lid;
|
||||||
|
}
|
||||||
|
// Query WhatsApp to resolve the correct JID/LID for this phone number.
|
||||||
|
try {
|
||||||
|
console.log(`[JID Resolve] Resolving JID for ${phone}...`);
|
||||||
|
const [result] = await sock.onWhatsApp(phone);
|
||||||
|
if (result && result.exists) {
|
||||||
|
console.log(`[JID Resolve] Successfully resolved JID for ${phone}: ${result.jid}`);
|
||||||
|
if (result.jid.endsWith('@lid')) {
|
||||||
|
phoneToLid.set(phone, result.jid);
|
||||||
|
}
|
||||||
|
return result.jid;
|
||||||
|
}
|
||||||
|
} catch (err) {
|
||||||
|
console.error(`[JID Resolve] Failed to query onWhatsApp for ${phone}:`, err.message);
|
||||||
|
}
|
||||||
|
return `${phone}@s.whatsapp.net`;
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Internal raw send — performs the actual sock.sendMessage without throttling.
|
||||||
|
* Used inside the throttle queue's sendFn callback.
|
||||||
|
*/
|
||||||
|
async function _rawSend(sock, jid, message, mediaUrl, audioBase64, mimetype, imageBase64) {
|
||||||
let sentMsg;
|
let sentMsg;
|
||||||
|
|
||||||
if (imageBase64) {
|
if (imageBase64) {
|
||||||
@@ -631,7 +628,6 @@ async function sendMessage(session_key, phone, message, mediaUrl = null, audioBa
|
|||||||
finalMime = 'audio/ogg; codecs=opus';
|
finalMime = 'audio/ogg; codecs=opus';
|
||||||
} catch (err) {
|
} catch (err) {
|
||||||
console.error(`[Baileys] FFmpeg conversion failed:`, err.message);
|
console.error(`[Baileys] FFmpeg conversion failed:`, err.message);
|
||||||
// Fallback to sending as normal audio if conversion fails
|
|
||||||
finalMime = 'audio/mp4';
|
finalMime = 'audio/mp4';
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -642,7 +638,7 @@ async function sendMessage(session_key, phone, message, mediaUrl = null, audioBa
|
|||||||
sentMsg = await sock.sendMessage(jid, {
|
sentMsg = await sock.sendMessage(jid, {
|
||||||
audio: buffer,
|
audio: buffer,
|
||||||
mimetype: finalMime,
|
mimetype: finalMime,
|
||||||
ptt: !isMp3 // PTT enabled for OGG/MP4, disabled for raw MP3
|
ptt: !isMp3
|
||||||
});
|
});
|
||||||
} else if (mediaUrl) {
|
} else if (mediaUrl) {
|
||||||
const ext = mediaUrl.split('.').pop().toLowerCase();
|
const ext = mediaUrl.split('.').pop().toLowerCase();
|
||||||
@@ -669,6 +665,37 @@ async function sendMessage(session_key, phone, message, mediaUrl = null, audioBa
|
|||||||
return sentMsg;
|
return sentMsg;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Send a message using an active session — routed through ThrottleManager
|
||||||
|
* for human-like delays, typing indicators, and anti-ban protection.
|
||||||
|
*/
|
||||||
|
async function sendMessage(session_key, phone, message, mediaUrl = null, audioBase64 = null, mimetype = null, imageBase64 = null) {
|
||||||
|
const sock = sessions.get(session_key);
|
||||||
|
if (!sock) {
|
||||||
|
throw new Error(`Session ${session_key} is not active or connected`);
|
||||||
|
}
|
||||||
|
|
||||||
|
const jid = await resolveJid(sock, phone);
|
||||||
|
|
||||||
|
// Determine message type for throttle delay calculation
|
||||||
|
let msgType = 'auto_reply';
|
||||||
|
if (audioBase64) msgType = 'voice_reply';
|
||||||
|
|
||||||
|
// Build the content descriptor for typing delay calculation
|
||||||
|
const content = { text: message || '' };
|
||||||
|
|
||||||
|
// Enqueue through ThrottleManager for anti-ban protection
|
||||||
|
const sentMsg = await throttleManager.enqueue(session_key, {
|
||||||
|
jid,
|
||||||
|
sock,
|
||||||
|
content,
|
||||||
|
type: msgType,
|
||||||
|
sendFn: () => _rawSend(sock, jid, message, mediaUrl, audioBase64, mimetype, imageBase64)
|
||||||
|
});
|
||||||
|
|
||||||
|
return sentMsg;
|
||||||
|
}
|
||||||
|
|
||||||
async function checkContact(session_key, phone) {
|
async function checkContact(session_key, phone) {
|
||||||
const sock = sessions.get(session_key);
|
const sock = sessions.get(session_key);
|
||||||
if (!sock) {
|
if (!sock) {
|
||||||
@@ -770,6 +797,7 @@ module.exports = {
|
|||||||
sendMessage,
|
sendMessage,
|
||||||
getActiveSessions,
|
getActiveSessions,
|
||||||
checkContact,
|
checkContact,
|
||||||
exportChatHistory
|
exportChatHistory,
|
||||||
|
throttleManager
|
||||||
};
|
};
|
||||||
|
|
||||||
|
|||||||
@@ -20,7 +20,7 @@ for (const p of envPaths) {
|
|||||||
|
|
||||||
const express = require('express');
|
const express = require('express');
|
||||||
const cors = require('cors');
|
const cors = require('cors');
|
||||||
const { startSession, disconnectSession, sendMessage, getActiveSessions, checkContact, exportChatHistory } = require('./baileys-client');
|
const { startSession, disconnectSession, sendMessage, getActiveSessions, checkContact, exportChatHistory, throttleManager } = require('./baileys-client');
|
||||||
|
|
||||||
const app = express();
|
const app = express();
|
||||||
app.use(cors());
|
app.use(cors());
|
||||||
@@ -82,6 +82,11 @@ app.get('/api/sessions/active', (req, res) => {
|
|||||||
res.json({ status: 'success', active_sessions: getActiveSessions() });
|
res.json({ status: 'success', active_sessions: getActiveSessions() });
|
||||||
});
|
});
|
||||||
|
|
||||||
|
// Get throttle queue status for all sessions (Anti-Ban monitoring)
|
||||||
|
app.get('/api/queue/status', (req, res) => {
|
||||||
|
res.json({ status: 'success', data: throttleManager.getAllStatus() });
|
||||||
|
});
|
||||||
|
|
||||||
// Check if contact is on WhatsApp
|
// Check if contact is on WhatsApp
|
||||||
app.post('/api/contacts/check', async (req, res) => {
|
app.post('/api/contacts/check', async (req, res) => {
|
||||||
const { session_key, phone } = req.body;
|
const { session_key, phone } = req.body;
|
||||||
|
|||||||
@@ -0,0 +1,381 @@
|
|||||||
|
/**
|
||||||
|
* ThrottleManager — Smart Anti-Ban Message Throttling Engine for Nabeh WhatsApp Gateway
|
||||||
|
*
|
||||||
|
* Simulates human-like messaging behavior to prevent WhatsApp bans:
|
||||||
|
* - Per-session message queues with intelligent delays
|
||||||
|
* - Typing indicators before sending
|
||||||
|
* - Random jitter to avoid pattern detection
|
||||||
|
* - Rate limiting per minute/hour/day
|
||||||
|
* - Automatic slowdown on warnings
|
||||||
|
* - Emergency pause on critical risk
|
||||||
|
*/
|
||||||
|
|
||||||
|
class ThrottleManager {
|
||||||
|
constructor() {
|
||||||
|
// Per-session message queues
|
||||||
|
this.queues = new Map(); // sessionKey -> Array<QueueItem>
|
||||||
|
this.processing = new Map(); // sessionKey -> boolean (is currently processing)
|
||||||
|
this.stats = new Map(); // sessionKey -> { sentLastMinute, sentLastHour, sentToday, warnings }
|
||||||
|
this.slowdownLevel = new Map(); // sessionKey -> 0 (normal), 1 (caution), 2 (danger)
|
||||||
|
|
||||||
|
// Default limits (randomized per session to avoid fingerprinting)
|
||||||
|
this.limits = {
|
||||||
|
perMinute: { min: 6, max: 10 },
|
||||||
|
perHour: { min: 80, max: 120 },
|
||||||
|
perDay: { min: 400, max: 600 },
|
||||||
|
burstCooldownAfter: 5, // After 5 quick messages, take a break
|
||||||
|
burstCooldownMs: { min: 30000, max: 90000 } // 30-90 second break
|
||||||
|
};
|
||||||
|
|
||||||
|
// Cleanup old stats every minute
|
||||||
|
this._cleanupInterval = setInterval(() => this._cleanupStats(), 60000);
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Get or initialize stats for a session
|
||||||
|
*/
|
||||||
|
_getStats(sessionKey) {
|
||||||
|
if (!this.stats.has(sessionKey)) {
|
||||||
|
this.stats.set(sessionKey, {
|
||||||
|
sentTimestamps: [], // Array of timestamps for rate tracking
|
||||||
|
warnings: 0,
|
||||||
|
lastWarningAt: 0,
|
||||||
|
burstCount: 0, // Messages sent in quick succession
|
||||||
|
lastSendAt: 0,
|
||||||
|
dailyReset: this._todayKey()
|
||||||
|
});
|
||||||
|
}
|
||||||
|
const stats = this.stats.get(sessionKey);
|
||||||
|
// Reset daily counter at midnight
|
||||||
|
if (stats.dailyReset !== this._todayKey()) {
|
||||||
|
stats.sentTimestamps = stats.sentTimestamps.filter(
|
||||||
|
ts => Date.now() - ts < 3600000
|
||||||
|
);
|
||||||
|
stats.dailyReset = this._todayKey();
|
||||||
|
stats.warnings = Math.max(0, stats.warnings - 1); // Decay warnings daily
|
||||||
|
}
|
||||||
|
return stats;
|
||||||
|
}
|
||||||
|
|
||||||
|
_todayKey() {
|
||||||
|
return new Date().toISOString().slice(0, 10);
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Calculate messages sent within a time window
|
||||||
|
*/
|
||||||
|
_countInWindow(timestamps, windowMs) {
|
||||||
|
const cutoff = Date.now() - windowMs;
|
||||||
|
return timestamps.filter(ts => ts > cutoff).length;
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Clean up old timestamps from stats to prevent memory bloat
|
||||||
|
*/
|
||||||
|
_cleanupStats() {
|
||||||
|
const oneDayAgo = Date.now() - 86400000;
|
||||||
|
for (const [key, stats] of this.stats) {
|
||||||
|
stats.sentTimestamps = stats.sentTimestamps.filter(ts => ts > oneDayAgo);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Random integer between min and max (inclusive)
|
||||||
|
*/
|
||||||
|
_randomBetween(min, max) {
|
||||||
|
return Math.floor(Math.random() * (max - min + 1)) + min;
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Calculate human-like typing delay based on message length
|
||||||
|
* Average human types ~40 words/minute = ~200 chars/minute
|
||||||
|
* We simulate faster (AI is "quick") but still believable: ~400 chars/min
|
||||||
|
*/
|
||||||
|
_calculateTypingDelay(messageLength) {
|
||||||
|
if (!messageLength) return this._randomBetween(800, 1500);
|
||||||
|
// Base: 150ms per character, capped between 1-5 seconds
|
||||||
|
const baseDelay = Math.min(messageLength * 40, 5000);
|
||||||
|
return Math.max(800, baseDelay + this._randomBetween(-300, 500));
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Calculate delay before reading/responding to simulate human behavior
|
||||||
|
*/
|
||||||
|
_calculateReadDelay(messageType) {
|
||||||
|
switch (messageType) {
|
||||||
|
case 'auto_reply':
|
||||||
|
return this._randomBetween(1500, 4000); // 1.5-4s for AI replies
|
||||||
|
case 'voice_reply':
|
||||||
|
return this._randomBetween(3000, 7000); // 3-7s for voice (longer processing feel)
|
||||||
|
case 'broadcast':
|
||||||
|
return this._randomBetween(15000, 45000); // 15-45s between broadcast messages
|
||||||
|
case 'reminder':
|
||||||
|
return this._randomBetween(3000, 8000); // 3-8s for reminders
|
||||||
|
default:
|
||||||
|
return this._randomBetween(1500, 3500);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Check if we're within rate limits for this session
|
||||||
|
*/
|
||||||
|
_isWithinLimits(sessionKey) {
|
||||||
|
const stats = this._getStats(sessionKey);
|
||||||
|
const slowdown = this.slowdownLevel.get(sessionKey) || 0;
|
||||||
|
|
||||||
|
// Apply multiplier based on slowdown level
|
||||||
|
const multiplier = slowdown === 0 ? 1 : (slowdown === 1 ? 0.5 : 0.2);
|
||||||
|
|
||||||
|
const perMinLimit = Math.floor(this._randomBetween(this.limits.perMinute.min, this.limits.perMinute.max) * multiplier);
|
||||||
|
const perHourLimit = Math.floor(this._randomBetween(this.limits.perHour.min, this.limits.perHour.max) * multiplier);
|
||||||
|
const perDayLimit = Math.floor(this._randomBetween(this.limits.perDay.min, this.limits.perDay.max) * multiplier);
|
||||||
|
|
||||||
|
const sentLastMinute = this._countInWindow(stats.sentTimestamps, 60000);
|
||||||
|
const sentLastHour = this._countInWindow(stats.sentTimestamps, 3600000);
|
||||||
|
const sentToday = this._countInWindow(stats.sentTimestamps, 86400000);
|
||||||
|
|
||||||
|
if (sentLastMinute >= perMinLimit) {
|
||||||
|
console.log(`[Throttle] ${sessionKey} — Rate limit: ${sentLastMinute}/${perMinLimit} per minute. Waiting...`);
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
if (sentLastHour >= perHourLimit) {
|
||||||
|
console.log(`[Throttle] ${sessionKey} — Rate limit: ${sentLastHour}/${perHourLimit} per hour. Waiting...`);
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
if (sentToday >= perDayLimit) {
|
||||||
|
console.log(`[Throttle] ${sessionKey} — Rate limit: ${sentToday}/${perDayLimit} per day. PAUSED.`);
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Enqueue a message for throttled sending
|
||||||
|
* @param {string} sessionKey - Session identifier
|
||||||
|
* @param {object} messageData - { jid, content, type, sendFn }
|
||||||
|
* @returns {Promise<object>} - Resolves when the message is actually sent
|
||||||
|
*/
|
||||||
|
enqueue(sessionKey, messageData) {
|
||||||
|
return new Promise((resolve, reject) => {
|
||||||
|
if (!this.queues.has(sessionKey)) {
|
||||||
|
this.queues.set(sessionKey, []);
|
||||||
|
}
|
||||||
|
|
||||||
|
const queueItem = {
|
||||||
|
...messageData,
|
||||||
|
enqueuedAt: Date.now(),
|
||||||
|
resolve,
|
||||||
|
reject
|
||||||
|
};
|
||||||
|
|
||||||
|
this.queues.get(sessionKey).push(queueItem);
|
||||||
|
console.log(`[Throttle] ${sessionKey} — Message queued. Queue size: ${this.queues.get(sessionKey).length}`);
|
||||||
|
|
||||||
|
// Start processing if not already running
|
||||||
|
if (!this.processing.get(sessionKey)) {
|
||||||
|
this._processQueue(sessionKey);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Process the message queue for a session sequentially
|
||||||
|
*/
|
||||||
|
async _processQueue(sessionKey) {
|
||||||
|
if (this.processing.get(sessionKey)) return;
|
||||||
|
this.processing.set(sessionKey, true);
|
||||||
|
|
||||||
|
const queue = this.queues.get(sessionKey);
|
||||||
|
|
||||||
|
while (queue && queue.length > 0) {
|
||||||
|
const item = queue[0];
|
||||||
|
|
||||||
|
// Check rate limits
|
||||||
|
if (!this._isWithinLimits(sessionKey)) {
|
||||||
|
// Wait and retry
|
||||||
|
await this._sleep(this._randomBetween(5000, 15000));
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
|
||||||
|
// Check burst cooldown
|
||||||
|
const stats = this._getStats(sessionKey);
|
||||||
|
if (stats.burstCount >= this.limits.burstCooldownAfter) {
|
||||||
|
const cooldown = this._randomBetween(
|
||||||
|
this.limits.burstCooldownMs.min,
|
||||||
|
this.limits.burstCooldownMs.max
|
||||||
|
);
|
||||||
|
console.log(`[Throttle] ${sessionKey} — Burst cooldown: ${Math.round(cooldown / 1000)}s pause after ${stats.burstCount} quick messages.`);
|
||||||
|
stats.burstCount = 0;
|
||||||
|
await this._sleep(cooldown);
|
||||||
|
}
|
||||||
|
|
||||||
|
// Calculate pre-send delays
|
||||||
|
const messageType = item.type || 'auto_reply';
|
||||||
|
const readDelay = this._calculateReadDelay(messageType);
|
||||||
|
|
||||||
|
// 1. Send "read" receipt if applicable
|
||||||
|
if (item.sock && item.incomingMsgKey) {
|
||||||
|
try {
|
||||||
|
await item.sock.readMessages([item.incomingMsgKey]);
|
||||||
|
} catch (e) {
|
||||||
|
// Non-critical, continue
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// 2. Wait (simulating reading the message)
|
||||||
|
await this._sleep(readDelay);
|
||||||
|
|
||||||
|
// 3. Send typing indicator
|
||||||
|
if (item.sock && item.jid) {
|
||||||
|
try {
|
||||||
|
await item.sock.sendPresenceUpdate('composing', item.jid);
|
||||||
|
} catch (e) {
|
||||||
|
// Non-critical
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// 4. Wait for typing delay
|
||||||
|
const textLength = item.content?.text?.length || item.content?.caption?.length || 20;
|
||||||
|
const typingDelay = this._calculateTypingDelay(textLength);
|
||||||
|
await this._sleep(typingDelay);
|
||||||
|
|
||||||
|
// 5. Stop typing indicator
|
||||||
|
if (item.sock && item.jid) {
|
||||||
|
try {
|
||||||
|
await item.sock.sendPresenceUpdate('paused', item.jid);
|
||||||
|
} catch (e) {
|
||||||
|
// Non-critical
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// 6. Actually send the message
|
||||||
|
try {
|
||||||
|
const result = await item.sendFn();
|
||||||
|
stats.sentTimestamps.push(Date.now());
|
||||||
|
stats.lastSendAt = Date.now();
|
||||||
|
stats.burstCount++;
|
||||||
|
|
||||||
|
// Reset burst if last send was more than 10 seconds ago
|
||||||
|
if (Date.now() - stats.lastSendAt > 10000) {
|
||||||
|
stats.burstCount = 1;
|
||||||
|
}
|
||||||
|
|
||||||
|
queue.shift(); // Remove processed item
|
||||||
|
item.resolve(result);
|
||||||
|
console.log(`[Throttle] ${sessionKey} — Message sent successfully. Queue remaining: ${queue.length}`);
|
||||||
|
} catch (err) {
|
||||||
|
queue.shift();
|
||||||
|
item.reject(err);
|
||||||
|
console.error(`[Throttle] ${sessionKey} — Send failed:`, err.message);
|
||||||
|
|
||||||
|
// Check if it's a rate-limit related error
|
||||||
|
if (err.message?.includes('rate') || err.message?.includes('429') || err.message?.includes('too many')) {
|
||||||
|
this.onWarning(sessionKey);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Small random gap between processing queue items
|
||||||
|
if (queue.length > 0) {
|
||||||
|
await this._sleep(this._randomBetween(500, 2000));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
this.processing.set(sessionKey, false);
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Called when a warning signal is detected (rate limit, 429, etc.)
|
||||||
|
*/
|
||||||
|
onWarning(sessionKey) {
|
||||||
|
const stats = this._getStats(sessionKey);
|
||||||
|
stats.warnings++;
|
||||||
|
stats.lastWarningAt = Date.now();
|
||||||
|
|
||||||
|
const currentLevel = this.slowdownLevel.get(sessionKey) || 0;
|
||||||
|
|
||||||
|
if (stats.warnings >= 5) {
|
||||||
|
// EMERGENCY: Pause all sending for this session
|
||||||
|
this.slowdownLevel.set(sessionKey, 2);
|
||||||
|
console.warn(`[Throttle] ⛔ ${sessionKey} — EMERGENCY SLOWDOWN! ${stats.warnings} warnings detected. All sending severely throttled.`);
|
||||||
|
} else if (stats.warnings >= 2) {
|
||||||
|
this.slowdownLevel.set(sessionKey, 1);
|
||||||
|
console.warn(`[Throttle] ⚠️ ${sessionKey} — CAUTION slowdown. ${stats.warnings} warnings. Reducing send rate by 50%.`);
|
||||||
|
}
|
||||||
|
|
||||||
|
// Auto-recover after 30 minutes of no warnings
|
||||||
|
setTimeout(() => {
|
||||||
|
const current = this._getStats(sessionKey);
|
||||||
|
if (Date.now() - current.lastWarningAt > 1800000) {
|
||||||
|
const level = this.slowdownLevel.get(sessionKey) || 0;
|
||||||
|
if (level > 0) {
|
||||||
|
this.slowdownLevel.set(sessionKey, Math.max(0, level - 1));
|
||||||
|
console.log(`[Throttle] ✅ ${sessionKey} — Slowdown level decreased to ${level - 1} after 30min recovery.`);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}, 1800000);
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Get queue status for monitoring
|
||||||
|
*/
|
||||||
|
getStatus(sessionKey) {
|
||||||
|
const stats = this._getStats(sessionKey);
|
||||||
|
const queue = this.queues.get(sessionKey) || [];
|
||||||
|
return {
|
||||||
|
queueSize: queue.length,
|
||||||
|
isProcessing: this.processing.get(sessionKey) || false,
|
||||||
|
slowdownLevel: this.slowdownLevel.get(sessionKey) || 0,
|
||||||
|
sentLastMinute: this._countInWindow(stats.sentTimestamps, 60000),
|
||||||
|
sentLastHour: this._countInWindow(stats.sentTimestamps, 3600000),
|
||||||
|
sentToday: this._countInWindow(stats.sentTimestamps, 86400000),
|
||||||
|
warnings: stats.warnings
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Get status for all sessions
|
||||||
|
*/
|
||||||
|
getAllStatus() {
|
||||||
|
const result = {};
|
||||||
|
for (const sessionKey of this.stats.keys()) {
|
||||||
|
result[sessionKey] = this.getStatus(sessionKey);
|
||||||
|
}
|
||||||
|
return result;
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Clear queue for a session (e.g., when disconnecting)
|
||||||
|
*/
|
||||||
|
clearQueue(sessionKey) {
|
||||||
|
const queue = this.queues.get(sessionKey) || [];
|
||||||
|
for (const item of queue) {
|
||||||
|
item.reject(new Error('Queue cleared — session disconnecting'));
|
||||||
|
}
|
||||||
|
this.queues.delete(sessionKey);
|
||||||
|
this.processing.delete(sessionKey);
|
||||||
|
console.log(`[Throttle] ${sessionKey} — Queue cleared.`);
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Promise-based sleep
|
||||||
|
*/
|
||||||
|
_sleep(ms) {
|
||||||
|
return new Promise(resolve => setTimeout(resolve, ms));
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Cleanup on shutdown
|
||||||
|
*/
|
||||||
|
destroy() {
|
||||||
|
if (this._cleanupInterval) {
|
||||||
|
clearInterval(this._cleanupInterval);
|
||||||
|
}
|
||||||
|
for (const sessionKey of this.queues.keys()) {
|
||||||
|
this.clearQueue(sessionKey);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Export singleton instance
|
||||||
|
const throttleManager = new ThrottleManager();
|
||||||
|
module.exports = throttleManager;
|
||||||
Reference in New Issue
Block a user