mirror of
https://github.com/hansjone/dsh-im-ops.git
synced 2026-10-10 06:30:46 +08:00
feat: add WhatsApp QR channel
This commit is contained in:
parent
f8841da76f
commit
c54463843f
29 changed files with 4565 additions and 224 deletions
165
src/channels/whatsapp/config-store.mjs
Normal file
165
src/channels/whatsapp/config-store.mjs
Normal file
|
|
@ -0,0 +1,165 @@
|
|||
import { createHash } from 'node:crypto';
|
||||
import { mkdir, readFile, rename, unlink, writeFile } from 'node:fs/promises';
|
||||
import { dirname } from 'node:path';
|
||||
|
||||
const EMPTY_DOCUMENT = Object.freeze({ version: 2, bots: Object.freeze([]) });
|
||||
const BOT_ID_PATTERN = /^whatsapp_[a-f0-9]{24}$/;
|
||||
const AUTH_DIRECTORY_PATTERN = /^[a-f0-9-]{36}$/;
|
||||
|
||||
function cleanString(value) {
|
||||
return typeof value === 'string' && value.trim() ? value.trim() : null;
|
||||
}
|
||||
|
||||
export function normalizeWhatsappAccountJid(value) {
|
||||
const jid = cleanString(value)?.toLowerCase();
|
||||
return /^\d{5,32}@(s\.whatsapp\.net|lid)$/.test(jid ?? '') ? jid : null;
|
||||
}
|
||||
|
||||
export function deriveWhatsappBotId(accountJid) {
|
||||
const normalized = normalizeWhatsappAccountJid(accountJid);
|
||||
if (!normalized) throw new TypeError('A valid WhatsApp account JID is required');
|
||||
return `whatsapp_${createHash('sha256').update(normalized).digest('hex').slice(0, 24)}`;
|
||||
}
|
||||
|
||||
export function maskWhatsappAccount(accountJid) {
|
||||
const digits = normalizeWhatsappAccountJid(accountJid)?.split('@')[0] ?? '';
|
||||
if (!digits) return 'WhatsApp账号';
|
||||
if (digits.length <= 7) return `${digits.slice(0, 2)}•••${digits.slice(-2)}`;
|
||||
return `${digits.slice(0, 4)}••••${digits.slice(-4)}`;
|
||||
}
|
||||
|
||||
export class WhatsappConfigStore {
|
||||
#path;
|
||||
#value = EMPTY_DOCUMENT;
|
||||
#writeQueue = Promise.resolve();
|
||||
|
||||
constructor(path) {
|
||||
this.#path = path;
|
||||
}
|
||||
|
||||
async load() {
|
||||
try {
|
||||
const normalized = this.#normalizeDocument(JSON.parse(await readFile(this.#path, 'utf8')));
|
||||
if (!normalized) throw new Error('dsh-im WhatsApp config contains invalid account data');
|
||||
this.#value = normalized;
|
||||
} catch (error) {
|
||||
if (error?.code !== 'ENOENT') throw error;
|
||||
this.#value = EMPTY_DOCUMENT;
|
||||
}
|
||||
return this;
|
||||
}
|
||||
|
||||
list() {
|
||||
return structuredClone(this.#value.bots);
|
||||
}
|
||||
|
||||
get(botId) {
|
||||
const bot = this.#value.bots.find((candidate) => candidate.botId === botId);
|
||||
return bot ? structuredClone(bot) : null;
|
||||
}
|
||||
|
||||
getByAccountJid(accountJid) {
|
||||
const normalized = normalizeWhatsappAccountJid(accountJid);
|
||||
const bot = this.#value.bots.find((candidate) => candidate.accountJid === normalized);
|
||||
return bot ? structuredClone(bot) : null;
|
||||
}
|
||||
|
||||
async save(value) {
|
||||
const normalized = this.#normalizeBot(value);
|
||||
if (!normalized) throw new Error('Refusing to persist incomplete WhatsApp account data');
|
||||
return this.#mutate((bots) => {
|
||||
const duplicate = bots.find((bot) => bot.accountJid === normalized.accountJid
|
||||
&& bot.botId !== normalized.botId);
|
||||
const authCollision = bots.find((bot) => bot.authDirectory === normalized.authDirectory
|
||||
&& bot.botId !== normalized.botId);
|
||||
if (duplicate || authCollision) throw new Error('Duplicate WhatsApp account identity');
|
||||
const index = bots.findIndex((bot) => bot.botId === normalized.botId);
|
||||
if (index === -1) bots.push(normalized);
|
||||
else bots[index] = normalized;
|
||||
return structuredClone(normalized);
|
||||
});
|
||||
}
|
||||
|
||||
async remove(botId) {
|
||||
if (!BOT_ID_PATTERN.test(botId)) throw new TypeError('Invalid WhatsApp bot id');
|
||||
return this.#mutate((bots) => {
|
||||
const index = bots.findIndex((bot) => bot.botId === botId);
|
||||
if (index === -1) return null;
|
||||
const [removed] = bots.splice(index, 1);
|
||||
return structuredClone(removed);
|
||||
});
|
||||
}
|
||||
|
||||
async clear() {
|
||||
const operation = this.#writeQueue.then(async () => {
|
||||
try {
|
||||
await unlink(this.#path);
|
||||
} catch (error) {
|
||||
if (error?.code !== 'ENOENT') throw error;
|
||||
}
|
||||
this.#value = EMPTY_DOCUMENT;
|
||||
});
|
||||
this.#writeQueue = operation.then(() => undefined, () => undefined);
|
||||
await operation;
|
||||
}
|
||||
|
||||
#normalizeBot(value) {
|
||||
if (!value || typeof value !== 'object') return null;
|
||||
const accountJid = normalizeWhatsappAccountJid(value.accountJid);
|
||||
const botId = cleanString(value.botId);
|
||||
const authDirectory = cleanString(value.authDirectory);
|
||||
const name = cleanString(value.name);
|
||||
if (!accountJid || !botId || !authDirectory || !name
|
||||
|| !BOT_ID_PATTERN.test(botId) || !AUTH_DIRECTORY_PATTERN.test(authDirectory)
|
||||
|| deriveWhatsappBotId(accountJid) !== botId) return null;
|
||||
return Object.freeze({
|
||||
botId,
|
||||
accountJid,
|
||||
authDirectory,
|
||||
name,
|
||||
createdAt: cleanString(value.createdAt) ?? new Date().toISOString(),
|
||||
connectedAt: cleanString(value.connectedAt),
|
||||
});
|
||||
}
|
||||
|
||||
#normalizeDocument(value) {
|
||||
if (!value || value.version !== 2 || !Array.isArray(value.bots)) return null;
|
||||
const bots = value.bots.map((bot) => this.#normalizeBot(bot));
|
||||
if (bots.some((bot) => bot === null)) return null;
|
||||
const botIds = new Set();
|
||||
const accountJids = new Set();
|
||||
const authDirectories = new Set();
|
||||
for (const bot of bots) {
|
||||
if (botIds.has(bot.botId) || accountJids.has(bot.accountJid)
|
||||
|| authDirectories.has(bot.authDirectory)) return null;
|
||||
botIds.add(bot.botId);
|
||||
accountJids.add(bot.accountJid);
|
||||
authDirectories.add(bot.authDirectory);
|
||||
}
|
||||
return Object.freeze({ version: 2, bots: Object.freeze(bots) });
|
||||
}
|
||||
|
||||
async #mutate(mutator) {
|
||||
let result;
|
||||
const operation = this.#writeQueue.then(async () => {
|
||||
const bots = [...this.#value.bots];
|
||||
result = mutator(bots);
|
||||
const document = Object.freeze({ version: 2, bots: Object.freeze(bots) });
|
||||
await this.#write(document);
|
||||
this.#value = document;
|
||||
});
|
||||
this.#writeQueue = operation.then(() => undefined, () => undefined);
|
||||
await operation;
|
||||
return result;
|
||||
}
|
||||
|
||||
async #write(document) {
|
||||
await mkdir(dirname(this.#path), { recursive: true, mode: 0o700 });
|
||||
const temporary = `${this.#path}.tmp`;
|
||||
await writeFile(temporary, `${JSON.stringify(document, null, 2)}\n`, {
|
||||
encoding: 'utf8',
|
||||
mode: 0o600,
|
||||
});
|
||||
await rename(temporary, this.#path);
|
||||
}
|
||||
}
|
||||
3
src/channels/whatsapp/harness-client.mjs
Normal file
3
src/channels/whatsapp/harness-client.mjs
Normal file
|
|
@ -0,0 +1,3 @@
|
|||
import { HarnessClient } from '../weixin/harness-client.mjs';
|
||||
|
||||
export class WhatsappHarnessClient extends HarnessClient {}
|
||||
3
src/channels/whatsapp/state-store.mjs
Normal file
3
src/channels/whatsapp/state-store.mjs
Normal file
|
|
@ -0,0 +1,3 @@
|
|||
import { ConversationStateStore } from '../shared/conversation-state-store.mjs';
|
||||
|
||||
export class WhatsappStateStore extends ConversationStateStore {}
|
||||
15
src/channels/whatsapp/whatsapp-bridge.mjs
Normal file
15
src/channels/whatsapp/whatsapp-bridge.mjs
Normal file
|
|
@ -0,0 +1,15 @@
|
|||
import { createTextBridgeStatus, TextHarnessBridge } from '../shared/text-harness-bridge.mjs';
|
||||
|
||||
export const WHATSAPP_DESCRIPTOR = Object.freeze({
|
||||
key: 'whatsapp',
|
||||
label: 'WhatsApp',
|
||||
connectionLabel: ' Web 关联设备',
|
||||
});
|
||||
|
||||
export class WhatsappHarnessBridge extends TextHarnessBridge {
|
||||
constructor(options) {
|
||||
super({ ...options, descriptor: WHATSAPP_DESCRIPTOR });
|
||||
}
|
||||
}
|
||||
|
||||
export { createTextBridgeStatus as createWhatsappBridgeStatus };
|
||||
388
src/channels/whatsapp/whatsapp-controller.mjs
Normal file
388
src/channels/whatsapp/whatsapp-controller.mjs
Normal file
|
|
@ -0,0 +1,388 @@
|
|||
import { randomUUID } from 'node:crypto';
|
||||
|
||||
import { deriveWhatsappBotId, maskWhatsappAccount } from './config-store.mjs';
|
||||
|
||||
const ACTIVE_ATTEMPT_STATES = new Set(['starting', 'pending', 'connecting']);
|
||||
const TERMINAL_ATTEMPT_STATES = new Set(['connected', 'failed', 'cancelled']);
|
||||
const QR_TTL_MS = 60_000;
|
||||
|
||||
function safeError(code, message) {
|
||||
return Object.freeze({ code, message });
|
||||
}
|
||||
|
||||
function publicAttempt(record) {
|
||||
if (!record) return null;
|
||||
return {
|
||||
attemptId: record.id,
|
||||
status: record.state,
|
||||
qrRevision: record.qrRevision,
|
||||
pollIntervalMs: 1_000,
|
||||
...(record.qrValue ? { qrValue: record.qrValue } : {}),
|
||||
...(record.expiresAt ? { expiresAt: record.expiresAt } : {}),
|
||||
...(record.botId ? { botId: record.botId } : {}),
|
||||
...(record.error ? { error: structuredClone(record.error) } : {}),
|
||||
};
|
||||
}
|
||||
|
||||
export class WhatsappController {
|
||||
#configStore;
|
||||
#authPath;
|
||||
#createSession;
|
||||
#createRuntime;
|
||||
#deleteAuth;
|
||||
#deleteState;
|
||||
#logger;
|
||||
#runtimes = new Map();
|
||||
#errors = new Map();
|
||||
#attempts = new Map();
|
||||
#transitions = new Map();
|
||||
#activeAttemptId = null;
|
||||
#revision = 0;
|
||||
#closed = false;
|
||||
|
||||
constructor({
|
||||
configStore,
|
||||
authPath,
|
||||
createSession,
|
||||
createRuntime,
|
||||
deleteAuth = async () => {},
|
||||
deleteState = async () => {},
|
||||
logger = console,
|
||||
}) {
|
||||
if (!configStore || typeof configStore.list !== 'function'
|
||||
|| typeof configStore.save !== 'function' || typeof configStore.remove !== 'function') {
|
||||
throw new TypeError('WhatsappController requires a config store');
|
||||
}
|
||||
if (typeof authPath !== 'function' || typeof createSession !== 'function'
|
||||
|| typeof createRuntime !== 'function') {
|
||||
throw new TypeError('WhatsappController dependencies are incomplete');
|
||||
}
|
||||
this.#configStore = configStore;
|
||||
this.#authPath = authPath;
|
||||
this.#createSession = createSession;
|
||||
this.#createRuntime = createRuntime;
|
||||
this.#deleteAuth = deleteAuth;
|
||||
this.#deleteState = deleteState;
|
||||
this.#logger = logger;
|
||||
}
|
||||
|
||||
async initialize() {
|
||||
if (this.#closed) return this.status();
|
||||
for (const config of this.#configStore.list()) {
|
||||
await this.#withBotTransition(config.botId, async () => {
|
||||
if (this.#closed || this.#runtimes.get(config.botId)?.status?.ready) return;
|
||||
try {
|
||||
await this.#startRuntime(config);
|
||||
this.#errors.delete(config.botId);
|
||||
} catch (error) {
|
||||
this.#errors.set(config.botId, safeError(
|
||||
error?.code === 'relink-required' ? 'relink-required' : 'connection-failed',
|
||||
error?.code === 'relink-required'
|
||||
? 'WhatsApp 关联设备已失效,请移除后重新扫码。'
|
||||
: 'WhatsApp 连接未就绪,插件会自动重试。',
|
||||
));
|
||||
this.#logger.warn?.(`[dsh-im:whatsapp] bot ${config.botId} failed to initialize`);
|
||||
} finally {
|
||||
this.#touch();
|
||||
}
|
||||
});
|
||||
}
|
||||
return this.status();
|
||||
}
|
||||
|
||||
async startProvisioning() {
|
||||
if (this.#closed) throw new Error('WhatsApp controller is closed');
|
||||
if (this.#activeAttemptId) await this.cancelProvisioning(this.#activeAttemptId);
|
||||
let resolveFirstQr;
|
||||
let rejectFirstQr;
|
||||
let firstQrSettled = false;
|
||||
const firstQr = new Promise((resolve, reject) => {
|
||||
resolveFirstQr = () => {
|
||||
if (firstQrSettled) return;
|
||||
firstQrSettled = true;
|
||||
resolve();
|
||||
};
|
||||
rejectFirstQr = (error) => {
|
||||
if (firstQrSettled) return;
|
||||
firstQrSettled = true;
|
||||
reject(error);
|
||||
};
|
||||
});
|
||||
const id = randomUUID();
|
||||
const record = {
|
||||
id,
|
||||
state: 'starting',
|
||||
authDirectory: id,
|
||||
createdAt: Date.now(),
|
||||
expiresAt: null,
|
||||
qrRevision: 0,
|
||||
qrValue: null,
|
||||
controller: new AbortController(),
|
||||
session: null,
|
||||
task: null,
|
||||
error: null,
|
||||
botId: null,
|
||||
};
|
||||
this.#attempts.set(id, record);
|
||||
this.#activeAttemptId = id;
|
||||
this.#touch();
|
||||
|
||||
try {
|
||||
const session = await this.#createSession({
|
||||
authDir: this.#authPath(record.authDirectory),
|
||||
signal: record.controller.signal,
|
||||
logger: this.#logger,
|
||||
onQr: (value) => {
|
||||
if (record.controller.signal.aborted || TERMINAL_ATTEMPT_STATES.has(record.state)
|
||||
|| typeof value !== 'string' || !value) return;
|
||||
record.qrValue = value;
|
||||
record.qrRevision += 1;
|
||||
record.expiresAt = Date.now() + QR_TTL_MS;
|
||||
record.state = 'pending';
|
||||
this.#touch();
|
||||
resolveFirstQr();
|
||||
},
|
||||
});
|
||||
record.session = session;
|
||||
record.task = session.ready.then((identity) => this.#completeProvisioning(record, identity))
|
||||
.catch((error) => this.#failProvisioning(record, error, rejectFirstQr));
|
||||
await firstQr;
|
||||
return publicAttempt(record);
|
||||
} catch (error) {
|
||||
await this.#failProvisioning(record, error, rejectFirstQr);
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
registrationStatus(attemptId) {
|
||||
return publicAttempt(this.#attempts.get(attemptId));
|
||||
}
|
||||
|
||||
async cancelProvisioning(attemptId) {
|
||||
const record = this.#attempts.get(attemptId);
|
||||
if (!record) return null;
|
||||
if (!TERMINAL_ATTEMPT_STATES.has(record.state)) {
|
||||
record.controller.abort();
|
||||
await record.session?.close().catch(() => undefined);
|
||||
await record.task?.catch(() => undefined);
|
||||
if (!TERMINAL_ATTEMPT_STATES.has(record.state)) {
|
||||
record.state = 'cancelled';
|
||||
record.error = safeError('cancelled', '扫码接入已取消。');
|
||||
}
|
||||
await this.#deleteAuth(record.authDirectory).catch(() => undefined);
|
||||
}
|
||||
if (this.#activeAttemptId === record.id) this.#activeAttemptId = null;
|
||||
this.#touch();
|
||||
return publicAttempt(record);
|
||||
}
|
||||
|
||||
async reconnectBot(botId) {
|
||||
const config = this.#configStore.get(botId);
|
||||
if (!config) throw new Error('Unknown WhatsApp bot');
|
||||
await this.#withBotTransition(botId, async () => {
|
||||
try {
|
||||
await this.#startRuntime(config);
|
||||
this.#errors.delete(botId);
|
||||
} catch (error) {
|
||||
this.#errors.set(botId, safeError(
|
||||
error?.code === 'relink-required' ? 'relink-required' : 'connection-failed',
|
||||
error?.code === 'relink-required'
|
||||
? 'WhatsApp 关联设备已失效,请移除后重新扫码。'
|
||||
: 'WhatsApp 连接仍未就绪,请稍后重试。',
|
||||
));
|
||||
throw error;
|
||||
} finally {
|
||||
this.#touch();
|
||||
}
|
||||
});
|
||||
return this.status();
|
||||
}
|
||||
|
||||
async deleteBot(botId) {
|
||||
const config = this.#configStore.get(botId);
|
||||
if (!config) throw new Error('Unknown WhatsApp bot');
|
||||
await this.#withBotTransition(botId, async () => {
|
||||
await this.#stopRuntime(botId);
|
||||
try {
|
||||
await this.#configStore.remove(botId);
|
||||
} catch (error) {
|
||||
await this.#startRuntime(config).catch(() => undefined);
|
||||
throw error;
|
||||
}
|
||||
await Promise.allSettled([
|
||||
this.#deleteAuth(config.authDirectory),
|
||||
this.#deleteState({ botId, config }),
|
||||
]);
|
||||
this.#errors.delete(botId);
|
||||
this.#touch();
|
||||
});
|
||||
return this.status();
|
||||
}
|
||||
|
||||
status() {
|
||||
const bots = this.#configStore.list().map((config) => {
|
||||
const runtimeStatus = this.#runtimes.get(config.botId)?.status ?? null;
|
||||
const connected = runtimeStatus?.ready === true
|
||||
&& runtimeStatus.connectionState === 'connected'
|
||||
&& runtimeStatus.harnessReachable === true;
|
||||
const state = connected ? 'connected'
|
||||
: runtimeStatus?.connectionState === 'connecting' ? 'connecting'
|
||||
: this.#errors.has(config.botId) || runtimeStatus?.connectionState === 'failed'
|
||||
? 'error' : 'offline';
|
||||
return {
|
||||
botId: config.botId,
|
||||
state,
|
||||
connected,
|
||||
configured: true,
|
||||
bot: { name: config.name, idMasked: maskWhatsappAccount(config.accountJid) },
|
||||
health: {
|
||||
status: connected ? 'healthy' : state === 'error' ? 'error' : 'offline',
|
||||
summary: connected ? 'WhatsApp Web 关联设备运行正常'
|
||||
: state === 'error' ? 'WhatsApp 连接未就绪' : 'WhatsApp 连接当前离线',
|
||||
lastCheckedAt: runtimeStatus?.lastCheckedAt ?? null,
|
||||
lastConnectedAt: runtimeStatus?.lastConnectedAt ?? null,
|
||||
},
|
||||
stats: {
|
||||
messagesReceived: runtimeStatus?.messagesReceived ?? 0,
|
||||
messagesReplied: runtimeStatus?.messagesReplied ?? 0,
|
||||
},
|
||||
error: structuredClone(this.#errors.get(config.botId) ?? null),
|
||||
};
|
||||
});
|
||||
const connectedCount = bots.filter((bot) => bot.connected).length;
|
||||
const active = this.#activeAttemptId ? this.#attempts.get(this.#activeAttemptId) : null;
|
||||
return {
|
||||
schemaVersion: 1,
|
||||
revision: this.#revision,
|
||||
state: active && ACTIVE_ATTEMPT_STATES.has(active.state) ? 'provisioning'
|
||||
: bots.length === 0 ? 'disconnected'
|
||||
: connectedCount === bots.length ? 'connected'
|
||||
: connectedCount > 0 ? 'degraded' : 'offline',
|
||||
bots,
|
||||
totals: { configured: bots.length, connected: connectedCount },
|
||||
...(active && ACTIVE_ATTEMPT_STATES.has(active.state)
|
||||
? { provisioning: publicAttempt(active) } : {}),
|
||||
};
|
||||
}
|
||||
|
||||
async close() {
|
||||
if (this.#closed) return;
|
||||
this.#closed = true;
|
||||
if (this.#activeAttemptId) await this.cancelProvisioning(this.#activeAttemptId);
|
||||
await Promise.allSettled([...this.#transitions.values()]);
|
||||
await Promise.allSettled([...this.#runtimes.keys()].map((botId) => this.#stopRuntime(botId)));
|
||||
}
|
||||
|
||||
async #completeProvisioning(record, identity) {
|
||||
if (record.controller.signal.aborted || this.#closed) return;
|
||||
record.state = 'connecting';
|
||||
record.qrValue = null;
|
||||
record.expiresAt = null;
|
||||
this.#touch();
|
||||
const botId = deriveWhatsappBotId(identity.accountJid);
|
||||
record.botId = botId;
|
||||
await record.session?.close();
|
||||
const previous = this.#configStore.get(botId);
|
||||
const config = {
|
||||
botId,
|
||||
accountJid: identity.accountJid,
|
||||
authDirectory: record.authDirectory,
|
||||
name: identity.name,
|
||||
createdAt: previous?.createdAt ?? new Date().toISOString(),
|
||||
connectedAt: new Date().toISOString(),
|
||||
};
|
||||
try {
|
||||
if (record.controller.signal.aborted || this.#closed) throw Object.assign(new Error(), { name: 'AbortError' });
|
||||
await this.#configStore.save(config);
|
||||
if (record.controller.signal.aborted || this.#closed) throw Object.assign(new Error(), { name: 'AbortError' });
|
||||
try {
|
||||
await this.#startRuntime(config);
|
||||
this.#errors.delete(botId);
|
||||
} catch (error) {
|
||||
this.#errors.set(botId, safeError('connection-failed', 'WhatsApp 已绑定,消息连接暂未就绪。'));
|
||||
this.#logger.warn?.(`[dsh-im:whatsapp] bot ${botId} did not reconnect after QR binding`);
|
||||
}
|
||||
if (previous?.authDirectory && previous.authDirectory !== config.authDirectory) {
|
||||
await this.#deleteAuth(previous.authDirectory).catch(() => undefined);
|
||||
}
|
||||
record.state = 'connected';
|
||||
record.error = null;
|
||||
} catch (error) {
|
||||
if (record.controller.signal.aborted || this.#closed || error?.name === 'AbortError') {
|
||||
await this.#stopRuntime(botId);
|
||||
if (previous) await this.#configStore.save(previous).catch(() => undefined);
|
||||
else await this.#configStore.remove(botId).catch(() => undefined);
|
||||
await this.#deleteAuth(record.authDirectory).catch(() => undefined);
|
||||
if (previous) await this.#startRuntime(previous).catch(() => undefined);
|
||||
record.state = 'cancelled';
|
||||
record.error = safeError('cancelled', '扫码接入已取消。');
|
||||
} else {
|
||||
await this.#deleteAuth(record.authDirectory).catch(() => undefined);
|
||||
record.state = 'failed';
|
||||
record.error = safeError('activation-failed', 'WhatsApp 已扫码,但无法保存关联设备。');
|
||||
this.#logger.error?.('[dsh-im:whatsapp] unable to persist linked-device session');
|
||||
}
|
||||
} finally {
|
||||
if (this.#activeAttemptId === record.id) this.#activeAttemptId = null;
|
||||
this.#touch();
|
||||
}
|
||||
}
|
||||
|
||||
async #failProvisioning(record, error, rejectFirstQr = () => {}) {
|
||||
if (TERMINAL_ATTEMPT_STATES.has(record.state)) return;
|
||||
if (record.controller.signal.aborted || error?.name === 'AbortError') {
|
||||
record.state = 'cancelled';
|
||||
record.error = safeError('cancelled', '扫码接入已取消。');
|
||||
} else {
|
||||
record.state = 'failed';
|
||||
record.error = safeError('qr-connect-failed', '无法连接 WhatsApp,请重新生成二维码。');
|
||||
}
|
||||
if (this.#activeAttemptId === record.id) this.#activeAttemptId = null;
|
||||
await record.session?.close().catch(() => undefined);
|
||||
await this.#deleteAuth(record.authDirectory).catch(() => undefined);
|
||||
this.#touch();
|
||||
rejectFirstQr(error);
|
||||
}
|
||||
|
||||
async #startRuntime(config) {
|
||||
if (this.#closed) throw new Error('WhatsApp controller is closed');
|
||||
await this.#stopRuntime(config.botId);
|
||||
if (this.#closed) throw new Error('WhatsApp controller is closed');
|
||||
const runtime = await this.#createRuntime({
|
||||
botId: config.botId,
|
||||
config,
|
||||
authDir: this.#authPath(config.authDirectory),
|
||||
});
|
||||
if (!runtime || typeof runtime.start !== 'function' || typeof runtime.stop !== 'function') {
|
||||
throw new TypeError('createRuntime returned an invalid WhatsApp runtime');
|
||||
}
|
||||
this.#runtimes.set(config.botId, runtime);
|
||||
try {
|
||||
await runtime.start();
|
||||
} catch (error) {
|
||||
await runtime.stop().catch(() => undefined);
|
||||
this.#runtimes.delete(config.botId);
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
async #stopRuntime(botId) {
|
||||
const runtime = this.#runtimes.get(botId);
|
||||
this.#runtimes.delete(botId);
|
||||
await runtime?.stop().catch(() => undefined);
|
||||
}
|
||||
|
||||
#withBotTransition(botId, operation) {
|
||||
const previous = this.#transitions.get(botId) ?? Promise.resolve();
|
||||
const current = previous.catch(() => undefined).then(operation);
|
||||
const settled = current.finally(() => {
|
||||
if (this.#transitions.get(botId) === settled) this.#transitions.delete(botId);
|
||||
});
|
||||
this.#transitions.set(botId, settled);
|
||||
return settled;
|
||||
}
|
||||
|
||||
#touch() {
|
||||
this.#revision += 1;
|
||||
}
|
||||
}
|
||||
299
src/channels/whatsapp/whatsapp-runtime.mjs
Normal file
299
src/channels/whatsapp/whatsapp-runtime.mjs
Normal file
|
|
@ -0,0 +1,299 @@
|
|||
import {
|
||||
areJidsSameUser,
|
||||
normalizeMessageContent,
|
||||
} from '@whiskeysockets/baileys';
|
||||
|
||||
import { splitMessageText } from '../shared/editable-message-stream.mjs';
|
||||
import { createWhatsappBridgeStatus, WhatsappHarnessBridge } from './whatsapp-bridge.mjs';
|
||||
import { createWhatsappWebSession } from './whatsapp-web-session.mjs';
|
||||
|
||||
function messageContext(content) {
|
||||
return content?.extendedTextMessage?.contextInfo
|
||||
?? content?.imageMessage?.contextInfo
|
||||
?? content?.videoMessage?.contextInfo
|
||||
?? content?.documentMessage?.contextInfo
|
||||
?? null;
|
||||
}
|
||||
|
||||
function messageText(content) {
|
||||
return content?.conversation
|
||||
?? content?.extendedTextMessage?.text
|
||||
?? content?.imageMessage?.caption
|
||||
?? content?.videoMessage?.caption
|
||||
?? content?.documentMessage?.caption
|
||||
?? '';
|
||||
}
|
||||
|
||||
export function normalizeWhatsappMessage(message, accountJid) {
|
||||
const remoteJid = typeof message?.key?.remoteJid === 'string' ? message.key.remoteJid : '';
|
||||
const alternateRemoteJid = typeof message?.key?.remoteJidAlt === 'string'
|
||||
? message.key.remoteJidAlt : '';
|
||||
const messageId = typeof message?.key?.id === 'string' ? message.key.id : '';
|
||||
if (!remoteJid || !messageId || remoteJid === 'status@broadcast'
|
||||
|| remoteJid.endsWith('@newsletter')) return null;
|
||||
const group = remoteJid.endsWith('@g.us');
|
||||
const fromMe = message.key.fromMe === true;
|
||||
const selfChat = fromMe && !group
|
||||
&& [remoteJid, alternateRemoteJid].some((jid) => jid && areJidsSameUser(jid, accountJid));
|
||||
if (fromMe && !selfChat) return null;
|
||||
const senderJid = selfChat ? accountJid : group ? message.key.participant : remoteJid;
|
||||
if (typeof senderJid !== 'string' || !senderJid) return null;
|
||||
const content = normalizeMessageContent(message.message);
|
||||
const context = messageContext(content);
|
||||
const mentioned = Array.isArray(context?.mentionedJid)
|
||||
&& context.mentionedJid.some((jid) => areJidsSameUser(jid, accountJid));
|
||||
const replyToSelf = typeof context?.participant === 'string'
|
||||
&& areJidsSameUser(context.participant, accountJid);
|
||||
return {
|
||||
messageId: `${remoteJid}:${messageId}`,
|
||||
providerMessageId: messageId,
|
||||
senderId: senderJid,
|
||||
senderIsBot: false,
|
||||
kind: group ? 'group' : 'direct',
|
||||
conversationId: remoteJid,
|
||||
content: messageText(content),
|
||||
addressed: !group || mentioned || replyToSelf,
|
||||
selfChat,
|
||||
replyTarget: { jid: remoteJid, quoted: message, selfChat },
|
||||
};
|
||||
}
|
||||
|
||||
class RecentWhatsappOutboundIds {
|
||||
#ids = new Map();
|
||||
|
||||
has(id) {
|
||||
this.#purge();
|
||||
return this.#ids.has(id);
|
||||
}
|
||||
|
||||
remember(id) {
|
||||
if (typeof id !== 'string' || !id) return;
|
||||
this.#purge();
|
||||
this.#ids.set(id, Date.now() + 5 * 60_000);
|
||||
while (this.#ids.size > 256) this.#ids.delete(this.#ids.keys().next().value);
|
||||
}
|
||||
|
||||
#purge() {
|
||||
const now = Date.now();
|
||||
for (const [id, expiresAt] of this.#ids) {
|
||||
if (expiresAt > now) continue;
|
||||
this.#ids.delete(id);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
class WhatsappBotClient {
|
||||
#socket;
|
||||
#outboundIds;
|
||||
#typingTimers = new Map();
|
||||
|
||||
constructor(socket, outboundIds) {
|
||||
this.#socket = socket;
|
||||
this.#outboundIds = outboundIds;
|
||||
}
|
||||
|
||||
async sendText(target, text) {
|
||||
await this.#stopTyping(target.jid);
|
||||
let result = null;
|
||||
for (const [index, chunk] of splitMessageText(text, 4_000).entries()) {
|
||||
result = await this.#socket.sendMessage(
|
||||
target.jid,
|
||||
{ text: chunk },
|
||||
index === 0 && target.quoted ? { quoted: target.quoted } : undefined,
|
||||
);
|
||||
this.#outboundIds.remember(result?.key?.id);
|
||||
}
|
||||
return result;
|
||||
}
|
||||
|
||||
async sendTyping(target) {
|
||||
if (!target.selfChat && target.quoted?.key) {
|
||||
await this.#socket.readMessages([target.quoted.key]).catch(() => undefined);
|
||||
}
|
||||
await this.#socket.sendPresenceUpdate('composing', target.jid);
|
||||
await this.#stopTyping(target.jid, false);
|
||||
const timer = setInterval(() => {
|
||||
void this.#socket.sendPresenceUpdate('composing', target.jid).catch(() => {
|
||||
void this.#stopTyping(target.jid);
|
||||
});
|
||||
}, 20_000);
|
||||
timer.unref?.();
|
||||
this.#typingTimers.set(target.jid, timer);
|
||||
}
|
||||
|
||||
async close() {
|
||||
const jids = [...this.#typingTimers.keys()];
|
||||
await Promise.allSettled(jids.map((jid) => this.#stopTyping(jid)));
|
||||
}
|
||||
|
||||
async #stopTyping(jid, sendPaused = true) {
|
||||
const timer = this.#typingTimers.get(jid);
|
||||
if (timer) clearInterval(timer);
|
||||
this.#typingTimers.delete(jid);
|
||||
if (sendPaused) await this.#socket.sendPresenceUpdate('paused', jid).catch(() => undefined);
|
||||
}
|
||||
}
|
||||
|
||||
export function createWhatsappRuntimeStatus() {
|
||||
return {
|
||||
startedAt: null,
|
||||
ready: false,
|
||||
connectionState: 'idle',
|
||||
harnessReachable: false,
|
||||
lastCheckedAt: null,
|
||||
lastConnectedAt: null,
|
||||
lastError: null,
|
||||
...createWhatsappBridgeStatus(),
|
||||
};
|
||||
}
|
||||
|
||||
export class WhatsappRuntime {
|
||||
#config;
|
||||
#authDir;
|
||||
#harness;
|
||||
#state;
|
||||
#logger;
|
||||
#replyTimeoutMs;
|
||||
#connectTimeoutMs;
|
||||
#createSession;
|
||||
#status = createWhatsappRuntimeStatus();
|
||||
#abortController = null;
|
||||
#session = null;
|
||||
#client = null;
|
||||
#bridge = null;
|
||||
#starting = null;
|
||||
|
||||
constructor({
|
||||
config,
|
||||
authDir,
|
||||
harness,
|
||||
state,
|
||||
logger = console,
|
||||
replyTimeoutMs = 600_000,
|
||||
connectTimeoutMs = 30_000,
|
||||
createSession = createWhatsappWebSession,
|
||||
}) {
|
||||
if (!config || !authDir || !harness || !state || typeof createSession !== 'function') {
|
||||
throw new TypeError('WhatsappRuntime requires config, auth directory, Harness, state, and session factory');
|
||||
}
|
||||
this.#config = config;
|
||||
this.#authDir = authDir;
|
||||
this.#harness = harness;
|
||||
this.#state = state;
|
||||
this.#logger = logger;
|
||||
this.#replyTimeoutMs = replyTimeoutMs;
|
||||
this.#connectTimeoutMs = connectTimeoutMs;
|
||||
this.#createSession = createSession;
|
||||
}
|
||||
|
||||
get status() {
|
||||
return structuredClone(this.#status);
|
||||
}
|
||||
|
||||
async start() {
|
||||
if (this.#status.ready && this.#session) return this.status;
|
||||
if (this.#starting) return this.#starting;
|
||||
this.#starting = this.#start().finally(() => { this.#starting = null; });
|
||||
return this.#starting;
|
||||
}
|
||||
|
||||
async #start() {
|
||||
await this.stop();
|
||||
this.#status.startedAt = new Date().toISOString();
|
||||
this.#status.connectionState = 'connecting';
|
||||
this.#status.lastError = null;
|
||||
await this.#harness.ensureRunning();
|
||||
this.#status.harnessReachable = true;
|
||||
const controller = new AbortController();
|
||||
this.#abortController = controller;
|
||||
const outboundIds = new RecentWhatsappOutboundIds();
|
||||
let rejectRelink;
|
||||
const relinkRequired = new Promise((_, reject) => { rejectRelink = reject; });
|
||||
void relinkRequired.catch(() => undefined);
|
||||
try {
|
||||
const session = await this.#createSession({
|
||||
authDir: this.#authDir,
|
||||
signal: controller.signal,
|
||||
logger: this.#logger,
|
||||
onQr: () => rejectRelink(Object.assign(
|
||||
new Error('WhatsApp linked-device session must be scanned again'),
|
||||
{ code: 'relink-required' },
|
||||
)),
|
||||
onMessage: async (raw) => {
|
||||
const message = normalizeWhatsappMessage(raw, this.#config.accountJid);
|
||||
if (!message || outboundIds.has(message.providerMessageId) || !this.#bridge) return;
|
||||
this.#status.lastCheckedAt = Date.now();
|
||||
await this.#bridge.accept(message);
|
||||
},
|
||||
onDisconnect: ({ error }) => {
|
||||
if (controller.signal.aborted) return;
|
||||
this.#status.ready = false;
|
||||
this.#status.connectionState = 'failed';
|
||||
this.#status.lastError = error?.message ?? 'WhatsApp Web connection closed';
|
||||
},
|
||||
});
|
||||
this.#session = session;
|
||||
let timer;
|
||||
const identity = await Promise.race([
|
||||
session.ready,
|
||||
relinkRequired,
|
||||
new Promise((_, reject) => {
|
||||
timer = setTimeout(
|
||||
() => reject(new Error('WhatsApp Web did not connect in time')),
|
||||
this.#connectTimeoutMs,
|
||||
);
|
||||
}),
|
||||
]).finally(() => clearTimeout(timer));
|
||||
if (!areJidsSameUser(identity.accountJid, this.#config.accountJid)) {
|
||||
throw new Error('WhatsApp linked account does not match the saved bot');
|
||||
}
|
||||
const client = new WhatsappBotClient(session.socket, outboundIds);
|
||||
this.#client = client;
|
||||
this.#bridge = new WhatsappHarnessBridge({
|
||||
bot: client,
|
||||
harness: this.#harness,
|
||||
state: this.#state,
|
||||
status: this.#status,
|
||||
logger: this.#logger,
|
||||
replyTimeoutMs: this.#replyTimeoutMs,
|
||||
});
|
||||
const now = Date.now();
|
||||
this.#status.ready = true;
|
||||
this.#status.connectionState = 'connected';
|
||||
this.#status.lastCheckedAt = now;
|
||||
this.#status.lastConnectedAt = now;
|
||||
return this.status;
|
||||
} catch (error) {
|
||||
this.#status.ready = false;
|
||||
this.#status.connectionState = 'failed';
|
||||
this.#status.lastError = error?.message ?? String(error);
|
||||
await this.stop();
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
async logout() {
|
||||
await this.#session?.logout().catch(() => undefined);
|
||||
return this.stop();
|
||||
}
|
||||
|
||||
async stop() {
|
||||
const session = this.#session;
|
||||
const client = this.#client;
|
||||
const bridge = this.#bridge;
|
||||
this.#abortController?.abort();
|
||||
this.#abortController = null;
|
||||
this.#session = null;
|
||||
this.#client = null;
|
||||
this.#bridge = null;
|
||||
await client?.close().catch(() => undefined);
|
||||
await session?.close().catch(() => undefined);
|
||||
await Promise.race([
|
||||
bridge?.waitForIdle() ?? Promise.resolve(),
|
||||
new Promise((resolve) => setTimeout(resolve, 2_000)),
|
||||
]);
|
||||
this.#status.ready = false;
|
||||
this.#status.connectionState = 'idle';
|
||||
return this.status;
|
||||
}
|
||||
}
|
||||
212
src/channels/whatsapp/whatsapp-web-session.mjs
Normal file
212
src/channels/whatsapp/whatsapp-web-session.mjs
Normal file
|
|
@ -0,0 +1,212 @@
|
|||
import { chmod, mkdir, readdir } from 'node:fs/promises';
|
||||
|
||||
import makeWASocket, {
|
||||
Browsers,
|
||||
DisconnectReason,
|
||||
jidNormalizedUser,
|
||||
useMultiFileAuthState,
|
||||
} from '@whiskeysockets/baileys';
|
||||
|
||||
const SILENT_LOGGER = Object.freeze({
|
||||
level: 'silent',
|
||||
trace() {},
|
||||
debug() {},
|
||||
info() {},
|
||||
warn() {},
|
||||
error() {},
|
||||
fatal() {},
|
||||
child() { return this; },
|
||||
});
|
||||
const APPEND_RECENT_GRACE_MS = 60_000;
|
||||
|
||||
function abortError() {
|
||||
return Object.assign(new Error('WhatsApp connection was cancelled'), { name: 'AbortError' });
|
||||
}
|
||||
|
||||
function disconnectStatus(error) {
|
||||
return error?.output?.statusCode ?? error?.data?.statusCode ?? error?.statusCode ?? null;
|
||||
}
|
||||
|
||||
function messageTimestampMs(value) {
|
||||
let seconds = value;
|
||||
if (typeof seconds === 'string') {
|
||||
if (!/^\d+$/.test(seconds)) return null;
|
||||
seconds = Number(seconds);
|
||||
} else if (typeof seconds === 'bigint') {
|
||||
seconds = Number(seconds);
|
||||
} else if (seconds && typeof seconds === 'object') {
|
||||
seconds = Number(seconds.valueOf());
|
||||
}
|
||||
return Number.isFinite(seconds) && seconds >= 0 ? seconds * 1_000 : null;
|
||||
}
|
||||
|
||||
async function hardenAuthDirectory(path) {
|
||||
await mkdir(path, { recursive: true, mode: 0o700 });
|
||||
await chmod(path, 0o700);
|
||||
const entries = await readdir(path, { withFileTypes: true }).catch(() => []);
|
||||
await Promise.all(entries.filter((entry) => entry.isFile())
|
||||
.map((entry) => chmod(`${path}/${entry.name}`, 0o600).catch(() => undefined)));
|
||||
}
|
||||
|
||||
function normalizeIdentity(socket, authState) {
|
||||
const source = socket.user ?? authState.creds.me;
|
||||
const accountJid = jidNormalizedUser(source?.id);
|
||||
if (!/^\d{5,32}@(s\.whatsapp\.net|lid)$/.test(accountJid ?? '')) {
|
||||
throw new Error('WhatsApp did not return a valid linked account');
|
||||
}
|
||||
return {
|
||||
accountJid,
|
||||
name: typeof source?.name === 'string' && source.name.trim()
|
||||
? source.name.trim().slice(0, 100) : 'WhatsApp机器人',
|
||||
};
|
||||
}
|
||||
|
||||
export async function createWhatsappWebSession({
|
||||
authDir,
|
||||
onQr,
|
||||
onMessage,
|
||||
onDisconnect,
|
||||
signal,
|
||||
logger = console,
|
||||
makeSocket = makeWASocket,
|
||||
loadAuthState = useMultiFileAuthState,
|
||||
} = {}) {
|
||||
if (!authDir || typeof onQr !== 'function') {
|
||||
throw new TypeError('WhatsApp Web session requires an auth directory and QR callback');
|
||||
}
|
||||
await hardenAuthDirectory(authDir);
|
||||
const sessionStartedAt = Date.now();
|
||||
const { state, saveCreds } = await loadAuthState(authDir);
|
||||
const originalKeySet = state.keys.set.bind(state.keys);
|
||||
state.keys.set = async (data) => {
|
||||
await originalKeySet(data);
|
||||
await hardenAuthDirectory(authDir);
|
||||
};
|
||||
|
||||
let closed = false;
|
||||
let readySettled = false;
|
||||
let resolveReady;
|
||||
let rejectReady;
|
||||
let saveQueue = Promise.resolve();
|
||||
let socket = null;
|
||||
let socketGeneration = 0;
|
||||
let restartTask = null;
|
||||
const ready = new Promise((resolve, reject) => {
|
||||
resolveReady = resolve;
|
||||
rejectReady = reject;
|
||||
});
|
||||
|
||||
const settleFailure = (error) => {
|
||||
if (readySettled) return;
|
||||
readySettled = true;
|
||||
rejectReady(error);
|
||||
};
|
||||
const close = async () => {
|
||||
if (closed) return;
|
||||
closed = true;
|
||||
socketGeneration += 1;
|
||||
settleFailure(abortError());
|
||||
await restartTask?.catch(() => undefined);
|
||||
await saveQueue.catch(() => undefined);
|
||||
await socket?.end(undefined).catch(() => undefined);
|
||||
};
|
||||
const logout = async () => {
|
||||
if (closed) return;
|
||||
closed = true;
|
||||
socketGeneration += 1;
|
||||
settleFailure(abortError());
|
||||
await restartTask?.catch(() => undefined);
|
||||
await saveQueue.catch(() => undefined);
|
||||
await socket?.logout('Removed from DeepSeek Harness').catch(() => undefined);
|
||||
};
|
||||
|
||||
const startSocket = () => {
|
||||
const generation = ++socketGeneration;
|
||||
let connectionOpen = false;
|
||||
const nextSocket = makeSocket({
|
||||
auth: state,
|
||||
browser: Browsers.macOS('DeepSeek Harness'),
|
||||
logger: SILENT_LOGGER,
|
||||
markOnlineOnConnect: false,
|
||||
syncFullHistory: false,
|
||||
shouldSyncHistoryMessage: () => false,
|
||||
getMessage: async () => undefined,
|
||||
generateHighQualityLinkPreview: false,
|
||||
});
|
||||
socket = nextSocket;
|
||||
|
||||
const resolveWhenLinked = () => {
|
||||
if (!connectionOpen || !state.creds.me || readySettled) return;
|
||||
void saveQueue.then(() => {
|
||||
if (closed || readySettled || generation !== socketGeneration
|
||||
|| !connectionOpen || !state.creds.me) return;
|
||||
readySettled = true;
|
||||
resolveReady(normalizeIdentity(nextSocket, state));
|
||||
}).catch((error) => settleFailure(error));
|
||||
};
|
||||
|
||||
nextSocket.ev.on('creds.update', () => {
|
||||
if (closed || generation !== socketGeneration) return;
|
||||
saveQueue = saveQueue.then(async () => {
|
||||
await saveCreds();
|
||||
await hardenAuthDirectory(authDir);
|
||||
});
|
||||
saveQueue.catch(() => logger.error?.('[dsh-im:whatsapp] failed to persist linked-device state'));
|
||||
resolveWhenLinked();
|
||||
});
|
||||
nextSocket.ev.on('connection.update', (update) => {
|
||||
if (closed || generation !== socketGeneration) return;
|
||||
if (typeof update.qr === 'string' && update.qr) onQr(update.qr);
|
||||
if (update.connection === 'open') {
|
||||
connectionOpen = true;
|
||||
resolveWhenLinked();
|
||||
}
|
||||
if (update.connection === 'close') {
|
||||
const status = disconnectStatus(update.lastDisconnect?.error);
|
||||
if (status === DisconnectReason.restartRequired) {
|
||||
restartTask ??= saveQueue.then(async () => {
|
||||
if (closed || generation !== socketGeneration) return;
|
||||
await nextSocket.end(undefined).catch(() => undefined);
|
||||
if (closed || generation !== socketGeneration) return;
|
||||
startSocket();
|
||||
}).catch((error) => settleFailure(error)).finally(() => {
|
||||
restartTask = null;
|
||||
});
|
||||
return;
|
||||
}
|
||||
const loggedOut = status === DisconnectReason.loggedOut;
|
||||
const error = Object.assign(new Error(loggedOut
|
||||
? 'WhatsApp linked device was removed from the phone'
|
||||
: 'WhatsApp Web connection closed'), { code: loggedOut ? 'logged-out' : 'connection-closed' });
|
||||
if (!readySettled) settleFailure(error);
|
||||
else onDisconnect?.({ error, loggedOut });
|
||||
}
|
||||
});
|
||||
nextSocket.ev.on('messages.upsert', ({ messages, type }) => {
|
||||
if (closed || generation !== socketGeneration
|
||||
|| (type !== 'notify' && type !== 'append') || typeof onMessage !== 'function') return;
|
||||
for (const message of Array.isArray(messages) ? messages : []) {
|
||||
if (type === 'append') {
|
||||
const timestamp = messageTimestampMs(message?.messageTimestamp);
|
||||
if (timestamp === null || timestamp < sessionStartedAt - APPEND_RECENT_GRACE_MS) continue;
|
||||
}
|
||||
Promise.resolve(onMessage(message)).catch(() => {
|
||||
logger.error?.('[dsh-im:whatsapp] failed to process an inbound WhatsApp message');
|
||||
});
|
||||
}
|
||||
});
|
||||
};
|
||||
|
||||
startSocket();
|
||||
|
||||
if (signal) {
|
||||
if (signal.aborted) await close();
|
||||
else signal.addEventListener('abort', () => void close(), { once: true });
|
||||
}
|
||||
return Object.freeze({
|
||||
get socket() { return socket; },
|
||||
ready,
|
||||
close,
|
||||
logout,
|
||||
});
|
||||
}
|
||||
Loading…
Add table
Add a link
Reference in a new issue