oclaw/runtime/operations/whatsapp_bridge/baileys_runner.ts
oliver 9cf026f605 Fix WhatsApp Admin QR bind and show the scan prompt only after Bind.
Use a current WA protocol version, start the sidecar via local tsx with inherited proxy, and stop auto-opening the QR UI on unbound connecting.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-17 23:54:48 +08:00

1519 lines
54 KiB
TypeScript
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

import fs from "node:fs";
import path from "node:path";
import process from "node:process";
import dns from "node:dns/promises";
import makeWASocket, {
Browsers,
DisconnectReason,
downloadMediaMessage,
extensionForMediaMessage,
fetchLatestBaileysVersion,
fetchLatestWaWebVersion,
getContentType,
jidNormalizedUser,
proto,
} from "@whiskeysockets/baileys";
import { HttpsProxyAgent } from "https-proxy-agent";
import { loadAuthState, resolveAuthDir } from "./auth";
import { printQrToTerminal, writeQrImage } from "./qr";
import { writeBridgeStatus } from "./status";
type Json = Record<string, unknown>;
const LOCAL_BASE_URL = (process.env.AIA_GATEWAY_BASE_URL || "http://127.0.0.1:8787").trim();
const ACCOUNT_ID = (process.env.AIA_WHATSAPP_ACCOUNT_ID || "wa-default").trim();
const STATE_DIR = (process.env.OCLAW_STATE_DIR || path.resolve(process.cwd(), "state")).trim();
const LOGIN_ONLY = process.argv.includes("--login") || String(process.env.AIA_WHATSAPP_LOGIN_ONLY || "").trim() === "1";
const VERBOSE = process.argv.includes("--verbose") || String(process.env.AIA_WHATSAPP_VERBOSE || "").trim() === "1";
const PROXY_URL = (
process.env.AIA_WHATSAPP_PROXY_URL ||
process.env.HTTPS_PROXY ||
process.env.HTTP_PROXY ||
process.env.https_proxy ||
process.env.http_proxy ||
""
).trim();
const MEDIA_MAX_BYTES = Number(process.env.OCLAW_WHATSAPP_MEDIA_MAX_BYTES || String(20 * 1024 * 1024)) || 20 * 1024 * 1024;
const MEDIA_DOWNLOAD_TIMEOUT_MS = Number(process.env.OCLAW_WHATSAPP_MEDIA_TIMEOUT_MS || "90000") || 90_000;
const TYPING_ENABLED =
String(process.env.OCLAW_WHATSAPP_TYPING || "1").trim() !== "0" &&
String(process.env.OCLAW_WHATSAPP_TYPING || "1").trim().toLowerCase() !== "false" &&
String(process.env.OCLAW_WHATSAPP_TYPING || "1").trim().toLowerCase() !== "off";
const TYPING_HEARTBEAT_MS = Math.max(
3000,
Number(process.env.OCLAW_WHATSAPP_TYPING_HEARTBEAT_MS || "8000") || 8000,
);
function log(msg: string): void {
process.stdout.write(`${new Date().toISOString()} [baileys-whatsapp] ${msg}\n`);
}
function ensureDir(dir: string): void {
fs.mkdirSync(dir, { recursive: true });
}
function pickText(m: proto.IMessage | null | undefined): string {
const u = unwrapMessage(m);
if (!u) return "";
const c = (u.conversation || "").trim();
if (c) return c;
const ext = (u.extendedTextMessage?.text || "").trim();
if (ext) return ext;
const imgCap = (u.imageMessage?.caption || "").trim();
if (imgCap) return imgCap;
const vidCap = (u.videoMessage?.caption || "").trim();
if (vidCap) return vidCap;
const docCap = (u.documentMessage?.caption || "").trim();
if (docCap) return docCap;
return "";
}
function unwrapMessage(m: proto.IMessage | null | undefined): proto.IMessage | null {
if (!m) return null;
const nested =
m.ephemeralMessage?.message ||
m.viewOnceMessage?.message ||
m.viewOnceMessageV2?.message ||
m.documentWithCaptionMessage?.message ||
m.editedMessage?.message ||
null;
if (nested) return unwrapMessage(nested);
return m;
}
function messageContextInfo(m: proto.IMessage | null | undefined): proto.IContextInfo | null | undefined {
const u = unwrapMessage(m);
if (!u) return null;
const extCtx = u.extendedTextMessage?.contextInfo;
// Reply+@ messages usually carry quote + mentions on extendedTextMessage.contextInfo.
if (extCtx && (extCtx.quotedMessage || extCtx.stanzaId)) return extCtx;
const top = (u as proto.IMessage & { messageContextInfo?: proto.IMessageContextInfo }).messageContextInfo;
const fromTop = top?.mentionedJid?.length ? top : null;
return (
(fromTop as unknown as proto.IContextInfo) ||
extCtx ||
u.imageMessage?.contextInfo ||
u.videoMessage?.contextInfo ||
u.documentMessage?.contextInfo ||
u.buttonsResponseMessage?.contextInfo ||
u.listResponseMessage?.contextInfo ||
u.templateButtonReplyMessage?.contextInfo ||
null
);
}
function extractMentionsFromUpsert(msg: proto.IWebMessageInfo): string[] {
const fromBody = extractMentions(msg.message);
const outer = (msg as proto.IWebMessageInfo & { messageContextInfo?: { mentionedJid?: string[] } })
.messageContextInfo?.mentionedJid;
const outerList = Array.isArray(outer) ? outer.map((j) => String(j || "").trim()).filter(Boolean) : [];
const merged = [...fromBody, ...outerList];
const seen = new Set<string>();
const out: string[] = [];
for (const j of merged) {
const key = jidPhone(j) || j;
if (seen.has(key)) continue;
seen.add(key);
out.push(j);
}
return out;
}
function jidPhone(jid: string): string {
const head = String(jid || "").split("@")[0]?.split(":")[0] || "";
return head.replace(/\D/g, "");
}
function jidsSameUser(a: string, b: string): boolean {
const na = jidNormalizedUser(String(a || "").trim());
const nb = jidNormalizedUser(String(b || "").trim());
if (na && nb && na === nb) return true;
const pa = jidPhone(a);
const pb = jidPhone(b);
return pa.length >= 6 && pa === pb;
}
function extractMentions(m: proto.IMessage | null | undefined): string[] {
const ctx = messageContextInfo(m);
const raw = ctx?.mentionedJid;
if (!Array.isArray(raw)) return [];
return raw.map((j) => String(j || "").trim()).filter(Boolean);
}
function pickQuotedText(m: proto.IMessage | null | undefined): string {
const ctx = messageContextInfo(m);
const qm = ctx?.quotedMessage;
if (!qm) return "";
return pickText(qm);
}
function extractQuoteContext(m: proto.IMessage | null | undefined): {
participant: string;
stanzaId: string;
quotedText: string;
} {
const ctx = messageContextInfo(m);
return {
participant: String(ctx?.participant || "").trim(),
stanzaId: String(ctx?.stanzaId || "").trim(),
quotedText: pickQuotedText(m),
};
}
function resolveSenderJid(key: proto.IMessageKey): string {
const participant = String(key.participant || "").trim();
const participantAlt = String((key as any).participantAlt || "").trim();
if (participantAlt && participant.toLowerCase().endsWith("@lid")) {
return jidNormalizedUser(participantAlt);
}
if (participant) return jidNormalizedUser(participant);
return "";
}
function resolvePnFromLid(sock: ReturnType<typeof makeWASocket> | null, lidJid: string): string {
const lid = String(lidJid || "").trim();
if (!lid || !sock) return "";
const lidMapping = (sock as any)?.signalRepository?.lidMapping;
if (!lidMapping || typeof lidMapping.getPNForLID !== "function") return "";
try {
const pn = lidMapping.getPNForLID(lid);
return pn ? jidNormalizedUser(String(pn)) : "";
} catch {
return "";
}
}
function resolveDmUserJid(key: proto.IMessageKey, sock: ReturnType<typeof makeWASocket> | null): string {
const remoteJid = jidNormalizedUser(String(key.remoteJid || ""));
const remoteJidAlt = String((key as any).remoteJidAlt || "").trim();
if (remoteJidAlt) {
const alt = jidNormalizedUser(remoteJidAlt);
if (remoteJid.toLowerCase().endsWith("@lid") && alt.toLowerCase().includes("@s.whatsapp")) {
return alt;
}
if (alt.toLowerCase().includes("@s.whatsapp")) return alt;
}
if (remoteJid.toLowerCase().endsWith("@lid")) {
const pn = resolvePnFromLid(sock, remoteJid);
if (pn) return pn;
}
return remoteJid;
}
function jidBaseLocal(jid: string): string {
const s = String(jid || "").trim().toLowerCase();
if (!s) return "";
return s.split("@")[0]?.split(":")[0] || "";
}
function mentionMatchesBotIdentity(mention: string, botId: string): boolean {
const m = String(mention || "").trim();
const b = String(botId || "").trim();
if (!m || !b) return false;
if (jidNormalizedUser(m) === jidNormalizedUser(b)) return true;
const mLid = m.toLowerCase().endsWith("@lid");
const bLid = b.toLowerCase().endsWith("@lid");
if (mLid || bLid) {
const ml = jidBaseLocal(m);
const bl = jidBaseLocal(b);
if (ml && bl && ml === bl && ml.replace(/\D/g, "").length >= 6) return true;
return false;
}
const mp = jidPhone(m);
const bp = jidPhone(b);
return (
mp.length >= 6 &&
mp === bp &&
m.toLowerCase().includes("@s.whatsapp") &&
b.toLowerCase().includes("@s.whatsapp")
);
}
function collectBotIdentityJids(
sock: ReturnType<typeof makeWASocket> | null,
botPn: string,
authCreds?: { me?: { id?: string; lid?: string } | null } | null,
): string[] {
const out = new Set<string>();
const pn = String(botPn || "").trim();
if (pn) {
out.add(pn);
out.add(jidNormalizedUser(pn));
const base = jidBaseLocal(pn);
if (base) out.add(`${base}@s.whatsapp.net`);
}
const user = (sock as any)?.user;
if (user?.id) out.add(String(user.id).trim());
if (user?.lid) out.add(String(user.lid).trim());
const creds = authCreds || (sock as any)?.authState?.creds;
const me = creds?.me;
if (me?.id) out.add(String(me.id).trim());
if (me?.lid) out.add(String(me.lid).trim());
const lidMapping = (sock as any)?.signalRepository?.lidMapping;
if (lidMapping && pn) {
for (const variant of Array.from(out)) {
if (!variant.includes("@s.whatsapp")) continue;
try {
if (typeof lidMapping.getLIDForPN === "function") {
const lid = lidMapping.getLIDForPN(variant);
if (lid) out.add(String(lid));
}
} catch {
// ignore
}
}
}
return Array.from(out).filter(Boolean);
}
function pickBotLid(botIds: string[]): string {
for (const id of botIds) {
if (String(id || "").toLowerCase().includes("@lid")) return String(id);
}
return "";
}
function messageMentionsBot(
sock: ReturnType<typeof makeWASocket> | null,
mentions: string[],
botPn: string,
botIds?: string[],
authCreds?: { me?: { id?: string; lid?: string } | null } | null,
): boolean {
const ids = botIds?.length ? botIds : collectBotIdentityJids(sock, botPn, authCreds);
if (!ids.length || !mentions.length) return false;
for (const m of mentions) {
const mention = String(m || "").trim();
if (!mention) continue;
for (const botId of ids) {
if (mentionMatchesBotIdentity(mention, botId)) return true;
}
}
return false;
}
function isStatusOrBroadcastJid(jid: string): boolean {
const low = String(jid || "").toLowerCase();
return low === "status@broadcast" || low.endsWith("@broadcast");
}
function inferExtension(contentType: string, originalFilename: string): string {
const explicit = path.extname(String(originalFilename || "").trim()).toLowerCase();
if (explicit && explicit.length <= 8) return explicit;
const low = String(contentType || "").trim().toLowerCase();
if (low.includes("jpeg") || low.includes("jpg")) return ".jpg";
if (low.includes("png")) return ".png";
if (low.includes("gif")) return ".gif";
if (low.includes("webp")) return ".webp";
if (low.includes("bmp")) return ".bmp";
if (low.includes("tiff")) return ".tiff";
if (low.includes("mp4")) return ".mp4";
if (low.includes("webm")) return ".webm";
if (low.includes("quicktime") || low.includes("mov")) return ".mov";
if (low.includes("pdf")) return ".pdf";
if (low.includes("msword") || low.includes("wordprocessingml")) return ".docx";
if (low.includes("spreadsheetml") || low.includes("ms-excel")) return ".xlsx";
if (low.includes("presentationml") || low.includes("ms-powerpoint")) return ".pptx";
if (low.includes("zip")) return ".zip";
if (low.includes("gzip")) return ".gz";
if (low.includes("json")) return ".json";
if (low.includes("csv")) return ".csv";
if (low.includes("plain")) return ".txt";
if (low.includes("opus") || low.includes("ogg")) return ".ogg";
if (low.includes("mpeg") || low.includes("mp3")) return ".mp3";
if (low.includes("wav")) return ".wav";
if (low.includes("m4a") || low.includes("mp4a")) return ".m4a";
if (low.includes("aac")) return ".aac";
return ".bin";
}
function isVoiceNote(m: proto.IMessage | null | undefined): boolean {
const u = unwrapMessage(m);
if (!u?.audioMessage) return false;
return Boolean(u.audioMessage.ptt);
}
function mediaMimeFromMessage(m: proto.IMessage | null | undefined): string {
const u = unwrapMessage(m);
if (!u) return "";
const raw = String(
u.imageMessage?.mimetype ||
u.videoMessage?.mimetype ||
u.documentMessage?.mimetype ||
u.audioMessage?.mimetype ||
u.stickerMessage?.mimetype ||
"",
).trim();
if (raw) return raw;
if (u.audioMessage) return isVoiceNote(m) ? "audio/ogg; codecs=opus" : "audio/mpeg";
if (u.imageMessage || u.stickerMessage) return "image/jpeg";
if (u.videoMessage) return "video/mp4";
return "";
}
function mediaFileNameFromMessage(m: proto.IMessage | null | undefined): string {
const u = unwrapMessage(m);
if (!u) return "";
const docName = String(u.documentMessage?.fileName || u.documentMessage?.title || "").trim();
if (docName) return docName;
if (u.audioMessage) return isVoiceNote(m) ? "voice.ogg" : "audio.mp3";
if (u.imageMessage) return "image.jpg";
if (u.videoMessage) return "video.mp4";
if (u.stickerMessage) return "sticker.webp";
return "";
}
function sanitizeFileName(name: string): string {
const base = String(name || "").trim().replace(/[^\w.\-()+\u4e00-\u9fff]/g, "_");
return base.slice(0, 120) || "attachment.bin";
}
function buildSavedFileName(params: {
msgId: string;
mime: string;
originalName: string;
ext: string;
}): string {
const original = sanitizeFileName(params.originalName);
if (original.includes(".")) return `${Date.now()}-${original}`;
const ext = params.ext.startsWith(".") ? params.ext : `.${params.ext || "bin"}`;
return `${Date.now()}-${params.msgId}${ext}`;
}
function resolveMediaKind(contentType: string | undefined, mime: string, m?: proto.IMessage | null): string | null {
const t = String(contentType || "").trim();
const mimeLower = String(mime || "").trim().toLowerCase();
if (t === "imageMessage" || t === "stickerMessage") return "image";
if (t === "videoMessage") return "video";
if (t === "audioMessage") return isVoiceNote(m || null) ? "voice" : "file";
if (t === "documentMessage") {
if (mimeLower.startsWith("image/")) return "image";
if (mimeLower.startsWith("video/")) return "video";
if (mimeLower.startsWith("audio/")) return "file";
return "file";
}
return null;
}
function hasInboundMedia(m: proto.IMessage | null | undefined): boolean {
const u = unwrapMessage(m);
if (!u) return false;
return Boolean(u.imageMessage || u.videoMessage || u.documentMessage || u.audioMessage || u.stickerMessage);
}
async function withTimeout<T>(promise: Promise<T>, ms: number, label: string): Promise<T> {
let timer: ReturnType<typeof setTimeout> | undefined;
const timeout = new Promise<never>((_, reject) => {
timer = setTimeout(() => reject(new Error(`${label} timeout after ${ms}ms`)), ms);
});
try {
return await Promise.race([promise, timeout]);
} finally {
if (timer) clearTimeout(timer);
}
}
async function downloadInboundAttachments(params: {
sock: ReturnType<typeof makeWASocket>;
msg: proto.IWebMessageInfo;
}): Promise<Json[]> {
const inner = unwrapMessage(params.msg.message);
if (!inner || !hasInboundMedia(params.msg.message)) return [];
const contentType = getContentType(inner);
const mime = mediaMimeFromMessage(params.msg.message) || "application/octet-stream";
const kind = resolveMediaKind(contentType, mime, params.msg.message);
if (!kind) return [];
const msgId = String(params.msg.key?.id || "media").replace(/[^\w.-]/g, "_").slice(0, 32);
const originalName = mediaFileNameFromMessage(params.msg.message);
try {
if (VERBOSE) log(`media download start id=${msgId} kind=${kind} mime=${mime} name=${originalName || "-"}`);
const buffer = await withTimeout(
downloadMediaMessage(
params.msg,
"buffer",
{},
{ reuploadRequest: params.sock.updateMediaMessage },
),
MEDIA_DOWNLOAD_TIMEOUT_MS,
"whatsapp media download",
);
if (!buffer?.length) return [];
if (MEDIA_MAX_BYTES > 0 && buffer.length > MEDIA_MAX_BYTES) {
log(`media skip id=${msgId}: ${buffer.length} bytes exceeds max ${MEDIA_MAX_BYTES}`);
return [];
}
const ext = extensionForMediaMessage(inner) || inferExtension(mime, originalName);
const dir = path.join(STATE_DIR, "media", "inbound");
ensureDir(dir);
const savedName = buildSavedFileName({ msgId, mime, originalName, ext });
const filePath = path.join(dir, savedName);
await fs.promises.writeFile(filePath, buffer);
if (VERBOSE) log(`media download done id=${msgId} bytes=${buffer.length} path=${filePath}`);
return [
{
kind,
local_path: filePath,
mime,
media_type: mime,
name: originalName || savedName,
filename: originalName || savedName,
},
];
} catch (err) {
log(`media download failed id=${msgId} err=${String(err)}`);
return [];
}
}
function buildInboundPayload(params: {
chatId: string;
userId: string;
text: string;
raw: unknown;
isGroup: boolean;
mentions: string[];
attachments?: Json[];
groupName?: string;
botJid?: string;
botLid?: string;
mentionsBot?: boolean;
}): Json {
const metadata: Json = {
source: "whatsapp_baileys",
raw: params.raw,
mentions_bot: params.mentionsBot === true,
};
if (params.groupName) metadata.group_name = params.groupName;
if (params.botJid) metadata.bot_jid = params.botJid;
if (params.botLid) metadata.bot_lid = params.botLid;
if (params.mentions.length) metadata.mentioned_jids = params.mentions;
const payload: Json = {
channel: "whatsapp",
account_id: ACCOUNT_ID,
user_id: params.userId,
chat_id: params.chatId,
text: params.text,
is_group: params.isGroup,
mentions: params.mentions,
metadata,
};
if (params.attachments?.length) payload.attachments = params.attachments;
return payload;
}
function decodeBase64Payload(raw: string): { mime: string; data: Buffer } | null {
const s = String(raw || "").trim();
if (!s) return null;
let mime = "application/octet-stream";
let payload = s;
const m = s.match(/^data:([^;,]+);base64,(.*)$/i);
if (m) {
mime = String(m[1] || mime).trim() || mime;
payload = String(m[2] || "").trim();
}
try {
const data = Buffer.from(payload.replace(/\s+/g, ""), "base64");
if (!data.length) return null;
return { mime, data };
} catch {
return null;
}
}
function readReplyMetadata(reply: Json): Json {
const meta = (reply as any).metadata;
return meta && typeof meta === "object" ? (meta as Json) : {};
}
function readMentionJids(meta: Json): string[] {
const raw = (meta as any).mention_jids ?? (meta as any).mentionJids;
if (!Array.isArray(raw)) {
const single = String((meta as any).reply_to_user_id || (meta as any).replyToUserId || "").trim();
return single ? [single] : [];
}
return raw.map((j) => String(j || "").trim()).filter(Boolean);
}
function readMentionNames(meta: Json): string[] {
const raw = (meta as any).mention_names ?? (meta as any).mentionNames;
if (!Array.isArray(raw)) return [];
return raw.map((n) => String(n || "").trim()).filter(Boolean);
}
function escapeRegExp(text: string): string {
return String(text || "").replace(/[.*+?^${}()|[\]\\]/g, "\\$&");
}
function mentionTokenForJid(jid: string): string {
const local = jidBaseLocal(jid);
return local ? `@${local}` : "";
}
function resolveOutboundMentionJids(
sock: ReturnType<typeof makeWASocket> | null,
chatId: string,
mentionJids: string[],
): string[] {
const isGroup = String(chatId || "").toLowerCase().endsWith("@g.us");
const lidMapping = (sock as any)?.signalRepository?.lidMapping;
const out: string[] = [];
const seen = new Set<string>();
const push = (jid: string) => {
const j = String(jid || "").trim();
if (!j) return;
const key = jidBaseLocal(j);
if (!key || seen.has(key)) return;
seen.add(key);
out.push(j);
};
for (let raw of mentionJids) {
let jid = String(raw || "").trim();
if (!jid) continue;
if (!jid.includes("@") && /^\d+$/.test(jid)) {
jid = `${jid}@lid`;
}
const low = jid.toLowerCase();
if (low.endsWith("@lid")) {
push(jid);
continue;
}
if (isGroup && lidMapping && typeof lidMapping.getLIDForPN === "function") {
try {
const pn = low.includes("@s.whatsapp") ? jidNormalizedUser(jid) : jid;
const lid = lidMapping.getLIDForPN(pn);
if (lid) {
push(String(lid));
continue;
}
} catch {
// ignore
}
}
push(jid.includes("@") ? jid : jidNormalizedUser(jid));
}
return out;
}
function alignMentionTextWithJids(text: string, mentionJids: string[], mentionNames: string[]): string {
let body = String(text || "").trim();
const tokens = mentionJids.map((j) => mentionTokenForJid(j)).filter(Boolean);
if (!tokens.length) return body;
for (let i = 0; i < mentionJids.length; i++) {
const token = tokens[i] || "";
if (!token) continue;
const name = String(mentionNames[i] || "").trim();
if (name) {
body = body.replace(new RegExp(`@${escapeRegExp(name)}(?=\\s|$|[,。!?!?,.])`, "gu"), token);
}
}
const missing = tokens.filter((token) => !body.includes(token));
if (missing.length) {
body = body.replace(/^(@\S+\s*)+/, "").trim();
body = `${missing.join(" ")} ${body}`.trim();
}
return body;
}
function decodeOutboundSource(raw: string): Json {
const text = String(raw || "").trim();
if (!text) return {};
if (text.startsWith("{")) {
try {
const data = JSON.parse(text);
return data && typeof data === "object" ? (data as Json) : {};
} catch {
return {};
}
}
return { kind: text };
}
function isOutboundMentionTextReady(meta: Json): boolean {
return Boolean((meta as any).mention_text_ready ?? (meta as any).mentionTextReady);
}
function buildMentionedOutboundText(
body: string,
meta: Json,
sock?: ReturnType<typeof makeWASocket> | null,
chatId?: string,
): { text: string; mentions?: string[] } {
const trimmed = String(body || "").trim();
if (!trimmed) return { text: trimmed };
let mentionJids = readMentionJids(meta);
if (!mentionJids.length) return { text: trimmed };
const chat = String(chatId || "").trim();
if (sock && chat) {
mentionJids = resolveOutboundMentionJids(sock, chat, mentionJids);
}
const mentionNames = readMentionNames(meta);
if (isOutboundMentionTextReady(meta)) {
const aligned = alignMentionTextWithJids(trimmed, mentionJids, mentionNames);
return { text: aligned, mentions: mentionJids };
}
const cleaned = trimmed
.replace(/@\+?\d+(?:\s+\d+)*\s*/g, "")
.replace(/@\S+\s*/g, "")
.trim();
const prefix = mentionJids.map((j) => mentionTokenForJid(j)).filter(Boolean).join(" ");
const outText = prefix ? `${prefix} ${cleaned}`.trim() : trimmed;
return { text: outText, mentions: mentionJids };
}
function buildOutboundSendContent(
text: string,
sourceRaw: string,
sock?: ReturnType<typeof makeWASocket> | null,
chatId?: string,
): { text: string; mentions?: string[] } {
const meta = decodeOutboundSource(sourceRaw);
return buildMentionedOutboundText(text, meta, sock, chatId);
}
function buildQuotedMessage(params: {
chatId: string;
stanzaId: string;
participant: string;
quoteText: string;
}): proto.IWebMessageInfo | undefined {
const stanzaId = String(params.stanzaId || "").trim();
const chatId = String(params.chatId || "").trim();
if (!stanzaId || !chatId) return undefined;
const participant = String(params.participant || "").trim();
const quoteText = String(params.quoteText || "").trim() || "...";
return {
key: {
remoteJid: chatId,
fromMe: false,
id: stanzaId,
...(participant ? { participant } : {}),
},
message: {
conversation: quoteText,
},
} as proto.IWebMessageInfo;
}
function buildTextSendOptions(params: {
deliverTo: string;
text: string;
reply: Json;
sock?: ReturnType<typeof makeWASocket> | null;
}): { content: { text: string; mentions?: string[] }; quoted?: proto.IWebMessageInfo } {
const meta = readReplyMetadata(params.reply);
const content = buildMentionedOutboundText(params.text, meta, params.sock, params.deliverTo);
const quoted = buildQuotedMessage({
chatId: String((meta as any).quote_remote_jid || (meta as any).quoteRemoteJid || params.deliverTo).trim(),
stanzaId: String((meta as any).quote_stanza_id || (meta as any).quoteStanzaId || "").trim(),
participant: String((meta as any).quote_participant || (meta as any).quoteParticipant || "").trim(),
quoteText: String((meta as any).quote_text || (meta as any).quoteText || "").trim(),
});
return quoted ? { content, quoted } : { content };
}
function shouldSendOutboundText(text: string): boolean {
const t = String(text || "").trim();
if (!t) return false;
const normalized = t.toLowerCase().replace(/(/g, "(").replace(/)/g, ")");
if (normalized === "(silent)" || normalized === "[silent]" || normalized === "no_reply" || normalized === "no reply" || normalized === "静默") {
return false;
}
return true;
}
function guessMimeFromPath(filePath: string): string {
const ext = path.extname(filePath).toLowerCase();
const map: Record<string, string> = {
".jpg": "image/jpeg",
".jpeg": "image/jpeg",
".png": "image/png",
".gif": "image/gif",
".webp": "image/webp",
".mp4": "video/mp4",
".mov": "video/quicktime",
".webm": "video/webm",
".mp3": "audio/mpeg",
".ogg": "audio/ogg",
".wav": "audio/wav",
".m4a": "audio/mp4",
".pdf": "application/pdf",
".txt": "text/plain",
".csv": "text/csv",
".json": "application/json",
".doc": "application/msword",
".docx": "application/vnd.openxmlformats-officedocument.wordprocessingml.document",
".xls": "application/vnd.ms-excel",
".xlsx": "application/vnd.openxmlformats-officedocument.spreadsheetml.sheet",
".ppt": "application/vnd.ms-powerpoint",
".pptx": "application/vnd.openxmlformats-officedocument.presentationml.presentation",
".zip": "application/zip",
};
return map[ext] || "application/octet-stream";
}
function buildOutboundMediaMessage(params: {
data: Buffer;
mime: string;
fileName: string;
caption?: string;
mentions?: string[];
}): Json {
const mime = String(params.mime || "application/octet-stream").trim() || "application/octet-stream";
const m = mime.toLowerCase();
const caption = String(params.caption || "").trim();
const captionOpt = caption ? { caption } : {};
const mentionOpt = params.mentions?.length ? { mentions: params.mentions } : {};
if (m.startsWith("image/")) {
return { image: params.data as any, mimetype: mime, ...captionOpt, ...mentionOpt };
}
if (m.startsWith("video/")) {
return { video: params.data as any, mimetype: mime, ...captionOpt, ...mentionOpt };
}
if (m.startsWith("audio/")) {
return { audio: params.data as any, mimetype: mime, ...captionOpt, ...mentionOpt };
}
return {
document: params.data as any,
mimetype: mime,
fileName: sanitizeFileName(params.fileName || "attachment.bin"),
...captionOpt,
...mentionOpt,
};
}
async function sendReplyWithAttachments(params: {
sock: ReturnType<typeof makeWASocket> | null;
deliverTo: string;
text: string;
reply: Json;
}): Promise<void> {
const s = params.sock;
if (!s) return;
const reply = params.reply || {};
const outText = String(params.text || "").trim();
const mediaPath = String((reply as any).media_path || (reply as any).mediaPath || "").trim();
const mediaUrl = String((reply as any).media_url || (reply as any).mediaUrl || "").trim();
const attachments = Array.isArray((reply as any).attachments) ? ((reply as any).attachments as Json[]) : [];
const textOpts = buildTextSendOptions({
deliverTo: params.deliverTo,
text: outText,
reply,
sock: params.sock,
});
const sendOpts = textOpts.quoted ? { quoted: textOpts.quoted } : undefined;
const sendMediaBuffer = async (data: Buffer, mime: string, fileName: string): Promise<boolean> => {
if (!data.length) return false;
const msg = buildOutboundMediaMessage({
data,
mime,
fileName,
caption: textOpts.content.text,
mentions: textOpts.content.mentions,
});
await s.sendMessage(params.deliverTo, msg as any, sendOpts);
return true;
};
const sendMediaRef = async (source: string): Promise<boolean> => {
const src = String(source || "").trim();
if (!src) return false;
try {
const data = await fs.promises.readFile(src);
const mime = guessMimeFromPath(src);
const fileName = path.basename(src);
return await sendMediaBuffer(data, mime, fileName);
} catch (err) {
log(`send media file failed path=${src} err=${String(err)}`);
return false;
}
};
if (await sendMediaRef(mediaPath || mediaUrl)) return;
for (const att of attachments) {
if (!att || typeof att !== "object") continue;
const raw = String(
(att as any).data_base64 ||
(att as any).media_base64 ||
(att as any).image_base64 ||
(att as any).video_base64 ||
(att as any).audio_base64 ||
(att as any).data ||
"",
).trim();
if (!raw) continue;
const decoded = decodeBase64Payload(raw);
if (!decoded) continue;
const fileName = String((att as any).name || (att as any).filename || "attachment.bin");
const mime = String((att as any).mime || (att as any).media_type || (att as any).mime_type || decoded.mime).trim() || decoded.mime;
if (await sendMediaBuffer(decoded.data, mime, fileName)) return;
}
if (outText) {
try {
await s.sendMessage(
params.deliverTo,
textOpts.content as any,
textOpts.quoted ? { quoted: textOpts.quoted } : undefined,
);
} catch (err) {
log(`send reply failed (${String(err)}); retry plain text`);
await s.sendMessage(params.deliverTo, { text: outText });
}
}
}
async function postInbound(payload: Json): Promise<Json> {
const url = `${LOCAL_BASE_URL.replace(/\/+$/, "")}/inbound/whatsapp`;
const res = await fetch(url, {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify(payload),
});
const text = await res.text();
if (!res.ok) {
throw new Error(`inbound ${res.status}: ${text.slice(0, 300)}`);
}
return text ? (JSON.parse(text) as Json) : {};
}
async function pollOutboundQueue(sock: ReturnType<typeof makeWASocket>): Promise<void> {
const url = `${LOCAL_BASE_URL.replace(/\/+$/, "")}/whatsapp/outbound/pending?account_id=${encodeURIComponent(ACCOUNT_ID)}`;
try {
const res = await fetch(url);
if (!res.ok) return;
const body = (await res.json()) as Json;
const items = Array.isArray(body.items) ? (body.items as Json[]) : [];
for (const item of items) {
if (!item || typeof item !== "object") continue;
const id = String((item as any).id || "").trim();
const chatId = String((item as any).chat_id || "").trim();
const text = String((item as any).text || "").trim();
const source = String((item as any).source || "").trim();
const attachments = Array.isArray((item as any).attachments) ? ((item as any).attachments as Json[]) : [];
const mediaPath = String((item as any).media_path || (item as any).mediaPath || "").trim();
const hasMedia = attachments.length > 0 || Boolean(mediaPath);
if (!id || !chatId || (!text && !hasMedia)) continue;
let ok = true;
let err = "";
try {
let stanzaId = "";
if (hasMedia) {
const reply: Json = {
attachments,
metadata: decodeOutboundSource(source),
...(mediaPath ? { media_path: mediaPath } : {}),
};
await sendReplyWithAttachments({ sock, deliverTo: chatId, text, reply });
log(
`outbound sent id=${id} chat=${chatId} attachments=${attachments.length} mediaPath=${mediaPath ? "yes" : "no"} text=${text.slice(0, 80)}`,
);
} else {
const meta = decodeOutboundSource(source);
const content = buildMentionedOutboundText(text, meta, sock, chatId);
const quoted = buildQuotedMessage({
chatId: String((meta as any).quote_remote_jid || (meta as any).quoteRemoteJid || chatId).trim(),
stanzaId: String((meta as any).quote_stanza_id || (meta as any).quoteStanzaId || "").trim(),
participant: String((meta as any).quote_participant || (meta as any).quoteParticipant || "").trim(),
quoteText: String((meta as any).quote_text || (meta as any).quoteText || "").trim(),
});
const sent = await sock.sendMessage(
chatId,
content as any,
quoted ? { quoted } : undefined,
);
stanzaId = String((sent as any)?.key?.id || "").trim();
log(
`outbound sent id=${id} chat=${chatId} stanza=${stanzaId || "?"} mentions=${(content.mentions || []).length} mention0=${(content.mentions || [])[0] || ""} text=${content.text.slice(0, 80)}`,
);
}
try {
await fetch(`${LOCAL_BASE_URL.replace(/\/+$/, "")}/whatsapp/outbound/ack`, {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify({ id, ok, error: err, stanza_id: stanzaId }),
});
} catch {
// best-effort ack
}
} catch (sendErr) {
ok = false;
err = String(sendErr);
log(`outbound send failed id=${id} chat=${chatId} err=${err.slice(0, 160)}`);
try {
await fetch(`${LOCAL_BASE_URL.replace(/\/+$/, "")}/whatsapp/outbound/ack`, {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify({ id, ok, error: err }),
});
} catch {
// best-effort ack
}
}
}
} catch (err) {
if (VERBOSE) log(`outbound poll error: ${String(err)}`);
}
}
function startOutboundPoller(getSock: () => ReturnType<typeof makeWASocket> | null): void {
const intervalMs = Number(process.env.OCLAW_WHATSAPP_OUTBOUND_POLL_MS || "1000") || 1000;
// Overlapping polls used to claim the same pending row twice (send → duplicate WhatsApp
// bubbles) once inbound replies also went through the outbound queue.
let pollInFlight = false;
setInterval(() => {
if (pollInFlight) return;
const s = getSock();
if (!s) return;
pollInFlight = true;
void pollOutboundQueue(s).finally(() => {
pollInFlight = false;
});
}, Math.max(1000, intervalMs));
}
function sleep(ms: number): Promise<void> {
return new Promise((r) => setTimeout(r, ms));
}
type TypingEntry = { refs: number; timer: ReturnType<typeof setInterval> | null };
const typingByChat = new Map<string, TypingEntry>();
async function sendTypingPresence(
sock: ReturnType<typeof makeWASocket> | null,
chatId: string,
state: "composing" | "paused",
): Promise<void> {
const s = sock;
const jid = String(chatId || "").trim();
if (!s || !jid || !TYPING_ENABLED) return;
try {
await s.sendPresenceUpdate(state, jid);
} catch (err) {
if (VERBOSE) log(`presence ${state} failed chat=${jid}: ${String(err).slice(0, 120)}`);
}
}
function startTypingHeartbeat(sock: ReturnType<typeof makeWASocket> | null, chatId: string): void {
const jid = String(chatId || "").trim();
if (!sock || !jid || !TYPING_ENABLED) return;
let entry = typingByChat.get(jid);
if (!entry) {
entry = { refs: 0, timer: null };
typingByChat.set(jid, entry);
}
entry.refs += 1;
if (entry.refs === 1) {
void sendTypingPresence(sock, jid, "composing");
entry.timer = setInterval(() => {
void sendTypingPresence(sock, jid, "composing");
}, TYPING_HEARTBEAT_MS);
}
}
function stopTypingHeartbeat(sock: ReturnType<typeof makeWASocket> | null, chatId: string): void {
const jid = String(chatId || "").trim();
if (!jid || !TYPING_ENABLED) return;
const entry = typingByChat.get(jid);
if (!entry) return;
entry.refs = Math.max(0, entry.refs - 1);
if (entry.refs > 0) return;
if (entry.timer) {
clearInterval(entry.timer);
entry.timer = null;
}
typingByChat.delete(jid);
void sendTypingPresence(sock, jid, "paused");
}
function _ipLooksHijackedOrUnroutable(ip: string): boolean {
const s = String(ip || "").trim();
if (!s) return false;
if (s.startsWith("198.18.") || s.startsWith("198.19.")) return true; // reserved benchmark range
if (s.startsWith("0.") || s.startsWith("127.") || s.startsWith("10.")) return true;
if (s.startsWith("192.168.") || s.startsWith("169.254.")) return true;
if (/^172\.(1[6-9]|2\d|3[0-1])\./.test(s)) return true;
return false;
}
const STALE_WA_BUILD = 1027934701;
const FALLBACK_WA_VERSION: [number, number, number] = [2, 3000, 1043857760];
function isUsableWaVersion(version: unknown): version is [number, number, number] {
if (!Array.isArray(version) || version.length < 3) return false;
const build = Number(version[2]);
return Number.isFinite(build) && build > STALE_WA_BUILD;
}
async function proxyFetchInit(): Promise<Record<string, unknown>> {
if (!PROXY_URL) return {};
try {
const undici = await import("node:undici");
return { dispatcher: new undici.ProxyAgent(PROXY_URL) };
} catch {
return {};
}
}
async function resolveWhatsAppVersion(): Promise<[number, number, number]> {
const fromEnv = String(process.env.AIA_WHATSAPP_WA_VERSION || "").trim();
if (fromEnv) {
const parts = fromEnv.split(/[.,]/).map((n) => Number(n));
if (isUsableWaVersion(parts)) {
log(`wa version from AIA_WHATSAPP_WA_VERSION=${parts.join(".")}`);
return [parts[0], parts[1], parts[2]];
}
log(`ignore invalid AIA_WHATSAPP_WA_VERSION=${fromEnv}`);
}
const fetchInit = await proxyFetchInit();
try {
const latest = await withTimeout(fetchLatestBaileysVersion(fetchInit as any), 8000, "fetchLatestBaileysVersion");
if (isUsableWaVersion(latest.version)) {
log(`wa version from Baileys master=${latest.version.join(".")} isLatest=${String(latest.isLatest)}`);
return [latest.version[0], latest.version[1], latest.version[2]];
}
log(
`fetchLatestBaileysVersion unusable version=${JSON.stringify(latest.version)} error=${String((latest as any).error || "")}`,
);
} catch (err) {
log(`fetchLatestBaileysVersion failed: ${String(err)}`);
}
try {
const web = await withTimeout(fetchLatestWaWebVersion(fetchInit as any), 8000, "fetchLatestWaWebVersion");
if (isUsableWaVersion(web.version)) {
log(`wa version from WhatsApp Web=${web.version.join(".")}`);
return [web.version[0], web.version[1], web.version[2]];
}
log(`fetchLatestWaWebVersion unusable version=${JSON.stringify(web.version)}`);
} catch (err) {
log(`fetchLatestWaWebVersion failed: ${String(err)}`);
}
log(`wa version fallback=${FALLBACK_WA_VERSION.join(".")} (rc.9 default ${STALE_WA_BUILD} is rejected with 405)`);
return FALLBACK_WA_VERSION;
}
async function logNetworkHints(): Promise<void> {
try {
const a = await dns.resolve4("web.whatsapp.com").catch(() => []);
const aaaa = await dns.resolve6("web.whatsapp.com").catch(() => []);
const sample = [...a, ...aaaa].slice(0, 6);
if (sample.length > 0) {
log(`dns web.whatsapp.com => ${sample.join(", ")}`);
const bad = sample.find(_ipLooksHijackedOrUnroutable);
if (bad) {
log(
`WARNING: DNS looks suspicious (e.g. ${bad}). This often causes WhatsApp Web handshake failures. Try switching system DNS (1.1.1.1/8.8.8.8) or check proxy/hosts rules.`,
);
}
}
} catch {
// best-effort
}
const httpProxy = process.env.HTTP_PROXY || process.env.http_proxy || "";
const httpsProxy = process.env.HTTPS_PROXY || process.env.https_proxy || "";
if (httpProxy || httpsProxy) {
log(`proxy env detected HTTP_PROXY=${httpProxy ? "set" : "unset"} HTTPS_PROXY=${httpsProxy ? "set" : "unset"}`);
}
if (PROXY_URL) {
log(`proxy configured for Baileys via PROXY_URL=${PROXY_URL}`);
try {
fs.writeFileSync(path.join(STATE_DIR, "proxy.url"), PROXY_URL, "utf8");
} catch {
// ignore
}
}
}
function makeDeduper(params: { ttlMs: number; max: number }) {
const seen = new Map<string, number>();
const prune = (now: number): void => {
for (const [k, ts] of seen) {
if (now - ts > params.ttlMs) {
seen.delete(k);
}
}
if (seen.size <= params.max) return;
const ordered = Array.from(seen.entries()).sort((a, b) => a[1] - b[1]);
const drop = ordered.slice(0, Math.max(0, ordered.length - params.max));
for (const [k] of drop) seen.delete(k);
};
return {
has: (id: string): boolean => {
const now = Date.now();
prune(now);
const ts = seen.get(id);
return typeof ts === "number" && now - ts <= params.ttlMs;
},
add: (id: string): void => {
const now = Date.now();
seen.set(id, now);
prune(now);
},
};
}
async function main(): Promise<void> {
ensureDir(STATE_DIR);
log(
`runner booting local=${LOCAL_BASE_URL} stateDir=${STATE_DIR} verbose=${VERBOSE} node=${process.version} proxy=${PROXY_URL ? "set" : "unset"}`,
);
writeBridgeStatus(STATE_DIR, {
connection: "connecting",
me: "",
phone: "",
qr: "",
qr_data_url: "",
qr_png: "",
last_error: "runner_starting",
reconnect_attempt: 0,
login_only: LOGIN_ONLY,
});
let { state, saveCreds } = await loadAuthState(STATE_DIR);
const version = await resolveWhatsAppVersion();
const dedupe = makeDeduper({ ttlMs: 10 * 60 * 1000, max: 20_000 });
let reconnectAttempt = 0;
let sock: ReturnType<typeof makeWASocket> | null = null;
let cachedBotIdentityJids: string[] = [];
const wsAgent = PROXY_URL ? new HttpsProxyAgent(PROXY_URL) : undefined;
const groupNameCache = new Map<string, { name: string; ts: number }>();
const GROUP_NAME_TTL_MS = 10 * 60 * 1000;
const initialMe = String((state.creds as any)?.me?.id || "").trim();
writeBridgeStatus(STATE_DIR, {
connection: "connecting",
me: initialMe ? jidNormalizedUser(initialMe) : "",
phone: jidPhone(initialMe),
qr: "",
qr_data_url: "",
qr_png: "",
last_error: "",
reconnect_attempt: 0,
login_only: LOGIN_ONLY,
});
const resolveGroupName = async (chatId: string): Promise<string> => {
const key = String(chatId || "").trim();
if (!key) return "";
const cached = groupNameCache.get(key);
const now = Date.now();
if (cached && now - cached.ts < GROUP_NAME_TTL_MS) return cached.name;
const s = sock;
if (!s) return cached?.name || "";
try {
const meta = await s.groupMetadata(key);
const name = String(meta?.subject || "").trim();
groupNameCache.set(key, { name, ts: now });
return name;
} catch {
return cached?.name || "";
}
};
const connectOnce = async () => {
sock = makeWASocket({
auth: state,
browser: Browsers.windows("oclaw"),
version,
agent: wsAgent as any,
printQRInTerminal: false,
generateHighQualityLinkPreview: false,
});
sock.ev.on("creds.update", async () => {
await saveCreds();
const meId = sock?.user?.id ? String(sock.user.id) : "";
const ids = collectBotIdentityJids(sock, meId, state.creds);
if (ids.length) cachedBotIdentityJids = ids;
});
sock.ev.on("connection.update", async (update) => {
if (update.qr) {
reconnectAttempt = 0;
log("QR received. Scan it from WhatsApp -> Linked devices.");
printQrToTerminal(update.qr);
const qrArt = await writeQrImage(STATE_DIR, update.qr);
writeBridgeStatus(STATE_DIR, {
connection: "qr",
qr: String(update.qr || ""),
qr_data_url: qrArt.dataUrl,
qr_png: qrArt.pngPath ? path.basename(qrArt.pngPath) : "",
last_error: "",
last_disconnect_reason: "",
last_disconnect_status: null,
reconnect_attempt: reconnectAttempt,
login_only: LOGIN_ONLY,
});
}
if (update.connection === "open") {
reconnectAttempt = 0;
const me = sock?.user?.id ? jidNormalizedUser(sock.user.id) : "";
const botIds = collectBotIdentityJids(sock, sock?.user?.id ? String(sock.user.id) : "", state.creds);
cachedBotIdentityJids = botIds;
log(`connected. me=${me || "unknown"} botIds=${botIds.join(",") || "none"} loginOnly=${LOGIN_ONLY}`);
writeBridgeStatus(STATE_DIR, {
connection: "open",
me,
phone: jidPhone(me),
qr: "",
qr_data_url: "",
qr_png: "",
last_error: "",
last_disconnect_reason: "",
last_disconnect_status: null,
reconnect_attempt: 0,
login_only: LOGIN_ONLY,
});
try {
const qrPng = path.join(STATE_DIR, "qr.png");
if (fs.existsSync(qrPng)) fs.unlinkSync(qrPng);
} catch {
// ignore
}
if (LOGIN_ONLY) {
log("login-only mode: exiting after successful link.");
process.exit(0);
}
}
if (update.connection === "close") {
const statusCode = (update.lastDisconnect?.error as any)?.output?.statusCode as number | undefined;
const reason = statusCode ? (DisconnectReason as any)[statusCode] || String(statusCode) : "unknown";
const errText = String(update.lastDisconnect?.error || "").slice(0, 200);
log(`disconnected reason=${reason} statusCode=${String(statusCode || "")} err=${errText}`);
const hasLinkedSession = Boolean(
jidPhone(String((state.creds as any)?.me?.id || sock?.user?.id || "")),
);
if (statusCode === DisconnectReason.loggedOut) {
log("logged out: delete data/channel_sidecar/whatsapp/state/auth to re-link.");
writeBridgeStatus(STATE_DIR, {
connection: "logged_out",
qr: "",
qr_data_url: "",
qr_png: "",
last_disconnect_reason: String(reason),
last_disconnect_status: typeof statusCode === "number" ? statusCode : null,
last_error: errText,
reconnect_attempt: reconnectAttempt,
login_only: LOGIN_ONLY,
});
return;
}
const rebindCodes = new Set<number>([
DisconnectReason.forbidden,
DisconnectReason.badSession,
405,
]);
if (hasLinkedSession && typeof statusCode === "number" && rebindCodes.has(statusCode)) {
log(`session invalid statusCode=${statusCode}: clear state/auth and scan QR again.`);
writeBridgeStatus(STATE_DIR, {
connection: "needs_rebind",
qr: "",
qr_data_url: "",
qr_png: "",
last_disconnect_reason: String(reason),
last_disconnect_status: statusCode,
last_error: errText,
reconnect_attempt: reconnectAttempt,
login_only: LOGIN_ONLY,
});
return;
}
if (!hasLinkedSession && statusCode === 405) {
reconnectAttempt += 1;
log(`pairing 405: WhatsApp rejected this identity. Resetting auth keys attempt=${reconnectAttempt}`);
writeBridgeStatus(STATE_DIR, {
connection: "connecting",
me: "",
phone: "",
qr: "",
qr_data_url: "",
qr_png: "",
last_disconnect_reason: "405",
last_disconnect_status: 405,
last_error: "pairing_rejected_resetting_keys",
reconnect_attempt: reconnectAttempt,
login_only: LOGIN_ONLY,
});
try {
sock?.end(undefined as any);
} catch {
// ignore
}
sock = null;
try {
fs.rmSync(resolveAuthDir(STATE_DIR), { recursive: true, force: true });
} catch {
// ignore
}
const fresh = await loadAuthState(STATE_DIR);
state = fresh.state;
saveCreds = fresh.saveCreds;
await sleep(2000);
await connectOnce();
return;
}
reconnectAttempt += 1;
writeBridgeStatus(STATE_DIR, {
connection: hasLinkedSession ? "close" : "connecting",
qr: "",
qr_data_url: "",
qr_png: "",
last_disconnect_reason: String(reason),
last_disconnect_status: typeof statusCode === "number" ? statusCode : null,
last_error: errText,
reconnect_attempt: reconnectAttempt,
login_only: LOGIN_ONLY,
});
const base = hasLinkedSession
? Math.min(30_000, 1000 * Math.pow(2, Math.min(6, reconnectAttempt)))
: 1500;
const jitter = Math.floor(Math.random() * 500);
const delay = base + jitter;
log(`reconnecting in ${delay}ms attempt=${reconnectAttempt}`);
await sleep(delay);
await connectOnce();
}
});
sock.ev.on("messages.upsert", async (upsert) => {
const msgs = Array.isArray(upsert.messages) ? upsert.messages : [];
for (const msg of msgs) {
try {
const key = msg.key;
const id = String(key.id || "").trim();
const remoteJid = String(key.remoteJid || "").trim();
if (!id || !remoteJid) continue;
if (isStatusOrBroadcastJid(remoteJid)) continue;
if (key.fromMe) continue;
if (dedupe.has(id)) continue;
dedupe.add(id);
const s = sock;
if (!s) continue;
const attachments = await downloadInboundAttachments({ sock: s, msg });
const text = pickText(msg.message);
if (!text && !attachments.length) continue;
const isGroup = remoteJid.endsWith("@g.us");
const participantRaw = String(key.participant || "").trim();
const participantAlt = String((key as any).participantAlt || "").trim();
const remoteJidAlt = String((key as any).remoteJidAlt || "").trim();
const remoteJidNorm = jidNormalizedUser(remoteJid);
const userId = isGroup
? resolveSenderJid(key) || remoteJidNorm
: resolveDmUserJid(key, sock);
const chatId = remoteJidNorm;
const mentions = extractMentionsFromUpsert(msg);
const quote = extractQuoteContext(msg.message);
const botJidRaw = sock?.user?.id ? String(sock.user.id).trim() : "";
const botJid = botJidRaw ? jidNormalizedUser(botJidRaw) : "";
const botIdentityJids = Array.from(
new Set([...cachedBotIdentityJids, ...collectBotIdentityJids(sock, botJidRaw, state.creds)]),
);
if (botIdentityJids.length > cachedBotIdentityJids.length) {
cachedBotIdentityJids = botIdentityJids;
}
const botLid = pickBotLid(botIdentityJids);
const mentionsBot = messageMentionsBot(sock, mentions, botJidRaw, botIdentityJids);
const isReplyToBot = Boolean(
botJidRaw &&
quote.participant &&
botIdentityJids.some((botId) => mentionMatchesBotIdentity(quote.participant, botId)),
);
const groupName = isGroup ? await resolveGroupName(chatId) : "";
const raw = {
id,
remoteJid,
remoteJidAlt: remoteJidAlt || null,
participant: participantRaw || null,
participantAlt: participantAlt || null,
pushName: (msg as any).pushName || null,
messageTimestamp: (msg as any).messageTimestamp || null,
quotedParticipant: quote.participant || null,
quotedStanzaId: quote.stanzaId || null,
quotedText: quote.quotedText || null,
isReplyToBot,
mentionsBot,
mentionedJids: mentions,
botLid: botLid || null,
};
const inbound = buildInboundPayload({
chatId,
userId,
text,
raw,
isGroup,
mentions,
attachments,
groupName: groupName || undefined,
botJid: botJidRaw || botJid || undefined,
botLid: botLid || undefined,
mentionsBot,
});
if (VERBOSE || isGroup) {
log(
`inbound group=${isGroup} chat=${chatId} user=${userId} mentions=${mentions.length} mentionsBot=${mentionsBot} replyToBot=${isReplyToBot} textLen=${text.length} attachments=${attachments.length} mention0=${mentions[0] || ""} botLid=${botLid || ""}`,
);
} else if (VERBOSE || attachments.length) {
log(`inbound posting chat=${chatId} user=${userId} textLen=${text.length} attachments=${attachments.length}`);
}
startTypingHeartbeat(s, chatId);
let out: Json = {};
try {
out = await postInbound(inbound);
} finally {
stopTypingHeartbeat(s, chatId);
}
const replies = Array.isArray(out.replies) ? (out.replies as Json[]) : [];
const delivery = String((out as any).delivery || "").trim().toLowerCase();
if (isGroup || VERBOSE || delivery === "queued" || delivery === "accepted_queued") {
log(
`inbound ok chat=${chatId} replies=${replies.length} delivery=${delivery || "sync"} outboundId=${String((out as any).outbound_message_id || "")}`,
);
}
for (const r of replies) {
const outText = String((r as any).text || "").trim();
if (!shouldSendOutboundText(outText) && !(Array.isArray((r as any).attachments) && (r as any).attachments.length)) {
if (VERBOSE) log(`skip outbound reply chat=${chatId} (nonsend text)`);
continue;
}
const deliverTo = String((r as any).chat_id || chatId).trim() || chatId;
await sendReplyWithAttachments({ sock, deliverTo, text: outText, reply: r });
}
} catch (err) {
log(`handle message error: ${String(err)}`);
}
}
});
};
log(
`runner started local=${LOCAL_BASE_URL} stateDir=${STATE_DIR} verbose=${VERBOSE} node=${process.version} waVersion=${version.join(".")}`,
);
await logNetworkHints();
startOutboundPoller(() => sock);
await connectOnce();
}
void main();