feat: hooks, plugins, apps

- MaintainOnlineStatusPlugin (Activity)
- MessageSourceCachePlugin
- WebhookPlugin
- Wid*Plugin
This commit is contained in:
devlikepro committed 2026-08-31 10:44:22 +07:00
1 parent 99e4dfcfb5
commit 61f1f71343
30 files changed
+1508 -623

No files matched your search

+1
View File
@@ -120,6 +120,7 @@
"shell-quote": "^1.8.3",
"sqlite3": "^5.1.7",
"swagger-ui-express": "^4.1.4",
"tapable": "^2.2.2",
"ulid": "^2.3.0",
"undici": "^7.16.0",
"uniqid": "^5.4.0",
@@ -212,6 +212,11 @@ export class AppsEnabledService implements IAppsService {
if (!service && !AppRuntimeConfig.HasApp(app.app)) {
throw new AppDisableError(app.app);
}
const plugins = service.plugins(app, session);
for (const plugin of plugins) {
const key = `${plugin.constructor.name}:${app.id}`;
session.plugins[key] = plugin;
}
service.beforeSessionStart(app, session);
}
}
+7
View File
@@ -1,6 +1,7 @@
import { App } from '@waha/apps/app_sdk/dto/app.dto';
import { SessionManager } from '@waha/core/abc/manager.abc';
import { WhatsappSession } from '@waha/core/abc/session.abc';
import { SessionPlugin } from '@waha/core/abc/session.plugin';
/**
* Exact App service
@@ -54,6 +55,12 @@ export interface IAppService {
*/
enrich(manager: SessionManager, app: App): Promise<void>;
/**
* Session plugins contributed by this app.
* The session manager registers their hooks and events before session.start().
*/
plugins(app: App, session: WhatsappSession): SessionPlugin<any>[];
beforeSessionStart(app: App, session: WhatsappSession): void;
afterSessionStart(app: App, session: WhatsappSession): void;
+17 -9
View File
@@ -2,18 +2,13 @@ import { Injectable } from '@nestjs/common';
import { App } from '@waha/apps/app_sdk/dto/app.dto';
import { IAppService } from '@waha/apps/app_sdk/services/IAppService';
import { CallsAppConfig } from '@waha/apps/calls/dto/config.dto';
import { CallsPlugin } from '@waha/apps/calls/services/CallsPlugin';
import { SessionManager } from '@waha/core/abc/manager.abc';
import { WhatsappSession } from '@waha/core/abc/session.abc';
import { InjectPinoLogger, PinoLogger } from 'nestjs-pino';
import { CallsListener } from '@waha/apps/calls/services/CallsListener';
import { SessionPlugin } from '@waha/core/abc/session.plugin';
@Injectable()
export class CallsAppService implements IAppService {
constructor(
@InjectPinoLogger('CallsAppService')
private readonly logger: PinoLogger,
) {}
validate(app: App<CallsAppConfig>): void {
// The DTO validation covers structure; no extra validation rules.
void app;
@@ -87,9 +82,22 @@ export class CallsAppService implements IAppService {
void app;
}
plugins(
app: App<CallsAppConfig>,
session: WhatsappSession,
): SessionPlugin<any>[] {
const logger = session.loggerBuilder.child({
plugin: CallsPlugin.name,
app: app.app,
});
const plugin = new CallsPlugin(session, logger, app.config);
return [plugin];
}
beforeSessionStart(app: App<CallsAppConfig>, session: WhatsappSession): void {
const listener = new CallsListener(app, session, this.logger);
listener.attach();
void app;
void session;
return;
}
afterSessionStart(app: App<CallsAppConfig>, session: WhatsappSession): void {
@@ -1,62 +1,29 @@
import { Subscription } from 'rxjs';
import {
CallsAppChannelConfig,
CallsAppConfig,
} from '@waha/apps/calls/dto/config.dto';
import { Logger } from 'pino';
import { App } from '@waha/apps/app_sdk/dto/app.dto';
import { PinoLogger } from 'nestjs-pino';
import { WhatsappSession } from '@waha/core/abc/session.abc';
import { WAHAEvents, WAHAPresenceStatus } from '@waha/structures/enums.dto';
import { SessionPlugin } from '@waha/core/abc/session.plugin';
import { PluginEvent } from '@waha/core/abc/session.plugin.events';
import { CallData } from '@waha/structures/calls.dto';
import { MessageTextRequest } from '@waha/structures/chatting.dto';
import { WAHAEvents, WAHAPresenceStatus } from '@waha/structures/enums.dto';
import { sleep } from '@waha/utils/promiseTimeout';
import { Observable } from 'rxjs';
export class CallsListener {
private subscription?: Subscription;
private config: CallsAppConfig;
private readonly log: Logger;
private readonly session: WhatsappSession;
constructor(
app: App<CallsAppConfig>,
session: WhatsappSession,
logger: PinoLogger,
) {
this.session = session;
this.config = app.config;
this.log = logger.logger.child({
app: 'calls',
session: app.session,
});
}
attach(): void {
this.detach();
const observable = this.session.getEventObservable(
WAHAEvents.CALL_RECEIVED,
);
if (!observable) {
this.log.warn('CALL_RECEIVED event stream is not available, skipping');
return;
}
this.subscription = observable.subscribe((payload) => {
this.handleCall(payload as CallData).catch((error) => {
this.log.error(
{ err: error, callId: (payload as any)?.id },
/**
* Rejects incoming calls and optionally sends an auto-reply message.
*/
export class CallsPlugin extends SessionPlugin<CallsAppConfig> {
@PluginEvent(WAHAEvents.CALL_RECEIVED)
onCallReceived(calls$: Observable<CallData>) {
calls$.subscribe((call: CallData) => {
this.handleCall(call).catch((error) => {
this.logger.error(
{ err: error, callId: call?.id },
'Failed to handle incoming call',
);
});
});
this.log.info('Calls app listener is attached');
}
detach(): void {
this.subscription?.unsubscribe();
this.subscription = undefined;
}
private configFor(call: CallData): CallsAppChannelConfig {
@@ -65,17 +32,17 @@ export class CallsListener {
private async handleCall(call: CallData): Promise<void> {
if (!call.from) {
this.log.warn({ call: call?.id }, 'Incoming call has no chat id');
this.logger.warn({ call: call?.id }, 'Incoming call has no chat id');
return;
}
if (!call.id) {
this.log.warn({ from: call.from }, 'Incoming call has no from');
this.logger.warn({ from: call.from }, 'Incoming call has no from');
return;
}
const config = this.configFor(call);
if (!config) {
this.log.warn({ callId: call.id }, 'No calls config found, skipping');
this.logger.warn({ callId: call.id }, 'No calls config found, skipping');
return;
}
@@ -84,7 +51,7 @@ export class CallsListener {
const shouldMessage = message.length > 0;
if (!shouldReject && !shouldMessage) {
this.log.debug(
this.logger.debug(
{ callId: call.id, chatId: call.from },
'No actions configured for this call',
);
@@ -105,16 +72,19 @@ export class CallsListener {
}
private async rejectCall(call: CallData): Promise<void> {
this.log.debug({ from: call.from, id: call.id }, 'Rejecting incoming call');
this.logger.debug(
{ from: call.from, id: call.id },
'Rejecting incoming call',
);
await this.session.rejectCall(call.from, call.id);
this.log.info({ from: call.from, id: call.id }, 'Call rejected');
this.logger.info({ from: call.from, id: call.id }, 'Call rejected');
}
private async replyWithTyping(
chatId: string,
message: string,
): Promise<void> {
this.log.info(
this.logger.info(
{ chatId: chatId },
'Sending auto-response for rejected call',
);
@@ -132,7 +102,7 @@ export class CallsListener {
try {
await this.session.setPresence(WAHAPresenceStatus.TYPING, chatId);
} catch (error) {
this.log.warn(
this.logger.warn(
{ err: error, chatId: chatId },
'Failed to set typing presence before reply',
);
@@ -146,7 +116,7 @@ export class CallsListener {
try {
await this.session.setPresence(WAHAPresenceStatus.PAUSED, chatId);
} catch (error) {
this.log.warn(
this.logger.warn(
{ err: error, chatId: chatId },
'Failed to clear typing presence after reply',
);
@@ -8,6 +8,7 @@ import { ChatWootWAHAQueueService } from '@waha/apps/chatwoot/services/ChatWootW
import { App } from '@waha/apps/chatwoot/storage';
import { SessionManager } from '@waha/core/abc/manager.abc';
import { WhatsappSession } from '@waha/core/abc/session.abc';
import { SessionPlugin } from '@waha/core/abc/session.plugin';
import { InjectPinoLogger, PinoLogger } from 'nestjs-pino';
import { DIContainer } from '../di/DIContainer';
@@ -152,6 +153,15 @@ export class ChatWootAppService implements IAppService {
await conversation.incoming(updated);
}
plugins(
app: App<ChatWootAppConfig>,
session: WhatsappSession,
): SessionPlugin<any>[] {
void app;
void session;
return [];
}
beforeSessionStart(app: App<ChatWootAppConfig>, session: WhatsappSession) {
this.chatWootWAHAQueueService.listenEvents(app, session);
}
+10
View File
@@ -6,6 +6,7 @@ import { McpAppConfig } from '@waha/apps/mcp/dto/config.dto';
import { SessionManager } from '@waha/core/abc/manager.abc';
import { ApiKeyService } from '@waha/core/services/ApiKeyService';
import { WhatsappSession } from '@waha/core/abc/session.abc';
import { SessionPlugin } from '@waha/core/abc/session.plugin';
import { InjectPinoLogger, PinoLogger } from 'nestjs-pino';
@Injectable()
@@ -114,6 +115,15 @@ export class McpAppService implements IAppService {
app.config = { ...(app.config ?? {}), key: keyDto.key } as McpAppConfig;
}
plugins(
app: App<McpAppConfig>,
session: WhatsappSession,
): SessionPlugin<any>[] {
void app;
void session;
return [];
}
beforeSessionStart(app: App<McpAppConfig>, session: WhatsappSession): void {
void app;
void session;
+1 -1
View File
@@ -1,6 +1,6 @@
import { Injectable, Logger, OnApplicationBootstrap } from '@nestjs/common';
import { ConfigService } from '@nestjs/config';
import { GlobalWebhookConfigConfig } from '@waha/core/config/GlobalWebhookConfig';
import { GlobalWebhookConfigConfig } from '@waha/plugins/WebhookPlugin.env';
import { MediaDownloadOptions } from '@waha/core/media/IMediaManager';
import { IgnoreJidConfig } from '@waha/core/utils/jids';
import * as lodash from 'lodash';
+17 -109
View File
@@ -1,6 +1,7 @@
import { MessageCappingTracker } from '@waha/core/abc/MessageCappingTracker';
import { ReachoutTimelockTracker } from '@waha/core/abc/ReachoutTimelockTracker';
import { getBrowserExecutablePath as getBrowserExecutablePathAutodetect } from '@waha/core/abc/session.browser';
import type { SessionPlugin } from '@waha/core/abc/session.plugin';
import { IMediaConverter } from '@waha/core/media/IConverter';
import { Ffmpeg } from '@waha/core/utils/ffmpeg';
import { MessagesForRead } from '@waha/core/utils/convertors';
@@ -37,7 +38,7 @@ import { BinaryFile, RemoteFile } from '@waha/structures/files.dto';
import { Label, LabelDTO, LabelID } from '@waha/structures/labels.dto';
import { LidToPhoneNumber } from '@waha/structures/lids.dto';
import { PaginationParams } from '@waha/structures/pagination.dto';
import { MessageSource, WAMessage } from '@waha/structures/responses.dto';
import { WAMessage } from '@waha/structures/responses.dto';
import { BrowserTraceQuery } from '@waha/structures/server.debug.dto';
import { DefaultMap } from '@waha/utils/DefaultMap';
import { generatePrefixedId } from '@waha/utils/ids';
@@ -123,6 +124,7 @@ import {
ProxyConfig,
ReachoutTimelockData,
SessionConfig,
SessionInfo,
} from '../../structures/sessions.dto';
import {
DeleteStatusRequest,
@@ -143,10 +145,7 @@ import { IMediaManager, MediaDownloadOptions } from '../media/IMediaManager';
import { QR } from '../QR';
import { DataStore } from './DataStore';
import { fetchBuffer } from '@waha/utils/fetch';
import {
PRESENCE_AUTO_ONLINE,
PRESENCE_AUTO_ONLINE_DURATION_SECONDS,
} from '@waha/core/env';
import { SessionHooks } from './session.hooks';
// eslint-disable-next-line @typescript-eslint/no-var-requires
const qrcode = require('qrcode-terminal');
@@ -218,11 +217,6 @@ export abstract class WhatsappSession {
| WAHAPresenceStatus.ONLINE
| WAHAPresenceStatus.OFFLINE
| null = null;
private lastActivityTimestamp?: number;
protected presenceAutoOnlineConfig = {
enabled: PRESENCE_AUTO_ONLINE,
duration: PRESENCE_AUTO_ONLINE_DURATION_SECONDS * 1000,
};
private shouldPrintQR: boolean;
protected events2: DefaultMap<WAHAEvents, SwitchObservable<any>>;
@@ -231,15 +225,9 @@ export abstract class WhatsappSession {
stdTTL: 24 * 60 * 60, // 1 day
});
// Save sent messages ids in cache so we can determine if a message was sent
// via API or APP
private sentMessageIds: NodeCache = new NodeCache({
stdTTL: 10 * 60, // 10 minutes
});
private presenceOfflineTimeout?: ReturnType<typeof setTimeout>;
public mediaConverter: IMediaConverter;
public hooks: SessionHooks;
public plugins: Record<string, SessionPlugin<any>> = {};
public constructor({
name,
@@ -255,11 +243,12 @@ export abstract class WhatsappSession {
}: SessionParams) {
this._status = WAHASessionStatus.STOPPED;
this.status$ = new Subject<SessionStatusUpdate>();
this.name = name;
this.proxyConfig = proxyConfig;
this.loggerBuilder = loggerBuilder;
this.logger = loggerBuilder.child({ name: 'WhatsappSession' });
this.hooks = new SessionHooks();
this.mediaConverter = new Ffmpeg(this.name, this.logger);
this.reachoutTimelock = new ReachoutTimelockTracker(this.logger);
this.reachoutTimelock.changes$.subscribe((timelock) => {
@@ -436,7 +425,7 @@ export abstract class WhatsappSession {
return this._statusData;
}
protected set presence(value: WAHAPresenceStatus) {
public set presence(value: WAHAPresenceStatus) {
switch (value) {
case null:
this._presence = null;
@@ -591,6 +580,14 @@ export abstract class WhatsappSession {
return null;
}
/**
* Builds the runtime session info via the "session.info" hook.
*/
public getSessionInfo(): Promise<SessionInfo> {
const info = {} as SessionInfo;
return this.hooks.session.info.promise(info);
}
/**
* Profile methods
*/
@@ -707,75 +704,6 @@ export abstract class WhatsappSession {
abstract stopTyping(chat: ChatRequest);
/**
* Activity tracking and presence management
*/
/**
* Returns the timestamp of the last "activity" in the session
* @returns Timestamp in milliseconds or undefined if there was never any activity
*/
public getLastActivityTimestamp(): number | undefined {
return this.lastActivityTimestamp;
}
/**
* Maintains ONLINE presence active while there is activity
* Resets the timer on each activity, only goes OFFLINE after Xs without activity
*/
async maintainPresenceOnline(): Promise<void> {
if (!this.presenceAutoOnlineConfig.enabled) {
return;
}
if (this.status !== WAHASessionStatus.WORKING) {
return;
}
this.lastActivityTimestamp = Date.now();
// If not ONLINE yet, send ONLINE
if (this._presence !== WAHAPresenceStatus.ONLINE) {
try {
// Force set ONLINE in case of many requests comes at the same time
// So we'll set ONLINE exactly once
this.presence = WAHAPresenceStatus.ONLINE;
await this.setPresence(WAHAPresenceStatus.ONLINE);
this.logger.debug('Set presence to ONLINE due to activity');
} catch (error) {
this.logger.debug('Failed to set presence ONLINE', error);
return;
}
}
// Cancel the previous timeout (if exists)
this.cleanupPresenceTimeout();
// Schedule to go back OFFLINE after timeout without activity
this.presenceOfflineTimeout = setTimeout(async () => {
try {
const working = this.status === WAHASessionStatus.WORKING;
const online = this.presence === WAHAPresenceStatus.ONLINE;
if (!working || !online) {
// Nothing to do
return;
}
await this.setPresence(WAHAPresenceStatus.OFFLINE);
this.logger.debug(
'Auto-set presence to OFFLINE after time without activity',
);
} catch (error) {
this.presence = WAHAPresenceStatus.OFFLINE;
this.logger.debug('Failed to set presence OFFLINE', error);
}
this.cleanupPresenceTimeout();
}, this.presenceAutoOnlineConfig.duration);
}
/**
* Cleans up the timeout when the session stops
*/
protected cleanupPresenceTimeout() {
clearTimeout(this.presenceOfflineTimeout);
this.presenceOfflineTimeout = null;
}
abstract setReaction(request: MessageReactionRequest);
setStar(request: MessageStarRequest): Promise<void> {
@@ -1316,14 +1244,6 @@ export abstract class WhatsappSession {
* END - Methods for API
*/
/**
* Add WhatsApp suffix (@c.us) to the phone number if it doesn't have it yet
* @param phone
*/
protected ensureSuffix(phone) {
return ensureSuffix(phone);
}
protected deserializeId(messageId: string): MessageId {
const parts = messageId.split('_');
return {
@@ -1350,18 +1270,6 @@ export abstract class WhatsappSession {
qrcode.generate(qr.raw, { small: true });
}
protected saveSentMessageId(id: string) {
this.sentMessageIds.set(id, true);
}
protected getMessageSource(id: string): MessageSource {
if (!id) {
return MessageSource.APP;
}
const api = this.sentMessageIds.has(id);
return api ? MessageSource.API : MessageSource.APP;
}
/**
* Fetches the content from the specified URL and returns it as a Buffer.
*/
@@ -1,8 +1,7 @@
import type { WhatsappSession } from '@waha/core/abc/session.abc';
/**
* Decorator to mark a method as an activity that
* keeps the WhatsApp session online.
* Decorator to call "activity" hook
* @constructor
*/
export function Activity() {
@@ -16,7 +15,7 @@ export function Activity() {
this: WhatsappSession,
...args: Parameters<T>
): Promise<ReturnType<T>> {
await this.maintainPresenceOnline();
await this.hooks.activity.promise(String(propertyKey));
return await original.apply(this, args);
} as T;
+64
View File
@@ -0,0 +1,64 @@
import { MessageSource } from '@waha/structures/responses.dto';
import { SessionInfo } from '@waha/structures/sessions.dto';
import {
AsyncSeriesBailHook,
AsyncSeriesHook,
AsyncSeriesWaterfallHook,
SyncHook,
} from 'tapable';
/**
* Well-known stages for hook taps. Lower stage runs earlier, taps without a stage run in between (default stage is 0).
*/
export const Stage = {
FIRST: -100,
LAST: 100,
};
export class SessionHooks {
/**
* Called before an engine method that makes a network call to WhatsApp.
* method - the name of the method about to run.
*/
readonly activity = new AsyncSeriesHook<[string]>(['method'], 'activity');
readonly session = Object.freeze({
/**
* Builds the runtime session info (waterfall) - starts with an empty object, each tap merges its part in.
* Called when the API reports session state, e.g. GET /api/sessions.
*/
info: new AsyncSeriesWaterfallHook<[SessionInfo]>(['info'], 'session.info'),
});
/**
* Converts a WhatsApp identifier (wid) supplied by the API caller into the form the engine addresses.
* "wid.chat" converts the chat target, "wid.mention" converts a mention entry.
* method - the engine method the wid is being converted for, e.g. 'sendText'.
*/
readonly wid = Object.freeze({
chat: new AsyncSeriesWaterfallHook<[string, string]>(
['wid', 'method'],
'wid.chat',
),
mention: new AsyncSeriesWaterfallHook<[string, string]>(
['wid', 'method'],
'wid.mention',
),
});
readonly message = Object.freeze({
/**
* Called with the message id when the session sends a message via API.
* Fire and forget - return values are ignored.
*/
sent: new SyncHook<[string]>(['id'], 'message.sent'),
/**
* Resolves the MessageSource (api/app) for a message id. The first tap returning a non-undefined result wins,
* if all return undefined - callers fall back to MessageSource.APP.
*/
source: new AsyncSeriesBailHook<[string], MessageSource | undefined>(
['id'],
'message.source',
),
});
}
+62
View File
@@ -0,0 +1,62 @@
import type { SessionPlugin } from '@waha/core/abc/session.plugin';
import { WAHAEvents } from '@waha/structures/enums.dto';
import { Observable } from 'rxjs';
type EventFn = (events$: Observable<any>) => void;
interface EventSubscriptionMetadata {
event: WAHAEvents;
propertyKey: string;
}
const EVENT_SUBSCRIPTIONS = Symbol('PluginEventSubscriptions');
function getOwnEventSubscriptions(ctor: any): EventSubscriptionMetadata[] {
if (!Object.prototype.hasOwnProperty.call(ctor, EVENT_SUBSCRIPTIONS)) {
ctor[EVENT_SUBSCRIPTIONS] = [];
}
return ctor[EVENT_SUBSCRIPTIONS];
}
function collectEventSubscriptions(ctor: any): EventSubscriptionMetadata[] {
const subscriptions: EventSubscriptionMetadata[] = [];
let current = ctor;
while (current) {
if (Object.prototype.hasOwnProperty.call(current, EVENT_SUBSCRIPTIONS)) {
subscriptions.push(...current[EVENT_SUBSCRIPTIONS]);
}
current = Object.getPrototypeOf(current);
}
return subscriptions;
}
/**
* Passes the session event observable to the decorated (public) method once, the method subscribes itself:
*
* @PluginEvent(WAHAEvents.SESSION_STATUS)
* onSessionStatus(events$: Observable<WASessionStatusBody>) {
* events$.subscribe({ next: ..., complete: ... });
* }
*/
export function PluginEvent(event: WAHAEvents) {
return function <K extends string, T extends Record<K, EventFn>>(
target: T,
propertyKey: K,
descriptor: PropertyDescriptor,
): void {
getOwnEventSubscriptions(target.constructor).push({
event: event,
propertyKey: propertyKey,
});
};
}
/**
* Calls all @PluginEvent methods of the plugin's class with the session event observables.
*/
export function RegisterPluginEvents(plugin: SessionPlugin<any>) {
for (const meta of collectEventSubscriptions(plugin.constructor)) {
const method = (plugin as any)[meta.propertyKey].bind(plugin);
method(plugin.session.getEventObservable(meta.event));
}
}
+129
View File
@@ -0,0 +1,129 @@
import { SessionHooks } from '@waha/core/abc/session.hooks';
import type { SessionPlugin } from '@waha/core/abc/session.plugin';
import {
SyncBailHook,
SyncHook,
SyncLoopHook,
SyncWaterfallHook,
} from 'tapable';
/**
* Which tapable method the tap is registered with.
*/
export enum TapType {
/** tap() for sync hooks, tapPromise() for async ones - the default */
Auto = 'auto',
/** tap() - the method runs synchronously, its return value is used as is */
Sync = 'sync',
/** tapPromise() - the method result is awaited (works for sync methods too) */
Promise = 'promise',
}
/**
* Tapable tap options (tapable does not export its TapOptions type) and the tap type.
*/
export interface HookTapOptions {
before?: string;
stage?: number;
type?: TapType;
}
/**
* The method signature a hook accepts, sync or async - inferred from the hook's tap() method.
*/
type HookFn<H> = H extends {
tap(options: any, fn: (...args: infer A) => infer R): void;
}
? (...args: A) => R | Promise<R>
: never;
type TapableHook = {
tap(options: any, fn: any): void;
};
interface HookTapMetadata {
selector: (hooks: SessionHooks) => TapableHook;
propertyKey: string;
options?: HookTapOptions;
}
const HOOK_TAPS = Symbol('PluginHookTaps');
function getOwnHookTaps(ctor: any): HookTapMetadata[] {
if (!Object.prototype.hasOwnProperty.call(ctor, HOOK_TAPS)) {
ctor[HOOK_TAPS] = [];
}
return ctor[HOOK_TAPS];
}
function collectHookTaps(ctor: any): HookTapMetadata[] {
const taps: HookTapMetadata[] = [];
let current = ctor;
while (current) {
if (Object.prototype.hasOwnProperty.call(current, HOOK_TAPS)) {
taps.push(...current[HOOK_TAPS]);
}
current = Object.getPrototypeOf(current);
}
return taps;
}
// identity check, not instanceof - tapable reassigns hook.constructor while all hooks share the base Hook prototype
const SYNC_HOOK_CLASSES: any[] = [
SyncHook,
SyncBailHook,
SyncLoopHook,
SyncWaterfallHook,
];
function isSyncHook(hook: any): boolean {
return SYNC_HOOK_CLASSES.includes(hook.constructor);
}
function useSyncTap(type: TapType, hook: any): boolean {
if (type === TapType.Sync) {
return true;
}
if (type === TapType.Promise) {
return false;
}
return isSyncHook(hook);
}
/**
* Taps the decorated (public) method into a session hook, stackable:
* @PluginHook((hooks) => hooks.wid.chat, { stage: Stage.FIRST })
*/
export function PluginHook<H extends TapableHook>(
selector: (hooks: SessionHooks) => H,
options?: HookTapOptions,
) {
return function <K extends string, T extends Record<K, HookFn<H>>>(
target: T,
propertyKey: K,
descriptor: PropertyDescriptor,
): void {
getOwnHookTaps(target.constructor).push({
selector: selector,
propertyKey: propertyKey,
options: options,
});
};
}
/**
* Applies all @PluginHook taps of the plugin's class to its session hooks, named after the plugin class.
*/
export function RegisterPluginHooks(plugin: SessionPlugin<any>) {
for (const meta of collectHookTaps(plugin.constructor)) {
const hook: any = meta.selector(plugin.session.hooks);
const { type, ...tapOptions } = meta.options ?? {};
const options = { name: plugin.constructor.name, ...tapOptions };
const method = (plugin as any)[meta.propertyKey].bind(plugin);
if (useSyncTap(type ?? TapType.Auto, hook)) {
hook.tap(options, method);
} else {
hook.tapPromise(options, async (...args: any[]) => method(...args));
}
}
}
+10
View File
@@ -0,0 +1,10 @@
import type { WhatsappSession } from '@waha/core/abc/session.abc';
import { Logger } from 'pino';
export abstract class SessionPlugin<Config = void> {
constructor(
public readonly session: WhatsappSession,
protected logger: Logger,
protected config: Config,
) {}
}
+111 -44
View File
@@ -44,6 +44,7 @@ import {
} from '@waha/core/engines/noweb/session.noweb.core';
import { extractMediaContent } from '@waha/core/engines/noweb/utils';
import { NotImplementedByEngineError } from '@waha/core/exceptions';
import { WidToJIDPlugin } from '@waha/plugins/WidToJIDPlugin';
import { IMediaEngineProcessor } from '@waha/core/media/IMediaEngineProcessor';
import { LottieMediaProcessorWrapper } from '@waha/core/media/LottieMediaProcessorWrapper';
import { QR } from '@waha/core/QR';
@@ -53,9 +54,7 @@ import {
isJidBroadcast,
isJidGroup,
isJidNewsletter,
normalizeJid,
toCusFormat,
toJID,
} from '@waha/core/utils/jids';
import {
PasskeyChallenge,
@@ -212,7 +211,6 @@ import { IsEditedMessage } from '@waha/core/utils/pwa';
import { GoToJSWAProto } from '@waha/core/engines/gows/waproto';
import { extractWALocation } from '@waha/core/engines/waproto/locaiton';
import { extractVCards } from '@waha/core/engines/waproto/vcards';
import { Activity } from '@waha/core/abc/activity';
import { TmpDir } from '@waha/utils/tmpdir';
import { detectMimetype } from '@waha/utils/files';
import { WAMimeType } from '@waha/core/media/WAMimeType';
@@ -224,6 +222,8 @@ import MessageServiceClient = messages.MessageServiceClient;
import * as fsp from 'fs/promises';
import { MediaDownloadOptions } from '@waha/core/media/IMediaManager';
import { Activity } from '@waha/core/abc/session.hooks.activity';
axiosRetry(axios, { retries: 3 });
function getGowsStorageConfig(
@@ -301,6 +301,10 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
public constructor(config) {
super(config);
this.plugins[WidToJIDPlugin.name] = new WidToJIDPlugin(
this,
this.loggerBuilder.child({ plugin: WidToJIDPlugin.name }),
);
this.qr = new QR();
this.session = new messages.Session({ id: this.name });
this.presences = new NodeCache({
@@ -420,14 +424,12 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
events.on(WhatsMeowEvent.DISCONNECTED, () => {
if (this.status != WAHASessionStatus.STARTING) {
this.cleanupPresenceTimeout();
this.presence = null;
this.status = WAHASessionStatus.STARTING;
}
});
events.on(WhatsMeowEvent.KEEP_ALIVE_TIMEOUT, () => {
if (this.status != WAHASessionStatus.STARTING) {
this.cleanupPresenceTimeout();
this.presence = null;
this.status = WAHASessionStatus.STARTING;
}
@@ -866,7 +868,10 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
@Activity()
async fetchContactProfilePicture(id: string): Promise<string> {
const jid = normalizeJid(toJID(this.ensureSuffix(id)));
const jid = await this.hooks.wid.chat.promise(
id,
'fetchContactProfilePicture',
);
const request = new messages.ProfilePictureRequest({
jid: jid,
session: this.session,
@@ -883,7 +888,6 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
}
async stop(): Promise<void> {
this.cleanupPresenceTimeout();
this.status = WAHASessionStatus.STOPPED;
this.events?.stop();
this.stopEvents();
@@ -1049,6 +1053,7 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
return true;
}
@Activity()
protected async deleteProfilePicture(): Promise<boolean> {
const request = new messages.SetProfilePictureRequest({
session: this.session,
@@ -1071,9 +1076,10 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
@Activity()
async rejectCall(from: string, id: string): Promise<void> {
const jid = await this.hooks.wid.chat.promise(from, 'rejectCall');
const request = new messages.RejectCallRequest({
session: this.session,
from: normalizeJid(toJID(this.ensureSuffix(from))),
from: jid,
id: id,
});
await promisify(this.client.RejectCall)(request);
@@ -1081,7 +1087,15 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
@Activity()
async sendText(request: MessageTextRequest) {
const jid = normalizeJid(toJID(this.ensureSuffix(request.chatId)));
const jid = await this.hooks.wid.chat.promise(request.chatId, 'sendText');
let mentions: string[] | undefined;
if (request.mentions) {
mentions = await Promise.all(
request.mentions.map((mention) =>
this.hooks.wid.mention.promise(mention, 'sendText'),
),
);
}
const message = new messages.MessageRequest({
id: request.id,
jid: jid,
@@ -1090,9 +1104,7 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
linkPreview: request.linkPreview ?? true,
linkPreviewHighQuality: request.linkPreviewHighQuality,
replyTo: getMessageIdFromSerialized(request.reply_to),
mentions: request.mentions?.map((mention) =>
normalizeJid(toJID(mention)),
),
mentions: mentions,
});
const response = await promisify(this.client.SendMessage)(message);
const data = response.toObject();
@@ -1105,7 +1117,7 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
messageId: string,
request: EditMessageRequest,
) {
const jid = normalizeJid(toJID(this.ensureSuffix(chatId)));
const jid = await this.hooks.wid.chat.promise(chatId, 'editMessage');
const key = parseMessageIdSerialized(messageId, true);
const message = new messages.EditMessageRequest({
session: this.session,
@@ -1122,7 +1134,10 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
@Activity()
async sendContactVCard(request: MessageContactVcardRequest) {
const jid = normalizeJid(toJID(this.ensureSuffix(request.chatId)));
const jid = await this.hooks.wid.chat.promise(
request.chatId,
'sendContactVCard',
);
const contacts = request.contacts.map((el) => ({
displayName:
(el as any).fullName || parseVCardV3(el.vcard || '').fullName,
@@ -1142,7 +1157,7 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
@Activity()
async sendPoll(request: MessagePollRequest) {
const jid = normalizeJid(toJID(request.chatId));
const jid = await this.hooks.wid.chat.promise(request.chatId, 'sendPoll');
const message = new messages.MessageRequest({
id: request.id,
jid: jid,
@@ -1161,7 +1176,10 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
@Activity()
async sendPollVote(request: MessagePollVoteRequest) {
const jid = normalizeJid(toJID(this.ensureSuffix(request.chatId)));
const jid = await this.hooks.wid.chat.promise(
request.chatId,
'sendPollVote',
);
const key = parseMessageIdSerialized(request.pollMessageId, true);
const pollVote = new messages.PollVoteMessage({
pollMessageId: key.id,
@@ -1183,7 +1201,7 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
@Activity()
async sendList(request: SendListRequest): Promise<any> {
const jid = normalizeJid(toJID(this.ensureSuffix(request.chatId)));
const jid = await this.hooks.wid.chat.promise(request.chatId, 'sendList');
if (isJidGroup(jid) || isJidBroadcast(jid) || isJidNewsletter(jid)) {
throw new UnprocessableEntityException(
`List message can only be sent to a direct message chat.`,
@@ -1210,7 +1228,7 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
@Activity()
public async deleteMessage(chatId: string, messageId: string) {
const jid = normalizeJid(toJID(this.ensureSuffix(chatId)));
const jid = await this.hooks.wid.chat.promise(chatId, 'deleteMessage');
const key = parseMessageIdSerialized(messageId);
const message = new messages.RevokeMessageRequest({
session: this.session,
@@ -1227,7 +1245,11 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
if (!contacts || contacts.length == 0) {
return [];
}
return contacts.map((c) => normalizeJid(toJID(c)));
return await Promise.all(
contacts.map((c) =>
this.hooks.wid.chat.promise(c, 'prepareJidsForStatus'),
),
);
}
@Activity()
@@ -1253,6 +1275,7 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
return this.messageResponse(Jid.BROADCAST, data);
}
@Activity()
public async deleteStatus(request: DeleteStatusRequest) {
const participants = await this.prepareJidsForStatus(request.contacts);
const key = parseMessageIdSerialized(request.id, true);
@@ -1304,7 +1327,10 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
@Activity()
async sendLocation(request: MessageLocationRequest) {
const jid = normalizeJid(toJID(this.ensureSuffix(request.chatId)));
const jid = await this.hooks.wid.chat.promise(
request.chatId,
'sendLocation',
);
const message = new messages.MessageRequest({
id: request.id,
jid: jid,
@@ -1326,7 +1352,7 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
}
private async sendMedia(type: messages.MediaType, request: any) {
const jid = normalizeJid(toJID(this.ensureSuffix(request.chatId)));
const jid = await this.hooks.wid.chat.promise(request.chatId, 'sendMedia');
const media = await this.fileToMedia(request.file);
media.type = type;
if (type === messages.MediaType.IMAGE) {
@@ -1378,6 +1404,14 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
});
}
const participants = await this.prepareJidsForStatus(request.contacts);
let mentions: string[] | undefined;
if (request.mentions) {
mentions = await Promise.all(
request.mentions.map((mention) =>
this.hooks.wid.mention.promise(mention, 'sendMedia'),
),
);
}
const message = new messages.MessageRequest({
id: request.id,
jid: jid,
@@ -1385,9 +1419,7 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
session: this.session,
media: media,
backgroundColor: backgroundColor,
mentions: request.mentions?.map((mention) =>
normalizeJid(toJID(mention)),
),
mentions: mentions,
participants: participants,
});
@@ -1463,7 +1495,10 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
async sendLinkCustomPreview(
request: MessageLinkCustomPreviewRequest,
): Promise<any> {
const jid = normalizeJid(toJID(this.ensureSuffix(request.chatId)));
const jid = await this.hooks.wid.chat.promise(
request.chatId,
'sendLinkCustomPreview',
);
const media = await this.fileToMedia(request.preview.image as RemoteFile);
const preview = new messages.LinkPreview({
url: request.preview.url,
@@ -1490,7 +1525,10 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
throw new NotImplementedByEngineError();
// Doesn't work yet
const jid = normalizeJid(toJID(this.ensureSuffix(request.chatId)));
const jid = await this.hooks.wid.chat.promise(
request.chatId,
'sendButtonsReply',
);
const message = new messages.ButtonReplyRequest({
jid: jid,
session: this.session,
@@ -1573,10 +1611,15 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
*/
@Activity()
public async createGroup(request: CreateGroupRequest) {
const participants = await Promise.all(
request.participants.map((p) =>
this.hooks.wid.chat.promise(p.id, 'createGroup'),
),
);
const req = new messages.CreateGroupRequest({
session: this.session,
name: request.name,
participants: request.participants.map((p) => normalizeJid(toJID(p.id))),
participants: participants,
});
const response = await promisify(this.client.CreateGroup)(req);
const data = parseJson(response);
@@ -1763,7 +1806,11 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
participants: Array<Participant>,
action: messages.ParticipantRequestAction,
): Promise<GroupJoinRequestResult[]> {
const jids = participants.map((p) => normalizeJid(toJID(p.id)));
const jids = await Promise.all(
participants.map((p) =>
this.hooks.wid.chat.promise(p.id, 'updateRequestParticipants'),
),
);
const req = new messages.UpdateRequestParticipantsRequest({
session: this.session,
jid: id,
@@ -1869,7 +1916,11 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
participants: Array<Participant>,
action: messages.ParticipantAction,
): Promise<any> {
const jids = participants.map((p) => normalizeJid(toJID(p.id)));
const jids = await Promise.all(
participants.map((p) =>
this.hooks.wid.chat.promise(p.id, 'updateParticipants'),
),
);
const req = new messages.UpdateParticipantsRequest({
session: this.session,
jid: id,
@@ -1922,7 +1973,7 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
@Activity()
async sendEvent(request: EventMessageRequest): Promise<WAMessage> {
const jid = normalizeJid(toJID(this.ensureSuffix(request.chatId)));
const jid = await this.hooks.wid.chat.promise(request.chatId, 'sendEvent');
const event = request.event;
// Create EventLocation if provided
@@ -1978,7 +2029,10 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
public async setPresence(presence: WAHAPresenceStatus, chatId?: string) {
let request: any;
let method: any;
const jid = chatId ? normalizeJid(toJID(this.ensureSuffix(chatId))) : null;
let jid: string | null = null;
if (chatId) {
jid = await this.hooks.wid.chat.promise(chatId, 'setPresence');
}
switch (presence) {
case WAHAPresenceStatus.ONLINE:
request = new messages.PresenceRequest({
@@ -1995,7 +2049,7 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
method = this.client.SendPresence;
break;
case WAHAPresenceStatus.TYPING:
await this.maintainPresenceOnline();
await this.hooks.activity.promise('setPresence');
request = new messages.ChatPresenceRequest({
session: this.session,
jid: jid,
@@ -2004,7 +2058,7 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
method = this.client.SendChatPresence;
break;
case WAHAPresenceStatus.RECORDING:
await this.maintainPresenceOnline();
await this.hooks.activity.promise('setPresence');
request = new messages.ChatPresenceRequest({
session: this.session,
jid: jid,
@@ -2013,7 +2067,7 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
method = this.client.SendChatPresence;
break;
case WAHAPresenceStatus.PAUSED:
await this.maintainPresenceOnline();
await this.hooks.activity.promise('setPresence');
request = new messages.ChatPresenceRequest({
session: this.session,
jid: jid,
@@ -2039,7 +2093,7 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
}
public async getPresence(chatId: string): Promise<WAHAChatPresences> {
const jid = normalizeJid(toJID(chatId));
const jid = await this.hooks.wid.chat.promise(chatId, 'getPresence');
await this.subscribePresence(jid);
if (!(jid in this.presences.keys())) {
await sleep(1000);
@@ -2050,7 +2104,7 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
@Activity()
async subscribePresence(chatId: string) {
const jid = normalizeJid(toJID(chatId));
const jid = await this.hooks.wid.chat.promise(chatId, 'subscribePresence');
const req = new messages.SubscribePresenceRequest({
session: this.session,
jid: jid,
@@ -2256,6 +2310,7 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
};
}
@Activity()
public async channelsList(query: ListChannelsQuery): Promise<Channel[]> {
const request = new messages.NewsletterListRequest({
session: this.session,
@@ -2357,7 +2412,7 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
*/
@Activity()
public async upsertContact(chatId: string, body: ContactUpdateBody) {
const jid = normalizeJid(toJID(chatId));
const jid = await this.hooks.wid.chat.promise(chatId, 'upsertContact');
const request = new messages.UpdateContactRequest({
session: this.session,
jid: jid,
@@ -2377,7 +2432,10 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
}
public async getContact(query: ContactQuery) {
const jid = normalizeJid(toJID(query.contactId));
const jid = await this.hooks.wid.chat.promise(
query.contactId,
'getContact',
);
const request = new messages.EntityByIdRequest({
session: this.session,
id: jid,
@@ -2472,7 +2530,10 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
public async findLIDByPhoneNumber(
phoneNumber: string,
): Promise<LidToPhoneNumber> {
const pn = normalizeJid(toJID(phoneNumber));
const pn = await this.hooks.wid.chat.promise(
phoneNumber,
'findLIDByPhoneNumber',
);
const request = new messages.EntityByIdRequest({
session: this.session,
id: pn,
@@ -2551,7 +2612,9 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
}
let jids = [];
if (filter?.ids && filter.ids.length > 0) {
jids = filter.ids.map((id) => normalizeJid(toJID(id)));
jids = await Promise.all(
filter.ids.map((id) => this.hooks.wid.chat.promise(id, 'getChats')),
);
}
const request = new messages.GetChatsRequest({
session: this.session,
@@ -2586,8 +2649,12 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
if (chatId === 'all') {
jid = null;
} else {
const value = await this.hooks.wid.chat.promise(
chatId,
'getChatMessages',
);
jid = new messages.OptionalString({
value: normalizeJid(toJID(this.ensureSuffix(chatId))),
value: value,
});
}
@@ -2746,7 +2813,7 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
}
public async getChatLabels(chatId: string): Promise<Label[]> {
const jid = normalizeJid(toJID(chatId));
const jid = await this.hooks.wid.chat.promise(chatId, 'getChatLabels');
const request = new messages.EntityByIdRequest({
session: this.session,
id: jid,
@@ -2758,7 +2825,7 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
@Activity()
public async chatsUnreadChat(chatId: string): Promise<any> {
const jid = normalizeJid(toJID(this.ensureSuffix(chatId)));
const jid = await this.hooks.wid.chat.promise(chatId, 'chatsUnreadChat');
const request = new messages.ChatUnreadRequest({
session: this.session,
jid: jid,
@@ -2770,7 +2837,7 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
@Activity()
public async putLabelsToChat(chatId: string, labels: LabelID[]) {
const jid = normalizeJid(toJID(chatId));
const jid = await this.hooks.wid.chat.promise(chatId, 'putLabelsToChat');
const labelsIds = labels.map((label) => label.id);
const currentLabels = await this.getChatLabels(jid);
const currentLabelsIds = currentLabels.map((label) => label.id);
+185 -69
View File
@@ -68,6 +68,7 @@ import {
import { NowebAuthFactoryCore } from '@waha/core/engines/noweb/NowebAuthFactoryCore';
import { NowebInMemoryStore } from '@waha/core/engines/noweb/store/NowebInMemoryStore';
import { NotImplementedByEngineError } from '@waha/core/exceptions';
import { WidToJIDPlugin } from '@waha/plugins/WidToJIDPlugin';
import { toVcardV3 } from '@waha/core/vcard';
import { createAgentProxy } from '@waha/core/helpers.proxy';
import type { Agent } from 'https';
@@ -180,7 +181,11 @@ import {
WAHAChatPresences,
WAHAPresenceData,
} from '@waha/structures/presence.dto';
import { WAMessage, WAMessageReaction } from '@waha/structures/responses.dto';
import {
MessageSource,
WAMessage,
WAMessageReaction,
} from '@waha/structures/responses.dto';
import {
MeInfo,
MessageCappingData,
@@ -215,6 +220,7 @@ import * as lodash from 'lodash';
import * as NodeCache from 'node-cache';
import {
filter,
concatMap,
fromEvent,
groupBy,
identity,
@@ -246,7 +252,6 @@ import {
} from '@waha/core/utils/secretEncryptedMessageEdit';
import { extractWALocation } from '@waha/core/engines/waproto/locaiton';
import { extractVCards } from '@waha/core/engines/waproto/vcards';
import { Activity } from '@waha/core/abc/activity';
import { WAMimeType } from '@waha/core/media/WAMimeType';
import {
WAHA_CLIENT_BROWSER_NAME,
@@ -257,6 +262,8 @@ import esm from '@waha/vendor/esm';
import axios from 'axios';
import axiosRetry from 'axios-retry';
import { formatWaVersion } from '@waha/core/engines/noweb/waversion';
import { Activity } from '@waha/core/abc/session.hooks.activity';
// eslint-disable-next-line @typescript-eslint/no-var-requires
const promiseRetry = require('promise-retry');
@@ -313,6 +320,10 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
public constructor(config) {
super(config);
this.plugins[WidToJIDPlugin.name] = new WidToJIDPlugin(
this,
this.loggerBuilder.child({ plugin: WidToJIDPlugin.name }),
);
this.shouldRestart = true;
this.qr = new QR();
@@ -902,7 +913,6 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
}
private async end() {
this.cleanupPresenceTimeout();
this.presence = null;
this.autoRestartJob.stop();
const sock = this.sock;
@@ -1038,6 +1048,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
return true;
}
@Activity()
protected async deleteProfilePicture(): Promise<boolean> {
const me = this.getSessionMeInfo();
await this.sock.removeProfilePicture(me.id);
@@ -1102,16 +1113,27 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
@Activity()
async rejectCall(from: string, id: string): Promise<void> {
const jid = toJID(this.ensureSuffix(from));
const jid = await this.hooks.wid.chat.promise(from, 'rejectCall');
await this.sock.rejectCall(id, jid);
}
@Activity()
async sendText(request: MessageTextRequest) {
const chatId = toJID(this.ensureSuffix(request.chatId));
const chatId = await this.hooks.wid.chat.promise(
request.chatId,
'sendText',
);
let mentions: string[] | undefined;
if (request.mentions) {
mentions = await Promise.all(
request.mentions.map((mention) =>
this.hooks.wid.mention.promise(mention, 'sendText'),
),
);
}
const message = {
text: request.text,
mentions: request.mentions?.map(toJID),
mentions: mentions,
linkPreview: this.getLinkPreview(request),
};
const options: any = await this.getMessageOptions(request);
@@ -1120,8 +1142,8 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
}
@Activity()
public deleteMessage(chatId: string, messageId: string) {
const jid = toJID(this.ensureSuffix(chatId));
public async deleteMessage(chatId: string, messageId: string) {
const jid = await this.hooks.wid.chat.promise(chatId, 'deleteMessage');
const key = parseMessageIdSerialized(messageId);
const options = {
messageId: this.generateMessageID(),
@@ -1135,7 +1157,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
messageId: string,
request: EditMessageRequest,
) {
const jid = toJID(this.ensureSuffix(chatId));
const jid = await this.hooks.wid.chat.promise(chatId, 'editMessage');
const key = parseMessageIdSerialized(messageId);
const stored = await this.store
?.loadMessage(key.remoteJid, key.id)
@@ -1171,9 +1193,17 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
},
};
}
let mentions: string[] | undefined;
if (request.mentions) {
mentions = await Promise.all(
request.mentions.map((mention) =>
this.hooks.wid.mention.promise(mention, 'editMessage'),
),
);
}
let message: any = {
text: request.text,
mentions: request.mentions?.map(toJID),
mentions: mentions,
edit: key,
editedMessage: editedMessage,
linkPreview: this.getLinkPreview(request),
@@ -1191,7 +1221,10 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
@Activity()
async sendContactVCard(request: MessageContactVcardRequest) {
const chatId = toJID(this.ensureSuffix(request.chatId));
const chatId = await this.hooks.wid.chat.promise(
request.chatId,
'sendContactVCard',
);
const contacts = request.contacts.map((el) => ({ vcard: toVcardV3(el) }));
const options = await this.getMessageOptions(request);
const msg = { contacts: { contacts: contacts } };
@@ -1209,20 +1242,32 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
: 1,
};
const message = { poll: poll };
const remoteJid = toJID(request.chatId);
const remoteJid = await this.hooks.wid.chat.promise(
request.chatId,
'sendPoll',
);
const options = await this.getMessageOptions(request);
const result = await this.sock.sendMessage(remoteJid, message, options);
return this.toWAMessage(result);
return await this.toWAMessage(result);
}
@Activity()
async reply(request: MessageReplyRequest) {
const chatId = await this.hooks.wid.chat.promise(request.chatId, 'reply');
const options = await this.getMessageOptions(request);
let mentions: string[] | undefined;
if (request.mentions) {
mentions = await Promise.all(
request.mentions.map((mention) =>
this.hooks.wid.mention.promise(mention, 'reply'),
),
);
}
const message = {
text: request.text,
mentions: request.mentions?.map(toJID),
mentions: mentions,
};
return await this.sock.sendMessage(request.chatId, message, options);
return await this.sock.sendMessage(chatId, message, options);
}
@Activity()
@@ -1233,7 +1278,10 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
request.caption,
);
message.mimetype = message.mimetype || WAMimeType.IMAGE;
const chatId = toJID(this.ensureSuffix(request.chatId));
const chatId = await this.hooks.wid.chat.promise(
request.chatId,
'sendImage',
);
// Baileys' newsletter media path skips thumbnail and dimension computation.
// Pre-compute them so iOS renders the image with the correct aspect ratio.
if (isJidNewsletter(chatId)) {
@@ -1250,7 +1298,11 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
}
}
if (request.mentions?.length) {
message.mentions = request.mentions.map((mention) => toJID(mention));
message.mentions = await Promise.all(
request.mentions.map((mention) =>
this.hooks.wid.mention.promise(mention, 'sendImage'),
),
);
}
const options = await this.getMessageOptions(request);
return this.sock.sendMessage(chatId, message, options);
@@ -1267,9 +1319,16 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
message.mimetype = await detectMimetype(message['document']);
}
if (request.mentions?.length) {
message.mentions = request.mentions.map((mention) => toJID(mention));
message.mentions = await Promise.all(
request.mentions.map((mention) =>
this.hooks.wid.mention.promise(mention, 'sendFile'),
),
);
}
const chatId = toJID(this.ensureSuffix(request.chatId));
const chatId = await this.hooks.wid.chat.promise(
request.chatId,
'sendFile',
);
const options = await this.getMessageOptions(request);
return this.sock.sendMessage(chatId, message, options);
}
@@ -1282,7 +1341,10 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
message['audio'] = await this.mediaConverter.voice(message['audio']);
message.mimetype = WAMimeType.VOICE;
}
const chatId = toJID(this.ensureSuffix(request.chatId));
const chatId = await this.hooks.wid.chat.promise(
request.chatId,
'sendVoice',
);
const options = await this.getMessageOptions(request);
return this.sock.sendMessage(chatId, message, options);
}
@@ -1300,7 +1362,11 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
message.mimetype = WAMimeType.VIDEO;
}
if (request.mentions?.length) {
message.mentions = request.mentions.map((mention) => toJID(mention));
message.mentions = await Promise.all(
request.mentions.map((mention) =>
this.hooks.wid.mention.promise(mention, 'sendVideo'),
),
);
}
const duration = await esm.b
.getAudioDuration(message['video'])
@@ -1314,7 +1380,10 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
message.gifPlayback = true;
message.externalShareFullVideoDurationInSeconds = 0;
}
const chatId = toJID(this.ensureSuffix(request.chatId));
const chatId = await this.hooks.wid.chat.promise(
request.chatId,
'sendVideo',
);
const options = await this.getMessageOptions(request);
message.ptv = parseBool(request.asNote);
return this.sock.sendMessage(chatId, message, options);
@@ -1324,7 +1393,10 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
async sendLinkCustomPreview(
request: MessageLinkCustomPreviewRequest,
): Promise<any> {
const chatId = toJID(this.ensureSuffix(request.chatId));
const chatId = await this.hooks.wid.chat.promise(
request.chatId,
'sendLinkCustomPreview',
);
const options = await this.getMessageOptions(request);
const preview = request.preview;
const urlInfo = {
@@ -1440,7 +1512,10 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
@Activity()
async sendButtons(request: SendButtonsRequest) {
const chatId = toJID(this.ensureSuffix(request.chatId));
const chatId = await this.hooks.wid.chat.promise(
request.chatId,
'sendButtons',
);
const headerImage = await this.uploadMedia(request.headerImage, 'image');
return await sendButtonMessage(
this.sock,
@@ -1455,7 +1530,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
@Activity()
async sendList(request: SendListRequest): Promise<any> {
const jid = toJID(this.ensureSuffix(request.chatId));
const jid = await this.hooks.wid.chat.promise(request.chatId, 'sendList');
if (!isLidUser(jid) && !isPnUser(jid)) {
throw new UnprocessableEntityException(
`List message can only be sent to a direct message chat.`,
@@ -1475,7 +1550,10 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
@Activity()
async sendLocation(request: MessageLocationRequest) {
const chatId = toJID(this.ensureSuffix(request.chatId));
const chatId = await this.hooks.wid.chat.promise(
request.chatId,
'sendLocation',
);
const msg = {
location: {
name: request.title || null,
@@ -1496,20 +1574,26 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
`Message with id '${request.messageId}' not found`,
);
}
const chatId = toJID(this.ensureSuffix(request.chatId));
const chatId = await this.hooks.wid.chat.promise(
request.chatId,
'forwardMessage',
);
const message = {
forward: forwardMessage,
force: true,
};
const options = await this.getMessageOptions(request);
const result = await this.sock.sendMessage(chatId, message as any, options);
return this.toWAMessage(result);
return await this.toWAMessage(result);
}
@Activity()
async sendLinkPreview(request: MessageLinkPreviewRequest) {
const text = `${request.title}\n${request.url}`;
const chatId = toJID(this.ensureSuffix(request.chatId));
const chatId = await this.hooks.wid.chat.promise(
request.chatId,
'sendLinkPreview',
);
const msg = { text: text };
const options = await this.getMessageOptions(request);
return this.sock.sendMessage(chatId, msg, options);
@@ -1535,13 +1619,19 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
@Activity()
async startTyping(request: ChatRequest): Promise<void> {
const chatId = toJID(this.ensureSuffix(request.chatId));
const chatId = await this.hooks.wid.chat.promise(
request.chatId,
'startTyping',
);
await this.sock.sendPresenceUpdate('composing', chatId);
}
@Activity()
async stopTyping(request: ChatRequest) {
const chatId = toJID(this.ensureSuffix(request.chatId));
const chatId = await this.hooks.wid.chat.promise(
request.chatId,
'stopTyping',
);
return this.sock.sendPresenceUpdate('paused', chatId);
}
@@ -1552,8 +1642,9 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
) {
const pagination = query as PaginationParams;
const merge = query.merge ?? true;
const jid = await this.hooks.wid.chat.promise(chatId, 'getChatMessages');
const messages = await this.store.getMessagesByJid(
toJID(chatId),
jid,
filter,
pagination,
merge,
@@ -1588,11 +1679,8 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
): Promise<null | WAMessage> {
const key = parseMessageIdSerialized(messageId, true);
const merge = query.merge ?? true;
const message = await this.store.getMessageById(
toJID(chatId),
key.id,
merge,
);
const jid = await this.hooks.wid.chat.promise(chatId, 'getChatMessage');
const message = await this.store.getMessageById(jid, key.id, merge);
if (!message) return null;
const params = {
download: query.downloadMedia,
@@ -1608,7 +1696,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
messageId: string,
duration: PinDuration,
): Promise<boolean> {
const jid = toJID(chatId);
const jid = await this.hooks.wid.chat.promise(chatId, 'pinMessage');
const key = parseMessageIdSerialized(messageId);
await this.sock.sendMessage(jid, {
pin: key,
@@ -1623,7 +1711,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
chatId: string,
messageId: string,
): Promise<boolean> {
const jid = toJID(chatId);
const jid = await this.hooks.wid.chat.promise(chatId, 'unpinMessage');
const key = parseMessageIdSerialized(messageId);
await this.sock.sendMessage(jid, {
pin: key,
@@ -1668,6 +1756,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
@Activity()
async setStar(request: MessageStarRequest) {
const key = parseMessageIdSerialized(request.messageId);
const jid = await this.hooks.wid.chat.promise(request.chatId, 'setStar');
await this.sock.chatModify(
{
star: {
@@ -1675,7 +1764,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
star: request.star,
},
},
toJID(request.chatId),
jid,
);
}
@@ -1699,8 +1788,13 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
// Convert customer format IDs to JID format if filter is provided
let jidFilter;
if (filter?.ids && filter.ids.length > 0) {
const ids = await Promise.all(
filter.ids.map((id) =>
this.hooks.wid.chat.promise(id, 'getChatsOverview'),
),
);
jidFilter = {
ids: filter.ids.map((id) => toJID(id)),
ids: ids,
};
}
@@ -1756,7 +1850,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
chatId: string,
archive: boolean,
): Promise<any> {
const jid = toJID(chatId);
const jid = await this.hooks.wid.chat.promise(chatId, 'chatsPutArchive');
const messages = await this.store.getMessagesByJid(jid, {}, { limit: 1 });
return await this.sock.chatModify(
{ archive: archive, lastMessages: messages },
@@ -1776,7 +1870,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
@Activity()
public async chatsUnreadChat(chatId: string): Promise<any> {
const jid = toJID(chatId);
const jid = await this.hooks.wid.chat.promise(chatId, 'chatsUnreadChat');
const messages = await this.store.getMessagesByJid(jid, {}, { limit: 1 });
return await this.sock.chatModify(
{ markRead: false, lastMessages: messages },
@@ -1850,14 +1944,14 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
}
public async getChatLabels(chatId: string): Promise<Label[]> {
const jid = toJID(chatId);
const jid = await this.hooks.wid.chat.promise(chatId, 'getChatLabels');
const labels = await this.store.getChatLabels(jid);
return labels.map(this.toLabel);
}
@Activity()
public async putLabelsToChat(chatId: string, labels: LabelID[]) {
const jid = toJID(chatId);
const jid = await this.hooks.wid.chat.promise(chatId, 'putLabelsToChat');
const labelsIds = labels.map((label) => label.id);
const currentLabels = await this.store.getChatLabels(jid);
const currentLabelsIds = currentLabels.map((label) => label.id);
@@ -1899,7 +1993,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
@Activity()
public async upsertContact(chatId: string, body: ContactUpdateBody) {
const jid = toJID(chatId);
const jid = await this.hooks.wid.chat.promise(chatId, 'upsertContact');
let fullName = body.firstName;
if (body.lastName) {
fullName = `${body.firstName} ${body.lastName}`;
@@ -1920,7 +2014,10 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
}
async getContact(query: ContactQuery) {
const jid = toJID(query.contactId);
const jid = await this.hooks.wid.chat.promise(
query.contactId,
'getContact',
);
const contact = await this.store.getContactById(jid);
if (!contact) {
return null;
@@ -1935,7 +2032,10 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
@Activity()
public async fetchContactProfilePicture(id: string) {
const contact = this.ensureSuffix(id);
const contact = await this.hooks.wid.chat.promise(
id,
'fetchContactProfilePicture',
);
try {
const url = await this.sock.profilePictureUrl(contact, 'image');
return url;
@@ -1988,7 +2088,10 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
public async findLIDByPhoneNumber(
phoneNumber: string,
): Promise<LidToPhoneNumber> {
const pn = toJID(phoneNumber);
const pn = await this.hooks.wid.chat.promise(
phoneNumber,
'findLIDByPhoneNumber',
);
const lid = await this.store.findLidByPN(pn);
return {
lid: lid || null,
@@ -2215,7 +2318,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
case WAHAPresenceStatus.TYPING:
case WAHAPresenceStatus.RECORDING:
case WAHAPresenceStatus.PAUSED:
await this.maintainPresenceOnline();
await this.hooks.activity.promise('setPresence');
}
const enginePresence = ToEnginePresenceStatus[presence];
if (!enginePresence) {
@@ -2224,7 +2327,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
);
}
if (chatId) {
chatId = toJID(this.ensureSuffix(chatId));
chatId = await this.hooks.wid.chat.promise(chatId, 'setPresence');
}
await this.sock.sendPresenceUpdate(enginePresence, chatId);
this.presence = presence;
@@ -2240,7 +2343,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
}
public async getPresence(chatId: string): Promise<WAHAChatPresences> {
const jid = toJID(chatId);
const jid = await this.hooks.wid.chat.promise(chatId, 'getPresence');
await this.subscribePresence(jid);
if (!(jid in this.store.presences)) {
this.store.presences[jid] = {};
@@ -2251,8 +2354,8 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
}
@Activity()
public subscribePresence(id: string): Promise<void> {
const jid = toJID(id);
public async subscribePresence(id: string): Promise<void> {
const jid = await this.hooks.wid.chat.promise(id, 'subscribePresence');
return this.sock.presenceSubscribe(jid);
}
@@ -2430,7 +2533,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
protected prepareMessageIdForStatus(status: StatusRequest) {
if (status.id) {
this.saveSentMessageId(status.id);
this.hooks.message.sent.call(status.id);
return status.id;
}
return this.generateMessageID();
@@ -2439,7 +2542,11 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
protected async prepareJidsForStatus(contacts: string[]) {
let jids: string[];
if (contacts?.length > 0) {
jids = contacts.map(toJID);
jids = await Promise.all(
contacts.map((contact) =>
this.hooks.wid.chat.promise(contact, 'prepareJidsForStatus'),
),
);
} else {
jids = await this.fetchMyContactsJids();
}
@@ -2582,6 +2689,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
};
}
@Activity()
public async channelsList(query: ListChannelsQuery): Promise<Channel[]> {
const newsletters = await this.sock.newsletterSubscribed();
let channels = newsletters
@@ -2620,6 +2728,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
return channel;
}
@Activity()
public async channelsGetChannel(id: string) {
const newsletter = await this.sock.newsletterMetadata('jid', id);
return this.toChannel(toNewsletterMetadata(newsletter));
@@ -2708,7 +2817,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
proto.Message.ProtocolMessage.Type.REVOKE,
),
mergeMap(async (message): Promise<WAMessageRevokedBody> => {
const afterMessage = this.toWAMessage(message);
const afterMessage = await this.toWAMessage(message);
// Extract the revoked message ID from protocolMessage.key
const revokedMessageId = message.message.protocolMessage.key?.id;
return {
@@ -2729,7 +2838,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
IsSecretEncryptedMessageEdit(message.message),
),
mergeMap(async (message): Promise<WAMessageEditedBody> => {
const waMessage = this.toWAMessage(message);
const waMessage = await this.toWAMessage(message);
let body = '';
let editedMessageId: string | undefined;
if (IsEditedMessage(message.message)) {
@@ -2756,7 +2865,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
// Message Reactions
//
const messageReactions$ = messagesUpsert$.pipe(
map(this.processMessageReaction.bind(this)),
concatMap((message) => this.processMessageReaction(message)),
filter(Boolean),
);
this.events2.get(WAHAEvents.MESSAGE_REACTION).switch(messageReactions$);
@@ -3025,7 +3134,9 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
* END - Methods for API
*/
private processMessageReaction(message): WAMessageReaction | null {
private async processMessageReaction(
message,
): Promise<WAMessageReaction | null> {
if (!message) return null;
if (!message.message) return null;
if (!message.message.reactionMessage) return null;
@@ -3034,7 +3145,8 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
const fromToParticipant = getFromToParticipant(message.key);
const reactionMessage = message.message.reactionMessage;
const messageId = buildMessageId(reactionMessage.key);
const source = this.getMessageSource(message.key.id);
let source = await this.hooks.message.source.promise(message.key.id);
source = source ?? MessageSource.APP;
const reaction: WAMessageReaction = {
id: id,
timestamp: ensureNumber(message.messageTimestamp),
@@ -3248,7 +3360,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
return null;
}
// Convert
const wamessage = this.toWAMessageSafe(message);
const wamessage = await this.toWAMessageSafe(message);
if (!wamessage) {
return null;
}
@@ -3273,9 +3385,9 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
return wamessage;
}
protected toWAMessageSafe(message): WAMessage | null {
protected async toWAMessageSafe(message): Promise<WAMessage | null> {
try {
return this.toWAMessage(message);
return await this.toWAMessage(message);
} catch (error) {
this.logger.error('Failed to process incoming message');
this.logger.error(error);
@@ -3283,14 +3395,15 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
}
}
protected toWAMessage(message): WAMessage {
protected async toWAMessage(message): Promise<WAMessage> {
const fromToParticipant = getFromToParticipant(message.key);
const id = buildMessageId(message.key);
const body = extractBody(message.message);
const replyTo = this.extractReplyTo(message.message);
const ack = message.ack || StatusToAck(message.status);
const mediaContent = extractMediaContent(message.message);
const source = this.getMessageSource(message.key.id);
let source = await this.hooks.message.source.promise(message.key.id);
source = source ?? MessageSource.APP;
const waproto = message.message;
return {
id: id,
@@ -3546,7 +3659,10 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
chatId: string;
reply_to?: string;
}) {
const jid = toJID(request.chatId);
const jid = await this.hooks.wid.chat.promise(
request.chatId,
'getMessageOptions',
);
let quoted;
if (request.reply_to) {
@@ -3555,7 +3671,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
}
const chat = await this.store.getChat(jid);
const messageId = request.id ? request.id : this.generateMessageID();
this.saveSentMessageId(messageId);
this.hooks.message.sent.call(messageId);
return {
quoted: quoted,
ephemeralExpiration: chat?.ephemeralExpiration,
@@ -3581,7 +3697,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
protected generateMessageID() {
const id = generateMessageIDV2(this.sock.user?.id);
this.saveSentMessageId(id);
this.hooks.message.sent.call(id);
return id;
}
}
+154 -111
View File
@@ -130,6 +130,7 @@ import {
WAHAPresenceData,
} from '@waha/structures/presence.dto';
import {
MessageSource,
WALocation,
WAMessage,
WAMessageReaction,
@@ -167,6 +168,7 @@ import * as lodash from 'lodash';
import * as path from 'path';
import { Browser, ProtocolError } from 'puppeteer';
import {
concatMap,
filter,
fromEvent,
merge,
@@ -211,7 +213,6 @@ import {
normalizeJid,
toCusFormat,
} from '@waha/core/utils/jids';
import { Activity } from '@waha/core/abc/activity';
import { CallData } from '@waha/structures/calls.dto';
import { Jid } from '@waha/core/engines/const';
import {
@@ -222,6 +223,8 @@ import { removeSingletonFiles } from '@waha/core/utils/chrome';
import { killProcessesByPatterns } from '@waha/core/utils/processes';
import { IsChrome } from '@waha/version';
import { Activity } from '@waha/core/abc/session.hooks.activity';
export interface WebJSConfig {
webVersion?: string;
cacheType: 'local' | 'none';
@@ -479,7 +482,6 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
}
async stop() {
this.cleanupPresenceTimeout();
this.shouldRestart = false;
this.startDelayedJob.cancel();
await this.end();
@@ -513,7 +515,6 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
}
private async end() {
this.cleanupPresenceTimeout();
this.presence = null;
this.engineStateCheckDelayedJob.cancel();
this.whatsNewModalJob.stop();
@@ -894,6 +895,7 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
return await this.whatsapp.setProfilePicture(media);
}
@Activity()
protected async deleteProfilePicture(): Promise<boolean> {
return await this.whatsapp.deleteProfilePicture();
}
@@ -979,7 +981,8 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
*/
@Activity()
async rejectCall(from: string, id: string): Promise<void> {
const peerJid = normalizeJid(this.ensureSuffix(from));
const wid = await this.hooks.wid.chat.promise(from, 'rejectCall');
const peerJid = normalizeJid(wid);
const call = new CallInstance(this.whatsapp, null);
call.id = id;
call.from = peerJid;
@@ -988,13 +991,13 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
}
@Activity()
sendText(request: MessageTextRequest) {
const options = this.getMessageOptions(request);
return this.whatsapp.sendMessage(
this.ensureSuffix(request.chatId),
request.text,
options,
async sendText(request: MessageTextRequest) {
const options = await this.getMessageOptions(request);
const chatId = await this.hooks.wid.chat.promise(
request.chatId,
'sendText',
);
return this.whatsapp.sendMessage(chatId, request.text, options);
}
@Activity()
@@ -1020,9 +1023,12 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
@Activity()
async sendContactVCard(request: MessageContactVcardRequest) {
const chatId = this.ensureSuffix(request.chatId);
const chatId = await this.hooks.wid.chat.promise(
request.chatId,
'sendContactVCard',
);
const vcards = request.contacts.map((el) => toVcardV3(el as any));
const options = this.getMessageOptions(request);
const options = await this.getMessageOptions(request);
// Single vCard: pass raw vcard text as a message body.
// WEBJS will detect BEGIN:VCARD when parseVCards=true and send as a contact card.
@@ -1043,12 +1049,9 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
@Activity()
async reply(request: MessageReplyRequest) {
const options = this.getMessageOptions(request);
return this.whatsapp.sendMessage(
this.ensureSuffix(request.chatId),
request.text,
options,
);
const options = await this.getMessageOptions(request);
const chatId = await this.hooks.wid.chat.promise(request.chatId, 'reply');
return this.whatsapp.sendMessage(chatId, request.text, options);
}
@Activity()
@@ -1057,12 +1060,12 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
allowMultipleAnswers: request.poll.multipleAnswers,
messageSecret: undefined,
});
const options = this.getMessageOptions(request);
return this.whatsapp.sendMessage(
this.ensureSuffix(request.chatId),
poll,
options,
const options = await this.getMessageOptions(request);
const chatId = await this.hooks.wid.chat.promise(
request.chatId,
'sendPoll',
);
return this.whatsapp.sendMessage(chatId, poll, options);
}
@Activity()
@@ -1071,33 +1074,33 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
if (!media.mimetype) {
media.mimetype = await detectMimetype(Buffer.from(media.data, 'base64'));
}
let options = this.getMessageOptions(request);
let options = await this.getMessageOptions(request);
options = {
...options,
sendMediaAsDocument: true,
caption: request.caption,
};
return this.whatsapp.sendMessage(
this.ensureSuffix(request.chatId),
media,
options,
const chatId = await this.hooks.wid.chat.promise(
request.chatId,
'sendFile',
);
return this.whatsapp.sendMessage(chatId, media, options);
}
@Activity()
async sendImage(request: MessageImageRequest) {
const media = await this.fileToMedia(request.file);
media.mimetype = media.mimetype || WAMimeType.IMAGE;
let options = this.getMessageOptions(request);
let options = await this.getMessageOptions(request);
options = {
...options,
caption: request.caption,
};
return this.whatsapp.sendMessage(
this.ensureSuffix(request.chatId),
media,
options,
const chatId = await this.hooks.wid.chat.promise(
request.chatId,
'sendImage',
);
return this.whatsapp.sendMessage(chatId, media, options);
}
@Activity()
@@ -1107,16 +1110,16 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
await this.convertVoice(media);
}
media.mimetype = request.file.mimetype || WAMimeType.VOICE;
let options = this.getMessageOptions(request);
let options = await this.getMessageOptions(request);
options = {
...options,
sendAudioAsVoice: true,
};
return this.whatsapp.sendMessage(
this.ensureSuffix(request.chatId),
media,
options,
const chatId = await this.hooks.wid.chat.promise(
request.chatId,
'sendVoice',
);
return this.whatsapp.sendMessage(chatId, media, options);
}
@Activity()
@@ -1128,16 +1131,16 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
await this.convertVideo(media);
}
let options = this.getMessageOptions(request);
let options = await this.getMessageOptions(request);
options = {
...options,
caption: request.caption,
};
return this.whatsapp.sendMessage(
this.ensureSuffix(request.chatId),
media,
options,
const chatId = await this.hooks.wid.chat.promise(
request.chatId,
'sendVideo',
);
return this.whatsapp.sendMessage(chatId, media, options);
}
private async convertVideo(media: MessageMedia) {
@@ -1171,7 +1174,7 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
@Activity()
async sendButtonsReply(request: MessageButtonReply) {
const options = this.getMessageOptions(request);
const options = await this.getMessageOptions(request);
const extra: any = {
type: 'buttons_response',
kind: 'buttonsResponse',
@@ -1187,8 +1190,12 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
viewMode: 'VISIBLE',
};
options.extra = extra;
const chatId = await this.hooks.wid.chat.promise(
request.chatId,
'sendButtonsReply',
);
return this.whatsapp.sendMessage(
this.ensureSuffix(request.chatId),
chatId,
request.selectedDisplayText,
options,
);
@@ -1207,18 +1214,22 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
const location = new Location(request.latitude, request.longitude, {
name: request.title,
});
const options = this.getMessageOptions(request);
return this.whatsapp.sendMessage(
this.ensureSuffix(request.chatId),
location,
options,
const options = await this.getMessageOptions(request);
const chatId = await this.hooks.wid.chat.promise(
request.chatId,
'sendLocation',
);
return this.whatsapp.sendMessage(chatId, location, options);
}
@Activity()
async forwardMessage(request: MessageForwardRequest): Promise<WAMessage> {
const forwardMessage = this.recreateMessage(request.messageId);
const msg = await forwardMessage.forward(this.ensureSuffix(request.chatId));
const chatId = await this.hooks.wid.chat.promise(
request.chatId,
'forwardMessage',
);
const msg = await forwardMessage.forward(chatId);
// Return "sent: true" for now
// need to research how to get the data from WebJS
// @ts-ignore
@@ -1227,25 +1238,31 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
@Activity()
async sendSeen(request: SendSeenRequest) {
const chat: Chat = await this.whatsapp.getChatById(
this.ensureSuffix(request.chatId),
const chatId = await this.hooks.wid.chat.promise(
request.chatId,
'sendSeen',
);
const chat: Chat = await this.whatsapp.getChatById(chatId);
await chat.sendSeen();
}
@Activity()
async startTyping(request: ChatRequest): Promise<void> {
const chat: Chat = await this.whatsapp.getChatById(
this.ensureSuffix(request.chatId),
const chatId = await this.hooks.wid.chat.promise(
request.chatId,
'startTyping',
);
const chat: Chat = await this.whatsapp.getChatById(chatId);
await chat.sendStateTyping();
}
@Activity()
async stopTyping(request: ChatRequest) {
const chat: Chat = await this.whatsapp.getChatById(
this.ensureSuffix(request.chatId),
const chatId = await this.hooks.wid.chat.promise(
request.chatId,
'stopTyping',
);
const chat: Chat = await this.whatsapp.getChatById(chatId);
await chat.clearState();
}
@@ -1315,9 +1332,10 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
GetSerialized(chat.id),
false,
);
const lastMessage = chat.lastMessage
? this.toWAMessage(chat.lastMessage)
: null;
let lastMessage = null;
if (chat.lastMessage) {
lastMessage = await this.toWAMessage(chat.lastMessage);
}
return {
id: GetSerialized(chat.id),
name: chat.name || null,
@@ -1338,14 +1356,11 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
);
}
const id = await this.hooks.wid.chat.promise(chatId, 'getChatMessages');
// Test there's chat with id
await this.whatsapp.getChatById(this.ensureSuffix(chatId));
await this.whatsapp.getChatById(id);
const pagination: PaginationParams = query;
const messages = await this.whatsapp.getMessages(
this.ensureSuffix(chatId),
filter,
pagination,
);
const messages = await this.whatsapp.getMessages(id, filter, pagination);
const promises = [];
const params = {
download: query.downloadMedia,
@@ -1365,9 +1380,8 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
chatId: string,
request: ReadChatMessagesQuery,
): Promise<ReadChatMessagesResponse> {
const chat: Chat = await this.whatsapp.getChatById(
this.ensureSuffix(chatId),
);
const id = await this.hooks.wid.chat.promise(chatId, 'readChatMessages');
const chat: Chat = await this.whatsapp.getChatById(id);
await chat.sendSeen();
return { ids: null };
}
@@ -1377,7 +1391,7 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
messageId: string,
query: GetChatMessageQuery,
): Promise<null | WAMessage> {
chatId = this.ensureSuffix(chatId);
chatId = await this.hooks.wid.chat.promise(chatId, 'getChatMessage');
// WEBJS waits the serializer messageId
// {fromMe}_{chatId}_{id}[_{participant}]
@@ -1464,7 +1478,8 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
@Activity()
async deleteChat(chatId) {
const chat = await this.whatsapp.getChatById(this.ensureSuffix(chatId));
const id = await this.hooks.wid.chat.promise(chatId, 'deleteChat');
const chat = await this.whatsapp.getChatById(id);
return chat.delete();
}
@@ -1475,20 +1490,20 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
}
@Activity()
public chatsArchiveChat(chatId: string): Promise<any> {
const id = this.ensureSuffix(chatId);
public async chatsArchiveChat(chatId: string): Promise<any> {
const id = await this.hooks.wid.chat.promise(chatId, 'chatsArchiveChat');
return this.whatsapp.archiveChat(id);
}
@Activity()
public chatsUnarchiveChat(chatId: string): Promise<any> {
const id = this.ensureSuffix(chatId);
public async chatsUnarchiveChat(chatId: string): Promise<any> {
const id = await this.hooks.wid.chat.promise(chatId, 'chatsUnarchiveChat');
return this.whatsapp.unarchiveChat(id);
}
@Activity()
public chatsUnreadChat(chatId: string): Promise<any> {
const id = this.ensureSuffix(chatId);
public async chatsUnreadChat(chatId: string): Promise<any> {
const id = await this.hooks.wid.chat.promise(chatId, 'chatsUnreadChat');
return this.whatsapp.markChatUnread(id);
}
@@ -1529,7 +1544,7 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
}
public async getChatLabels(chatId: string): Promise<Label[]> {
const id = this.ensureSuffix(chatId);
const id = await this.hooks.wid.chat.promise(chatId, 'getChatLabels');
const labels = await this.whatsapp.getChatLabels(id);
return labels.map(this.toLabel);
}
@@ -1537,7 +1552,8 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
@Activity()
public async putLabelsToChat(chatId: string, labels: LabelID[]) {
const labelIds = labels.map((label) => label.id);
const chatIds = [this.ensureSuffix(chatId)];
const id = await this.hooks.wid.chat.promise(chatId, 'putLabelsToChat');
const chatIds = [id];
await this.whatsapp.addOrRemoveLabels(labelIds, chatIds);
}
@@ -1565,10 +1581,12 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
);
}
getContact(query: ContactQuery) {
return this.whatsapp
.getContactById(this.ensureSuffix(query.contactId))
.then(this.toWAContact);
async getContact(query: ContactQuery) {
const contactId = await this.hooks.wid.chat.promise(
query.contactId,
'getContact',
);
return this.whatsapp.getContactById(contactId).then(this.toWAContact);
}
async getContacts(pagination: PaginationParams) {
@@ -1579,32 +1597,42 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
}
public async getContactAbout(query: ContactQuery) {
const contact = await this.whatsapp.getContactById(
this.ensureSuffix(query.contactId),
const contactId = await this.hooks.wid.chat.promise(
query.contactId,
'getContactAbout',
);
const contact = await this.whatsapp.getContactById(contactId);
return { about: await contact.getAbout() };
}
@Activity()
public async fetchContactProfilePicture(id: string) {
const contact = await this.whatsapp.getContactById(this.ensureSuffix(id));
const contactId = await this.hooks.wid.chat.promise(
id,
'fetchContactProfilePicture',
);
const contact = await this.whatsapp.getContactById(contactId);
const url = await contact.getProfilePicUrl();
return url;
}
@Activity()
public async blockContact(request: ContactRequest) {
const contact = await this.whatsapp.getContactById(
this.ensureSuffix(request.contactId),
const contactId = await this.hooks.wid.chat.promise(
request.contactId,
'blockContact',
);
const contact = await this.whatsapp.getContactById(contactId);
await contact.block();
}
@Activity()
public async unblockContact(request: ContactRequest) {
const contact = await this.whatsapp.getContactById(
this.ensureSuffix(request.contactId),
const contactId = await this.hooks.wid.chat.promise(
request.contactId,
'unblockContact',
);
const contact = await this.whatsapp.getContactById(contactId);
await contact.unblock();
}
@@ -2125,17 +2153,17 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
await this.whatsapp.sendPresenceUnavailable();
break;
case WAHAPresenceStatus.TYPING:
await this.maintainPresenceOnline();
await this.hooks.activity.promise('setPresence');
chat = await this.whatsapp.getChatById(chatId);
await chat.sendStateTyping();
break;
case WAHAPresenceStatus.RECORDING:
await this.maintainPresenceOnline();
await this.hooks.activity.promise('setPresence');
chat = await this.whatsapp.getChatById(chatId);
await chat.sendStateRecording();
break;
case WAHAPresenceStatus.PAUSED:
await this.maintainPresenceOnline();
await this.hooks.activity.promise('setPresence');
chat = await this.whatsapp.getChatById(chatId);
await chat.clearState();
break;
@@ -2259,6 +2287,7 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
return this.whatsapp.sendMessage(Jid.BROADCAST, status.text, options);
}
@Activity()
public async deleteStatus(request: DeleteStatusRequest) {
this.checkStatusRequest(request);
@@ -2275,7 +2304,7 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
subscribeEngineEvents2() {
// Save sent message in cache
this.whatsapp.events.on('message.id', (data) => {
this.saveSentMessageId(data.id);
this.hooks.message.sent.call(data.id);
});
//
@@ -2349,11 +2378,15 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
filter((evt: any) =>
this.jids.include(evt?.after?.id?.remote || evt?.before?.id?.remote),
),
map((event): WAMessageRevokedBody => {
const afterMessage = event.after ? this.toWAMessage(event.after) : null;
const beforeMessage = event.before
? this.toWAMessage(event.before)
: null;
concatMap(async (event): Promise<WAMessageRevokedBody> => {
let afterMessage = null;
if (event.after) {
afterMessage = await this.toWAMessage(event.after);
}
let beforeMessage = null;
if (event.before) {
beforeMessage = await this.toWAMessage(event.before);
}
// Extract the revoked message ID from the protocolMessageKey.id field
const revokedMessageId = afterMessage?._data?.protocolMessageKey?.id;
return {
@@ -2368,7 +2401,7 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
const messageReaction$ = fromEvent(this.whatsapp, 'message_reaction');
const messagesReaction$ = messageReaction$.pipe(
filter((reaction: Reaction) => this.jids.include(reaction?.id?.remote)),
map(this.processMessageReaction.bind(this)),
concatMap((reaction: Reaction) => this.processMessageReaction(reaction)),
filter(Boolean),
);
this.events2.get(WAHAEvents.MESSAGE_REACTION).switch(messagesReaction$);
@@ -2382,8 +2415,8 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
);
const messagesEdit$ = messageEdit$.pipe(
filter((event: any) => this.jids.include(event?.message?.id?.remote)),
map((event): WAMessageEditedBody => {
const message = this.toWAMessage(event.message);
concatMap(async (event): Promise<WAMessageEditedBody> => {
const message = await this.toWAMessage(event.message);
return {
...message,
body: event.newBody,
@@ -2410,7 +2443,7 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
);
const messagesAckDM$ = messageAckWEBJS$.pipe(
map((event) => event.message),
map<any, WAMessage>(this.toWAMessage.bind(this)),
concatMap((message: any) => this.toWAMessage(message)),
filter((ack) => !isJidGroup(ack.to) && !isJidStatusBroadcast(ack.to)),
filter((ack) => this.jids.include(ack.to)),
);
@@ -2568,7 +2601,7 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
options: MediaDownloadOptions,
) {
// Convert
const wamessage = this.toWAMessage(message);
const wamessage = await this.toWAMessage(message);
// Media
const media = await this.downloadMediaSafe(message, options);
wamessage.media = media;
@@ -2601,7 +2634,9 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
};
}
private processMessageReaction(reaction: Reaction): WAMessageReaction {
private async processMessageReaction(
reaction: Reaction,
): Promise<WAMessageReaction> {
if (this.lastQRDate) {
// If it's timestamp before last qr - ignore it
// Fixes: https://github.com/devlikeapro/waha/issues/494
@@ -2618,7 +2653,8 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
return null;
}
const source = this.getMessageSource(reaction.id.id);
let source = await this.hooks.message.source.promise(reaction.id.id);
source = source ?? MessageSource.APP;
return {
id: GetSerialized(reaction.id),
from: normalizeJid(reaction.senderId),
@@ -2712,9 +2748,10 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
return acks;
}
protected toWAMessage(message: Message): WAMessage {
protected async toWAMessage(message: Message): Promise<WAMessage> {
const replyTo = this.extractReplyTo(message);
const source = this.getMessageSource(message.id.id);
let source = await this.hooks.message.source.promise(message.id.id);
source = source ?? MessageSource.APP;
const key = parseMessageIdSerialized(GetSerialized(message.id));
// @ts-ignore
return {
@@ -2816,9 +2853,15 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
return null;
}
protected getMessageOptions(request: any): any {
let mentions = request.mentions;
mentions = mentions ? mentions.map(this.ensureSuffix) : undefined;
protected async getMessageOptions(request: any): Promise<any> {
let mentions: string[] | undefined = request.mentions;
if (mentions) {
mentions = await Promise.all(
mentions.map((mention) =>
this.hooks.wid.mention.promise(mention, 'getMessageOptions'),
),
);
}
const quotedMessageId = request.reply_to || request.replyTo;
File diff suppressed because it is too large. Load diff
-12
View File
@@ -1,19 +1,7 @@
import { parseBool } from '@waha/helpers';
//
// Presence
//
// Automatically mark session as ONLINE on any messages activity
export const PRESENCE_AUTO_ONLINE = process.env.WAHA_PRESENCE_AUTO_ONLINE
? parseBool(process.env.WAHA_PRESENCE_AUTO_ONLINE)
: true;
// Duration (in seconds) to keep session ONLINE after activity
// 25 seconds is default web timeout with no activity
export const PRESENCE_AUTO_ONLINE_DURATION_SECONDS =
parseInt(process.env.WAHA_PRESENCE_AUTO_ONLINE_DURATION_SECONDS) || 25;
//
// Local - sqlite3 engine
//
+70 -26
View File
@@ -19,7 +19,6 @@ import { WhatsappSessionNoWebCore } from '@waha/core/engines/noweb/session.noweb
import { WhatsappSessionWPPCore } from '@waha/core/engines/wpp/session.wpp.core';
import { WhatsappSessionWebJSCore } from '@waha/core/engines/webjs/session.webjs.core';
import { getProxyConfig } from '@waha/core/helpers.proxy';
import { WebhookConductor } from '@waha/core/integrations/webhooks/WebhookConductor';
import { MediaManager } from '@waha/core/media/MediaManager';
import { MediaStorageFactory } from '@waha/core/media/MediaStorageFactory';
import { LocalSessionAuthRepository } from '@waha/core/storage/LocalSessionAuthRepository';
@@ -70,8 +69,21 @@ import {
} from '../structures/sessions.dto';
import { WebhookConfig } from '../structures/webhooks.config.dto';
import { populateSessionInfo, SessionManager } from './abc/manager.abc';
import type { SessionPlugin } from '@waha/core/abc/session.plugin';
import { RegisterPluginEvents } from '@waha/core/abc/session.plugin.events';
import { RegisterPluginHooks } from '@waha/core/abc/session.plugin.hooks';
import { SessionParams, WhatsappSession } from './abc/session.abc';
import { EngineConfigService } from './config/EngineConfigService';
import { WidEnsureSuffixPlugin } from '@waha/plugins/WidEnsureSuffixPlugin';
import { MaintainOnlineStatusPlugin } from '@waha/plugins/MaintainOnlineStatusPlugin';
import { MessageSourceCachePlugin } from '@waha/plugins/MessageSourceCachePlugin';
import { SessionRuntimeInfoPlugin } from '@waha/plugins/SessionRuntimeInfoPlugin';
import { WebhookPlugin } from '@waha/plugins/WebhookPlugin';
import {
PRESENCE_AUTO_ONLINE,
PRESENCE_AUTO_ONLINE_DURATION_SECONDS,
} from '@waha/plugins/MaintainOnlineStatusPlugin.env';
const ALL = '*';
@@ -348,7 +360,6 @@ export class SessionManagerCore
storage,
loggerBuilder.child({ name: 'MediaManager' }),
);
const webhook = new WebhookConductor(loggerBuilder);
const proxyConfig = this.getProxyConfig(name, config);
const sessionConfig: SessionParams = {
name,
@@ -376,11 +387,48 @@ export class SessionManagerCore
this.sessions[name] = session;
this.updateSessions();
// configure webhooks
const webhooks = this.getWebhooks(config);
webhook.configure(session, webhooks);
// Apps
// Plugins
{
const logger = loggerBuilder.child({
plugin: SessionRuntimeInfoPlugin.name,
});
session.plugins[SessionRuntimeInfoPlugin.name] =
new SessionRuntimeInfoPlugin(session, logger);
}
{
const webhooks = this.getWebhooks(config);
const logger = loggerBuilder.child({ plugin: WebhookPlugin.name });
session.plugins[WebhookPlugin.name] = new WebhookPlugin(session, logger, {
webhooks: webhooks,
});
}
if (PRESENCE_AUTO_ONLINE) {
const config = {
duration: PRESENCE_AUTO_ONLINE_DURATION_SECONDS * 1000,
};
const logger = loggerBuilder.child({
plugin: MaintainOnlineStatusPlugin.name,
});
session.plugins[MaintainOnlineStatusPlugin.name] =
new MaintainOnlineStatusPlugin(session, logger, config);
}
{
const logger = loggerBuilder.child({
plugin: MessageSourceCachePlugin.name,
});
session.plugins[MessageSourceCachePlugin.name] =
new MessageSourceCachePlugin(session, logger);
}
{
const logger = loggerBuilder.child({
plugin: WidEnsureSuffixPlugin.name,
});
session.plugins[WidEnsureSuffixPlugin.name] = new WidEnsureSuffixPlugin(
session,
logger,
);
}
// Apps (may contribute their own plugins to the session)
try {
await this.appsService.beforeSessionStart(session, this.store);
} catch (e) {
@@ -388,6 +436,12 @@ export class SessionManagerCore
session.status = WAHASessionStatus.FAILED;
}
const plugins: SessionPlugin<any>[] = Object.values(session.plugins);
for (const plugin of plugins) {
RegisterPluginHooks(plugin);
RegisterPluginEvents(plugin);
}
// start session
if (session.status !== WAHASessionStatus.FAILED) {
await session.start();
@@ -514,27 +568,17 @@ export class SessionManagerCore
/**
* Get all runtime sessions
*/
private getRuntimeSessions(name: string = null): SessionInfo[] {
private async getRuntimeSessions(
name: string = null,
): Promise<SessionInfo[]> {
let names = Object.keys(this.sessions);
if (name) {
names = names.filter((n) => n === name);
}
const sessions = names.map((sessionName) => {
const status = this.sessions[sessionName].status;
const sessionConfig = this.sessions[sessionName].sessionConfig;
const me = this.sessions[sessionName].getSessionMeInfo();
return {
name: sessionName,
status: status,
config: sessionConfig,
me: me,
presence: this.sessions[sessionName].presence,
timestamps: {
activity: this.sessions[sessionName].getLastActivityTimestamp(),
},
};
});
return sessions;
const sessions = names.map((sessionName) =>
this.sessions[sessionName].getSessionInfo(),
);
return await Promise.all(sessions);
}
/**
@@ -568,7 +612,7 @@ export class SessionManagerCore
}
async getSessions(all: boolean): Promise<SessionInfo[]> {
const runtimeSessions = this.getRuntimeSessions();
const runtimeSessions = await this.getRuntimeSessions();
let offlineSessions: SessionInfo[] = [];
if (all) {
offlineSessions = await this.getOfflineSessions();
@@ -595,7 +639,7 @@ export class SessionManagerCore
let session: SessionDetailedInfo = null;
// Try to find session in runtime sessions
const runtimeSessions = this.getRuntimeSessions(name);
const runtimeSessions = await this.getRuntimeSessions(name);
if (runtimeSessions.length === 1) {
session = runtimeSessions[0];
}
@@ -0,0 +1,10 @@
import { parseBool } from '@waha/helpers';
// Automatically mark session as ONLINE on any messages activity
export const PRESENCE_AUTO_ONLINE = process.env.WAHA_PRESENCE_AUTO_ONLINE
? parseBool(process.env.WAHA_PRESENCE_AUTO_ONLINE)
: true;
// Duration (in seconds) to keep session ONLINE after activity
// Web goes "unavailable" after ~90s of no user activity (60s idle grace + 30s timer)
export const PRESENCE_AUTO_ONLINE_DURATION_SECONDS =
parseInt(process.env.WAHA_PRESENCE_AUTO_ONLINE_DURATION_SECONDS) || 90;
+104
View File
@@ -0,0 +1,104 @@
import { PluginEvent } from '@waha/core/abc/session.plugin.events';
import { PluginHook } from '@waha/core/abc/session.plugin.hooks';
import { SessionPlugin } from '@waha/core/abc/session.plugin';
import {
WAHAEvents,
WAHAPresenceStatus,
WAHASessionStatus,
} from '@waha/structures/enums.dto';
import { SessionInfo } from '@waha/structures/sessions.dto';
import { WASessionStatusBody } from '@waha/structures/webhooks.dto';
import * as lodash from 'lodash';
import { Observable } from 'rxjs';
export class MaintainOnlineStatusPluginConfig {
duration: number; // in milliseconds
}
/**
* Sets the session presence to ONLINE before any activity goes to the server
* then auto-sets OFFLINE after `duration` ms without activity.
*/
export class MaintainOnlineStatusPlugin extends SessionPlugin<MaintainOnlineStatusPluginConfig> {
private lastOnlineTimestamp?: number;
private offlineTimeout?: ReturnType<typeof setTimeout>;
@PluginHook((hooks) => hooks.session.info)
attachActivityTimestamp(info: SessionInfo): SessionInfo {
return lodash.merge(info, {
timestamps: {
activity: this.getLastOnlineTimestamp() ?? null,
},
});
}
@PluginEvent(WAHAEvents.SESSION_STATUS)
onSessionStatus(events$: Observable<WASessionStatusBody>) {
events$.subscribe({
next: (body: WASessionStatusBody) => {
if (body.status !== WAHASessionStatus.WORKING) {
this.cleanupPresenceTimeout();
}
},
complete: () => {
this.cleanupPresenceTimeout();
},
});
}
/**
* Returns the timestamp of the last "activity" in the session
* @returns Timestamp in milliseconds or undefined if there was never any activity
*/
public getLastOnlineTimestamp(): number | undefined {
return this.lastOnlineTimestamp;
}
@PluginHook((hooks) => hooks.activity)
async maintainPresenceOnline(): Promise<void> {
if (this.session.status !== WAHASessionStatus.WORKING) {
return;
}
this.lastOnlineTimestamp = Date.now();
// If not ONLINE yet, send ONLINE
if (this.session.presence !== WAHAPresenceStatus.ONLINE) {
try {
// Force set ONLINE in case of many requests comes at the same time
// So we'll set ONLINE exactly once
this.session.presence = WAHAPresenceStatus.ONLINE;
await this.session.setPresence(WAHAPresenceStatus.ONLINE);
this.logger.debug('Set presence to ONLINE due to activity');
} catch (error) {
this.logger.debug('Failed to set presence ONLINE', error);
return;
}
}
// Cancel the previous timeout (if exists)
this.cleanupPresenceTimeout();
// Schedule to go back OFFLINE after timeout without activity
this.offlineTimeout = setTimeout(async () => {
try {
const working = this.session.status === WAHASessionStatus.WORKING;
const online = this.session.presence === WAHAPresenceStatus.ONLINE;
if (!working || !online) {
// Nothing to do
return;
}
await this.session.setPresence(WAHAPresenceStatus.OFFLINE);
this.logger.debug(
'Auto-set presence to OFFLINE after time without activity',
);
} catch (error) {
this.session.presence = WAHAPresenceStatus.OFFLINE;
this.logger.debug('Failed to set presence OFFLINE', error);
}
this.cleanupPresenceTimeout();
}, this.config.duration);
}
protected cleanupPresenceTimeout() {
clearTimeout(this.offlineTimeout);
this.offlineTimeout = null;
}
}
+54
View File
@@ -0,0 +1,54 @@
import { PluginEvent } from '@waha/core/abc/session.plugin.events';
import { PluginHook } from '@waha/core/abc/session.plugin.hooks';
import { SessionPlugin } from '@waha/core/abc/session.plugin';
import { WAHAEvents } from '@waha/structures/enums.dto';
import { MessageSource } from '@waha/structures/responses.dto';
import * as NodeCache from 'node-cache';
import { Observable } from 'rxjs';
/**
* Remembers ids of messages sent via API ("message.sent" hook)
* and answers "message.source" lookups:
* - MessageSource.API if the id is in the cache
* - undefined otherwise, so other taps (or the caller's APP fallback) decide
*/
export class MessageSourceCachePlugin extends SessionPlugin {
private sentMessageIds: NodeCache = new NodeCache({
stdTTL: 10 * 60, // 10 minutes
});
@PluginEvent(WAHAEvents.SESSION_STATUS)
onSessionStatus(events$: Observable<any>) {
events$.subscribe({
complete: () => {
this.close();
},
});
}
@PluginHook((hooks) => hooks.message.sent)
saveSentMessageId(id: string) {
if (!id) {
return;
}
this.sentMessageIds.set(id, true);
}
@PluginHook((hooks) => hooks.message.source)
getMessageSource(id: string): MessageSource | undefined {
if (!id) {
return undefined;
}
const api = this.sentMessageIds.has(id);
if (api) {
return MessageSource.API;
}
return undefined;
}
private close() {
this.logger.debug('Closing sent message ids cache');
this.sentMessageIds.flushAll();
this.sentMessageIds.close();
}
}
+24
View File
@@ -0,0 +1,24 @@
import { Stage } from '@waha/core/abc/session.hooks';
import { PluginHook } from '@waha/core/abc/session.plugin.hooks';
import { SessionPlugin } from '@waha/core/abc/session.plugin';
import { SessionInfo } from '@waha/structures/sessions.dto';
import * as lodash from 'lodash';
/**
* Provides the base runtime session info.
*/
export class SessionRuntimeInfoPlugin extends SessionPlugin {
@PluginHook((hooks) => hooks.session.info, { stage: Stage.FIRST })
getSessionInfo(info: SessionInfo): SessionInfo {
return lodash.merge(info, {
name: this.session.name,
status: this.session.status,
config: this.session.sessionConfig,
me: this.session.getSessionMeInfo(),
presence: this.session.presence,
timestamps: {
activity: null,
},
});
}
}
File renamed without changes.
@@ -1,34 +1,37 @@
import { populateSessionInfo } from '@waha/core/abc/manager.abc';
import { WhatsappSession } from '@waha/core/abc/session.abc';
import { WebhookSender } from '@waha/core/integrations/webhooks/WebhookSender';
import { SessionPlugin } from '@waha/core/abc/session.plugin';
import { WebhookSender } from '@waha/plugins/WebhookPlugin.sender';
import { WAHAEvents, WAHAEventsWild } from '@waha/structures/enums.dto';
import { WebhookConfig } from '@waha/structures/webhooks.config.dto';
import { EventWildUnmask } from '@waha/utils/events';
import { LoggerBuilder } from '@waha/utils/logging';
import { Logger } from 'pino';
export class WebhookConductor {
private logger: Logger;
export class WebhookPluginConfig {
webhooks: WebhookConfig[];
}
/**
* Sends session events to the configured webhooks (per session and global ones)
*/
export class WebhookPlugin extends SessionPlugin<WebhookPluginConfig> {
private eventUnmask = new EventWildUnmask(WAHAEvents, WAHAEventsWild);
constructor(protected loggerBuilder: LoggerBuilder) {
this.logger = loggerBuilder.child({ name: WebhookConductor.name });
}
protected buildSender(webhookConfig: WebhookConfig): WebhookSender {
return new WebhookSender(this.loggerBuilder, webhookConfig);
constructor(
session: WhatsappSession,
logger: Logger,
config: WebhookPluginConfig,
) {
super(session, logger, config);
for (const webhookConfig of config.webhooks) {
this.configureSingleWebhook(session, webhookConfig);
}
}
private getSuitableEvents(events: WAHAEvents[] | string[]): WAHAEvents[] {
return this.eventUnmask.unmask(events);
}
public configure(session: WhatsappSession, webhooks: WebhookConfig[]) {
for (const webhookConfig of webhooks) {
this.configureSingleWebhook(session, webhookConfig);
}
}
private configureSingleWebhook(
session: WhatsappSession,
webhook: WebhookConfig,
@@ -40,7 +43,7 @@ export class WebhookConductor {
const url = webhook.url;
this.logger.info(`Configuring webhooks for ${url}...`);
const events = this.getSuitableEvents(webhook.events);
const sender = this.buildSender(webhook);
const sender = new WebhookSender(this.logger, webhook);
for (const event of events) {
const obs$ = session.getEventObservable(event);
obs$.subscribe((payload) => {
+15
View File
@@ -0,0 +1,15 @@
import { ensureSuffix } from '@waha/core/abc/session.abc';
import { Stage } from '@waha/core/abc/session.hooks';
import { PluginHook } from '@waha/core/abc/session.plugin.hooks';
import { SessionPlugin } from '@waha/core/abc/session.plugin';
/**
* Ensures the wid has a suffix - 11111111111 => 11111111111@c.us
*/
export class WidEnsureSuffixPlugin extends SessionPlugin {
@PluginHook((hooks) => hooks.wid.chat, { stage: Stage.FIRST })
@PluginHook((hooks) => hooks.wid.mention, { stage: Stage.FIRST })
ensureWidSuffix(wid: string): string {
return ensureSuffix(wid);
}
}
+16
View File
@@ -0,0 +1,16 @@
import { Stage } from '@waha/core/abc/session.hooks';
import { PluginHook } from '@waha/core/abc/session.plugin.hooks';
import { SessionPlugin } from '@waha/core/abc/session.plugin';
import { normalizeJid, toJID } from '@waha/core/utils/jids';
/**
* Converts the wid to the engine-native JID format -
* 11111111111@c.us => 11111111111@s.whatsapp.net
*/
export class WidToJIDPlugin extends SessionPlugin {
@PluginHook((hooks) => hooks.wid.chat, { stage: Stage.LAST })
@PluginHook((hooks) => hooks.wid.mention, { stage: Stage.LAST })
toJID(wid: string): string {
return normalizeJid(toJID(wid));
}
}
+8
View File
@@ -13873,6 +13873,13 @@ __metadata:
languageName: node
linkType: hard
"tapable@npm:^2.2.2":
version: 2.3.3
resolution: "tapable@npm:2.3.3"
checksum: 10/21fb64a7ae1a0e11d855a6c33a22ae5ecf7e2f23170c942da673b44bf4c3aae8aa52451ef2792d0ce36c7feca13dceafa4f135105d66fc06912632488c0913fd
languageName: node
linkType: hard
"tar-fs@npm:^2.0.0":
version: 2.1.3
resolution: "tar-fs@npm:2.1.3"
@@ -14748,6 +14755,7 @@ __metadata:
sqlite3: "npm:^5.1.7"
supertest: "npm:^4.0.2"
swagger-ui-express: "npm:^4.1.4"
tapable: "npm:^2.2.2"
ts-jest: "npm:^29.1.3"
ts-loader: "npm:^6.2.1"
ts-node: "npm:^10.9.2"