feat: granular control for media - api, events, mimetypes

closes #2211

Co-authored-by: Elimeshi1 <eliyaume@edu.hac.ac.il>
This commit is contained in:
devlikeproandElimeshi1 committed 2026-08-31 10:44:22 +07:00
1 parent cc12380e48
commit 99e4dfcfb5
17 files changed
+399 -163

No files matched your search

+10
View File
@@ -64,7 +64,17 @@ WAHA_MEDIA_STORAGE=LOCAL
WHATSAPP_FILES_LIFETIME=0
WHATSAPP_FILES_FOLDER=/app/.media
#
# Media download settings
#
# Events (webhooks, websockets)
# WAHA_EVENTS_DOWNLOAD_MEDIA=true
# WAHA_EVENTS_DOWNLOAD_MEDIA_MIMETYPES=image/jpeg,image/png
# API default behavior (GET /api/{session}/chats/{chatId}/messages, etc.)
# WAHA_API_DOWNLOAD_MEDIA=true
# WAHA_API_DOWNLOAD_MEDIA_MIMETYPES=image/jpeg,image/png
# DEPRECATED - use WAHA_EVENTS_* and WAHA_API_* variables above
# WHATSAPP_DOWNLOAD_MEDIA=true
# WHATSAPP_FILES_MIMETYPES=image/jpeg,image/png
+10
View File
@@ -60,7 +60,17 @@ WAHA_MEDIA_STORAGE=LOCAL
WHATSAPP_FILES_LIFETIME=1800
WHATSAPP_FILES_FOLDER=/app/.media
#
# Media download settings
#
# Events (webhooks, websockets)
# WAHA_EVENTS_DOWNLOAD_MEDIA=true
# WAHA_EVENTS_DOWNLOAD_MEDIA_MIMETYPES=image/jpeg,image/png
# API default behavior (GET /api/{session}/chats/{chatId}/messages, etc.)
# WAHA_API_DOWNLOAD_MEDIA=true
# WAHA_API_DOWNLOAD_MEDIA_MIMETYPES=image/jpeg,image/png
# DEPRECATED - use WAHA_EVENTS_* and WAHA_API_* variables above
# WHATSAPP_DOWNLOAD_MEDIA=true
# WHATSAPP_FILES_MIMETYPES=image/jpeg,image/png
+4 -1
View File
@@ -107,7 +107,10 @@ export class ChannelTools extends McpController {
return this.textRequest({
method: 'GET',
url: `/api/${session}/channels/${id}/messages/preview`,
params: query,
params: {
...query,
downloadMediaMimetypes: query.downloadMediaMimetypes?.join(','),
},
});
}
+10 -3
View File
@@ -145,7 +145,8 @@ export class ChatTools extends McpController {
description:
'Get messages in a chat. ' +
'To retrieve all messages, paginate by incrementing the offset by limit until the returned array is empty or shorter than the limit. ' +
'To fetch media for a specific message, use chats-get-message with that message id and downloadMedia=true instead of fetching it here.',
'To fetch media for a specific message, use chats-get-message with that message id and downloadMedia=true instead of fetching it here. ' +
'Use downloadMediaMimetypes to download only media with the given mimetypes (prefix match).',
inputSchema: ChatMessagesInput,
annotations: {
readOnlyHint: true,
@@ -162,7 +163,10 @@ export class ChatTools extends McpController {
const response = await this.request({
method: 'GET',
url: `/api/${session}/chats/${chatId}/messages`,
params: query,
params: {
...query,
downloadMediaMimetypes: query.downloadMediaMimetypes?.join(','),
},
});
let messages = response.data;
if (!_data && Array.isArray(messages)) {
@@ -229,7 +233,10 @@ export class ChatTools extends McpController {
const result = await this.textRequest({
method: 'GET',
url: `/api/${session}/chats/${chatId}/messages/${messageId}`,
params: query,
params: {
...query,
downloadMediaMimetypes: query.downloadMediaMimetypes?.join(','),
},
});
if (query.downloadMedia) {
const mediaKey = await this.mediaApiKey(session);
+44 -9
View File
@@ -1,7 +1,9 @@
import { Injectable, Logger, OnApplicationBootstrap } from '@nestjs/common';
import { ConfigService } from '@nestjs/config';
import { GlobalWebhookConfigConfig } from '@waha/core/config/GlobalWebhookConfig';
import { MediaDownloadOptions } from '@waha/core/media/IMediaManager';
import { IgnoreJidConfig } from '@waha/core/utils/jids';
import * as lodash from 'lodash';
import { parseBool } from './helpers';
import { WebhookConfig } from './structures/webhooks.config.dto';
@@ -63,17 +65,50 @@ export class WhatsappConfigService implements OnApplicationBootstrap {
}
}
get mimetypes(): string[] {
if (!this.shouldDownloadMedia) {
return ['mimetype/ignore-all-media'];
}
const types = this.configService.get('WHATSAPP_FILES_MIMETYPES', '');
return types ? types.split(',') : [];
private get mediaGlobalParams() {
return {
download: this.configService.get('WHATSAPP_DOWNLOAD_MEDIA', 'true'),
mimetypes: this.configService.get('WHATSAPP_FILES_MIMETYPES', ''),
};
}
get shouldDownloadMedia(): boolean {
const value = this.configService.get('WHATSAPP_DOWNLOAD_MEDIA', 'true');
return parseBool(value);
private get mediaApiConfig(): MediaDownloadOptions {
const params = {
download: this.configService.get('WAHA_API_DOWNLOAD_MEDIA'),
mimetypes: this.configService.get('WAHA_API_DOWNLOAD_MEDIA_MIMETYPES'),
};
const config = lodash.defaults({}, params, this.mediaGlobalParams);
return this.parseMediaDownloadOptions(config);
}
private get mediaEventsConfig(): MediaDownloadOptions {
const params = {
download: this.configService.get('WAHA_EVENTS_DOWNLOAD_MEDIA'),
mimetypes: this.configService.get('WAHA_EVENTS_DOWNLOAD_MEDIA_MIMETYPES'),
};
const config = lodash.defaults({}, params, this.mediaGlobalParams);
return this.parseMediaDownloadOptions(config);
}
private parseMediaDownloadOptions(config: {
download: string;
mimetypes: string;
}): MediaDownloadOptions {
const types = config.mimetypes;
const mimetypes = types
? types.split(',').map((type: string) => type.trim())
: [];
return {
download: parseBool(config.download),
mimetypes: mimetypes,
};
}
get mediaConfig() {
return {
api: this.mediaApiConfig,
events: this.mediaEventsConfig,
};
}
get startSessions(): string[] {
+21 -2
View File
@@ -139,7 +139,7 @@ import {
AvailableInPlusVersion,
NotImplementedByEngineError,
} from '../exceptions';
import { IMediaManager } from '../media/IMediaManager';
import { IMediaManager, MediaDownloadOptions } from '../media/IMediaManager';
import { QR } from '../QR';
import { DataStore } from './DataStore';
import { fetchBuffer } from '@waha/utils/fetch';
@@ -147,7 +147,6 @@ import {
PRESENCE_AUTO_ONLINE,
PRESENCE_AUTO_ONLINE_DURATION_SECONDS,
} from '@waha/core/env';
import { Activity } from '@waha/core/abc/activity';
// eslint-disable-next-line @typescript-eslint/no-var-requires
const qrcode = require('qrcode-terminal');
@@ -162,6 +161,11 @@ export function ensureSuffix(phone) {
return phone + suffix;
}
interface MediaConfig {
api: MediaDownloadOptions;
events: MediaDownloadOptions;
}
export interface SessionParams {
name: string;
printQR: boolean;
@@ -174,6 +178,8 @@ export interface SessionParams {
engineConfig?: any;
// Ignore settings
ignore: IgnoreJidConfig;
// Media config
media: MediaConfig;
}
/**
@@ -199,6 +205,7 @@ export abstract class WhatsappSession {
protected sessionStore: DataStore;
protected proxyConfig?: ProxyConfig;
public sessionConfig?: SessionConfig;
protected media: MediaConfig;
protected engineConfig?: any;
protected unpairing: boolean = false;
protected jids: JidFilter;
@@ -244,6 +251,7 @@ export abstract class WhatsappSession {
sessionConfig,
engineConfig,
ignore,
media,
}: SessionParams) {
this._status = WAHASessionStatus.STOPPED;
this.status$ = new Subject<SessionStatusUpdate>();
@@ -362,6 +370,17 @@ export abstract class WhatsappSession {
'The session ignores the following chat ids',
);
this.jids = new JidFilter(ignore);
//
// Media options
//
this.media = media;
const mimetypes = this.media.events.mimetypes;
if (mimetypes && mimetypes.length > 0) {
const str = mimetypes.join(',');
const msg = `Only '${str}' mimetypes will be downloaded for the session`;
this.logger.info(msg);
}
}
public getEventObservable(event: WAHAEvents) {
+42 -32
View File
@@ -222,6 +222,7 @@ import axiosRetry from 'axios-retry';
import * as path from 'path';
import MessageServiceClient = messages.MessageServiceClient;
import * as fsp from 'fs/promises';
import { MediaDownloadOptions } from '@waha/core/media/IMediaManager';
axiosRetry(axios, { retries: 3 });
@@ -586,7 +587,7 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
msg.Status = MessageStatus.ServerAck;
return msg;
}),
mergeMap((msg) => this.processIncomingMessage(msg, true)),
mergeMap((msg) => this.processIncomingMessage(msg, this.media.events)),
filter(Boolean),
// Deduplicate messages by ID to prevent duplicate webhooks
// @see https://github.com/devlikeapro/waha/issues/1564
@@ -598,7 +599,7 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
msg.Status = MessageStatus.DeliveryAck;
return msg;
}),
mergeMap((msg) => this.processIncomingMessage(msg, true)),
mergeMap((msg) => this.processIncomingMessage(msg, this.media.events)),
filter(Boolean),
// Deduplicate messages by ID to prevent duplicate webhooks
// @see https://github.com/devlikeapro/waha/issues/1564
@@ -2173,7 +2174,6 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
inviteCode: string,
query: PreviewChannelMessages,
): Promise<ChannelMessage[]> {
const downloadMedia = query.downloadMedia;
const request = new messages.GetNewsletterMessagesByInviteRequest({
session: this.session,
invite: inviteCode,
@@ -2187,12 +2187,17 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
if (!resp.Messages) {
return [];
}
const params = {
download: query.downloadMedia,
mimetypes: query.downloadMediaMimetypes,
};
const options = lodash.defaults({}, params, this.media.api);
for (const msg of resp.Messages) {
promises.push(
this.GowsChannelMessageToChannelMessage(
resp.NewsletterJID,
msg,
downloadMedia,
options,
),
);
}
@@ -2204,7 +2209,7 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
private async GowsChannelMessageToChannelMessage(
jid: string,
channelMessage: any,
downloadMedia: boolean,
options: MediaDownloadOptions,
): Promise<ChannelMessage> {
const msg = {
Info: {
@@ -2217,7 +2222,7 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
},
Message: channelMessage.Message,
};
const message = await this.processIncomingMessage(msg, downloadMedia);
const message = await this.processIncomingMessage(msg, options);
const reactions: any =
sortObjectByValues(channelMessage.ReactionCounts) || {};
return {
@@ -2576,7 +2581,6 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
query: GetChatMessagesQuery,
filter: GetChatMessagesFilter,
) {
const downloadMedia = query.downloadMedia;
const merge = query.merge ?? true;
let jid: messages.OptionalString;
if (chatId === 'all') {
@@ -2620,8 +2624,13 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
const response = await promisify(this.client.GetMessages)(request);
const msgs = parseJsonList(response);
const promises = [];
const params = {
download: query.downloadMedia,
mimetypes: query.downloadMediaMimetypes,
};
const options = lodash.defaults({}, params, this.media.api);
for (const msg of msgs) {
promises.push(this.processIncomingMessage(msg, downloadMedia));
promises.push(this.processIncomingMessage(msg, options));
}
let result = await Promise.all(promises);
result = result.filter(Boolean);
@@ -2648,7 +2657,12 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
});
const response = await promisify(this.client.GetMessageById)(request);
const msg = parseJson(response);
return this.processIncomingMessage(msg, query.downloadMedia);
const params = {
download: query.downloadMedia,
mimetypes: query.downloadMediaMimetypes,
};
const options = lodash.defaults({}, params, this.media.api);
return this.processIncomingMessage(msg, options);
}
/**
@@ -2812,19 +2826,19 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
return true;
}
protected async processIncomingMessage(message, downloadMedia = true) {
protected async processIncomingMessage(
message,
options: MediaDownloadOptions,
) {
// Filter
if (!this.shouldProcessIncomingMessage(message)) {
return null;
}
// Convert
const wamessage = this.toWAMessage(message);
// Media
if (downloadMedia) {
const media = await this.downloadMediaSafe(message);
wamessage.media = media;
}
if (downloadMedia && wamessage.replyTo?.hasMedia) {
const media = await this.downloadMediaSafe(message, options);
wamessage.media = media;
if (wamessage.replyTo?.hasMedia) {
const msg = {
Message: wamessage.replyTo._data,
Info: {
@@ -2832,14 +2846,23 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
ID: wamessage.replyTo.id || '',
},
};
wamessage.replyTo.media = await this.downloadMediaSafe(msg);
wamessage.replyTo.media = await this.downloadMediaSafe(msg, options);
}
return wamessage;
}
protected async downloadMediaSafe(message) {
protected async downloadMediaSafe(message, options: MediaDownloadOptions) {
try {
return await this.downloadMedia(message);
let processor: IMediaEngineProcessor<any> = new GOWSEngineMediaProcessor(
this,
);
processor = new LottieMediaProcessorWrapper(processor, this.logger);
const media = await this.mediaManager.processMedia(
processor,
message,
options,
);
return media;
} catch (e) {
this.logger.error('Failed when tried to download media for a message');
this.logger.error(e, e.stack);
@@ -2847,19 +2870,6 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
}
}
protected async downloadMedia(message) {
let processor: IMediaEngineProcessor<any> = new GOWSEngineMediaProcessor(
this,
);
processor = new LottieMediaProcessorWrapper(processor, this.logger);
const media = await this.mediaManager.processMedia(
processor,
message,
this.name,
);
return media;
}
protected toWAMessage(message): WAMessage {
const fromToParticipant = getFromToParticipant(message);
const id = buildMessageId(message);
+37 -29
View File
@@ -72,6 +72,7 @@ import { toVcardV3 } from '@waha/core/vcard';
import { createAgentProxy } from '@waha/core/helpers.proxy';
import type { Agent } from 'https';
import { IMediaEngineProcessor } from '@waha/core/media/IMediaEngineProcessor';
import { MediaDownloadOptions } from '@waha/core/media/IMediaManager';
import { LottieMediaProcessorWrapper } from '@waha/core/media/LottieMediaProcessorWrapper';
import { QR } from '@waha/core/QR';
import { AckToStatus, StatusToAck } from '@waha/core/utils/acks';
@@ -1549,7 +1550,6 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
query: GetChatMessagesQuery,
filter: GetChatMessagesFilter,
) {
const downloadMedia = query.downloadMedia;
const pagination = query as PaginationParams;
const merge = query.merge ?? true;
const messages = await this.store.getMessagesByJid(
@@ -1560,8 +1560,13 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
);
const promises = [];
const params = {
download: query.downloadMedia,
mimetypes: query.downloadMediaMimetypes,
};
const options = lodash.defaults({}, params, this.media.api);
for (const msg of messages) {
promises.push(this.processIncomingMessage(msg, downloadMedia));
promises.push(this.processIncomingMessage(msg, options));
}
let result = await Promise.all(promises);
result = result.filter(Boolean);
@@ -1589,7 +1594,12 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
merge,
);
if (!message) return null;
return await this.processIncomingMessage(message, query.downloadMedia);
const params = {
download: query.downloadMedia,
mimetypes: query.downloadMediaMimetypes,
};
const options = lodash.defaults({}, params, this.media.api);
return await this.processIncomingMessage(message, options);
}
@Activity()
@@ -2510,7 +2520,6 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
inviteCode: string,
query: PreviewChannelMessages,
): Promise<ChannelMessage[]> {
const downloadMedia = query.downloadMedia;
const updates = await this.sock.newsletterFetchPreviewMessages(
'invite',
inviteCode,
@@ -2518,9 +2527,14 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
null,
);
const promises = [];
const params = {
download: query.downloadMedia,
mimetypes: query.downloadMediaMimetypes,
};
const options = lodash.defaults({}, params, this.media.api);
for (const update of updates) {
promises.push(
this.NewsletterFetchedUpdateToChannelMessage(update, downloadMedia),
this.NewsletterFetchedUpdateToChannelMessage(update, options),
);
}
let result = await Promise.all(promises);
@@ -2530,16 +2544,13 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
private async NewsletterFetchedUpdateToChannelMessage(
update: NewsletterFetchedUpdate,
downloadMedia: boolean,
options: MediaDownloadOptions,
): Promise<ChannelMessage> {
let reactions: any = Object.fromEntries(
update.reactions.map(({ code, count }) => [code, count]),
);
reactions = sortObjectByValues(reactions) || {};
const message = await this.processIncomingMessage(
update.message,
downloadMedia,
);
const message = await this.processIncomingMessage(update.message, options);
return {
message: message,
reactions: reactions,
@@ -2674,13 +2685,13 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
isMine,
);
messagesFromMe$ = messagesFromMe$.pipe(
mergeMap((msg) => this.processIncomingMessage(msg, true)),
mergeMap((msg) => this.processIncomingMessage(msg, this.media.events)),
filter(Boolean),
DistinctMessages(),
share(), // share it so we don't process twice in message.any
);
messagesFromOthers$ = messagesFromOthers$.pipe(
mergeMap((msg) => this.processIncomingMessage(msg, true)),
mergeMap((msg) => this.processIncomingMessage(msg, this.media.events)),
filter(Boolean),
DistinctMessages(),
share(), // share it so we don't process twice in message.any
@@ -3230,7 +3241,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
protected async processIncomingMessage(
message,
downloadMedia: boolean,
options: MediaDownloadOptions,
): Promise<WAMessage | null> {
// Filter
if (!this.shouldProcessIncomingMessage(message)) {
@@ -3242,11 +3253,9 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
return null;
}
// Media
if (downloadMedia && wamessage.hasMedia) {
wamessage.media = await this.downloadMediaSafe(message);
}
wamessage.media = await this.downloadMediaSafe(message, options);
if (downloadMedia && wamessage.replyTo?.hasMedia) {
if (wamessage.replyTo?.hasMedia) {
const mediaContent = extractMediaContent(wamessage.replyTo._data);
const m = {
message: wamessage.replyTo._data,
@@ -3259,7 +3268,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
remoteJid: message.key.remoteJid,
},
};
wamessage.replyTo.media = await this.downloadMediaSafe(m);
wamessage.replyTo.media = await this.downloadMediaSafe(m, options);
}
return wamessage;
}
@@ -3514,9 +3523,17 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
return { id: chatId, presences: presences };
}
protected async downloadMediaSafe(message): Promise<WAMedia | null> {
protected async downloadMediaSafe(
message,
options: MediaDownloadOptions,
): Promise<WAMedia | null> {
try {
return await this.downloadMedia(message);
let processor: IMediaEngineProcessor<any> = new NOWEBEngineMediaProcessor(
this,
this.loggerBuilder,
);
processor = new LottieMediaProcessorWrapper(processor, this.logger);
return await this.mediaManager.processMedia(processor, message, options);
} catch (e) {
this.logger.error('Failed when tried to download media for a message');
this.logger.error(e, e.stack);
@@ -3524,15 +3541,6 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
return null;
}
protected async downloadMedia(message): Promise<WAMedia | null> {
let processor: IMediaEngineProcessor<any> = new NOWEBEngineMediaProcessor(
this,
this.loggerBuilder,
);
processor = new LottieMediaProcessorWrapper(processor, this.logger);
return this.mediaManager.processMedia(processor, message, this.name);
}
protected async getMessageOptions(request: {
id?: string;
chatId: string;
+53 -31
View File
@@ -40,6 +40,7 @@ import { WAMimeType } from '@waha/core/media/WAMimeType';
import { detectMimetype } from '@waha/utils/files';
import { NotImplementedByEngineError } from '@waha/core/exceptions';
import { IMediaEngineProcessor } from '@waha/core/media/IMediaEngineProcessor';
import { MediaDownloadOptions } from '@waha/core/media/IMediaManager';
import { LottieMediaProcessorWrapper } from '@waha/core/media/LottieMediaProcessorWrapper';
import { QR } from '@waha/core/QR';
import { StatusToAck } from '@waha/core/utils/acks';
@@ -1337,7 +1338,6 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
);
}
const downloadMedia = query.downloadMedia;
// Test there's chat with id
await this.whatsapp.getChatById(this.ensureSuffix(chatId));
const pagination: PaginationParams = query;
@@ -1347,8 +1347,13 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
pagination,
);
const promises = [];
const params = {
download: query.downloadMedia,
mimetypes: query.downloadMediaMimetypes,
};
const options = lodash.defaults({}, params, this.media.api);
for (const msg of messages) {
promises.push(this.processIncomingMessage(msg, downloadMedia));
promises.push(this.processIncomingMessage(msg, options));
}
let result = await Promise.all(promises);
result = result.filter(Boolean);
@@ -1430,7 +1435,12 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
return null;
});
}
return await this.processIncomingMessage(message, query.downloadMedia);
const params = {
download: query.downloadMedia,
mimetypes: query.downloadMediaMimetypes,
};
const options = lodash.defaults({}, params, this.media.api);
return await this.processIncomingMessage(message, options);
}
@Activity()
@@ -1964,21 +1974,24 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
query.limit,
);
const promises = [];
const params = {
download: query.downloadMedia,
mimetypes: query.downloadMediaMimetypes,
};
const options = lodash.defaults({}, params, this.media.api);
for (const msg of channelMessages) {
promises.push(
this.WebjsChannelMessageToChannelMessage(msg, query.downloadMedia),
);
promises.push(this.WebjsChannelMessageToChannelMessage(msg, options));
}
return await Promise.all(promises);
}
private async WebjsChannelMessageToChannelMessage(
channelMessage: WebjsChannelMessage,
downloadMedia: boolean,
options: MediaDownloadOptions,
): Promise<ChannelMessage> {
const message = await this.processIncomingMessage(
channelMessage.message,
downloadMedia,
options,
);
const reactions = {};
for (const reaction of channelMessage.reactions.sort((x) => -x.count)) {
@@ -2292,7 +2305,9 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
const messageReceived$ = fromEvent(this.whatsapp, Events.MESSAGE_RECEIVED);
const messagesFromOthers$ = messageReceived$.pipe(
filter((msg: Message) => this.jids.include(msg?.id?.remote)),
mergeMap((msg: any) => this.processIncomingMessage(msg, true)),
mergeMap((msg: any) =>
this.processIncomingMessage(msg, this.media.events),
),
share(),
);
this.events2.get(WAHAEvents.MESSAGE).switch(messagesFromOthers$);
@@ -2300,7 +2315,9 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
const messageCreate$ = fromEvent(this.whatsapp, Events.MESSAGE_CREATE);
const messagesFromAll$ = messageCreate$.pipe(
filter((msg: Message) => this.jids.include(msg?.id?.remote)),
mergeMap((msg: any) => this.processIncomingMessage(msg, true)),
mergeMap((msg: any) =>
this.processIncomingMessage(msg, this.media.events),
),
share(),
);
this.events2.get(WAHAEvents.MESSAGE_ANY).switch(messagesFromAll$);
@@ -2311,7 +2328,12 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
);
const messagesWaiting$ = messageCiphertext$.pipe(
filter((msg: Message) => this.jids.include(msg?.id?.remote)),
mergeMap((msg: any) => this.processIncomingMessage(msg, false)),
mergeMap((msg: any) =>
this.processIncomingMessage(msg, {
...this.media.events,
download: false,
}),
),
share(),
);
this.events2.get(WAHAEvents.MESSAGE_WAITING).switch(messagesWaiting$);
@@ -2543,19 +2565,20 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
protected async processIncomingMessage(
message: Message,
downloadMedia = true,
options: MediaDownloadOptions,
) {
// Convert
const wamessage = this.toWAMessage(message);
// Media
if (downloadMedia) {
const media = await this.downloadMediaSafe(message);
wamessage.media = media;
}
if (downloadMedia && wamessage.replyTo?.hasMedia) {
const media = await this.downloadMediaSafe(message, options);
wamessage.media = media;
if (wamessage.replyTo?.hasMedia) {
const quotedMessage = await message.getQuotedMessage().catch(() => null);
if (quotedMessage) {
wamessage.replyTo.media = await this.downloadMediaSafe(quotedMessage);
wamessage.replyTo.media = await this.downloadMediaSafe(
quotedMessage,
options,
);
}
}
return wamessage;
@@ -2773,9 +2796,19 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
return contact;
}
protected async downloadMediaSafe(message): Promise<WAMedia | null> {
protected async downloadMediaSafe(
message: Message,
options: MediaDownloadOptions,
): Promise<WAMedia | null> {
try {
return await this.downloadMedia(message);
let processor = new WEBJSEngineMediaProcessor();
processor = new LottieMediaProcessorWrapper(processor, this.logger);
const media = await this.mediaManager.processMedia(
processor,
message,
options,
);
return media;
} catch (e) {
this.logger.error('Failed when tried to download media for a message');
this.logger.error(e, e.stack);
@@ -2783,17 +2816,6 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
return null;
}
protected async downloadMedia(message: Message) {
let processor = new WEBJSEngineMediaProcessor();
processor = new LottieMediaProcessorWrapper(processor, this.logger);
const media = await this.mediaManager.processMedia(
processor,
message,
this.name,
);
return media;
}
protected getMessageOptions(request: any): any {
let mentions = request.mentions;
mentions = mentions ? mentions.map(this.ensureSuffix) : undefined;
+32 -23
View File
@@ -132,6 +132,7 @@ import {
} from '@waha/core/engines/wpp/WppTypes';
import { NotImplementedByEngineError } from '@waha/core/exceptions';
import { IMediaEngineProcessor } from '@waha/core/media/IMediaEngineProcessor';
import { MediaDownloadOptions } from '@waha/core/media/IMediaManager';
import { LottieMediaProcessorWrapper } from '@waha/core/media/LottieMediaProcessorWrapper';
import { IWPPAuthManager } from '@waha/core/engines/wpp/IWPPAuthManager';
import { QR } from '@waha/core/QR';
@@ -1103,7 +1104,6 @@ export class WhatsappSessionWPPCore extends WhatsappSession {
const offset = query?.offset || 0;
const limit = query?.limit || 10;
const fetchCount = offset + limit;
const downloadMedia = query.downloadMedia;
const rawMessages = await this.wpp!.getMessages(id, {
count: fetchCount,
direction: 'before',
@@ -1127,10 +1127,11 @@ export class WhatsappSessionWPPCore extends WhatsappSession {
sortOrder: query?.sortOrder || SortOrder.DESC,
}).apply(messages);
if (!downloadMedia) {
return messages;
}
const params = {
download: query.downloadMedia,
mimetypes: query.downloadMediaMimetypes,
};
const options = lodash.defaults({}, params, this.media.api);
const promises = [];
for (const message of messages) {
const rawMessage = messagesById.get(message.id);
@@ -1138,7 +1139,7 @@ export class WhatsappSessionWPPCore extends WhatsappSession {
promises.push(Promise.resolve(message));
continue;
}
promises.push(this.processIncomingMessage(rawMessage, true));
promises.push(this.processIncomingMessage(rawMessage, options));
}
let result = await Promise.all(promises);
result = result.filter(Boolean);
@@ -1155,7 +1156,12 @@ export class WhatsappSessionWPPCore extends WhatsappSession {
if (!message) {
return null;
}
return this.processIncomingMessage(message, query.downloadMedia);
const params = {
download: query.downloadMedia,
mimetypes: query.downloadMediaMimetypes,
};
const options = lodash.defaults({}, params, this.media.api);
return this.processIncomingMessage(message, options);
}
@Activity()
@@ -2429,27 +2435,25 @@ export class WhatsappSessionWPPCore extends WhatsappSession {
};
}
protected async processIncomingMessage(message: any, downloadMedia = true) {
protected async processIncomingMessage(
message: any,
options: MediaDownloadOptions,
) {
const wamessage = this.toWAMessage(message);
if (downloadMedia) {
const media = await this.downloadMediaSafe(message);
wamessage.media = media;
}
if (downloadMedia && wamessage.replyTo?.hasMedia) {
const media = await this.downloadMediaSafe(message, options);
wamessage.media = media;
if (wamessage.replyTo?.hasMedia) {
const quotedMessage = message?.quotedMsg || message?._data?.quotedMsg;
if (quotedMessage) {
wamessage.replyTo.media = await this.downloadMediaSafe(quotedMessage);
wamessage.replyTo.media = await this.downloadMediaSafe(
quotedMessage,
options,
);
}
}
return wamessage;
}
protected async downloadMedia(message: any): Promise<WAMedia | null> {
let processor = new WPPEngineMediaProcessor(this.wpp);
processor = new LottieMediaProcessorWrapper(processor, this.logger);
return this.mediaManager.processMedia(processor, message, this.name);
}
protected checkStatusRequest(request: { contacts?: any[] }) {
if (request.contacts && request.contacts?.length > 0) {
const msg =
@@ -2458,9 +2462,14 @@ export class WhatsappSessionWPPCore extends WhatsappSession {
}
}
protected async downloadMediaSafe(message): Promise<WAMedia | null> {
protected async downloadMediaSafe(
message: any,
options: MediaDownloadOptions,
): Promise<WAMedia | null> {
try {
return await this.downloadMedia(message);
let processor = new WPPEngineMediaProcessor(this.wpp);
processor = new LottieMediaProcessorWrapper(processor, this.logger);
return await this.mediaManager.processMedia(processor, message, options);
} catch (error) {
this.logger.error('Failed when tried to download media for a message');
this.logger.error(error, error.stack);
@@ -2523,7 +2532,7 @@ export class WhatsappSessionWPPCore extends WhatsappSession {
if (msg?.fromMe) {
await sleep(3_000);
}
return this.processIncomingMessage(msg, true);
return this.processIncomingMessage(msg, this.media.events);
}
private async refreshMeInfo() {
+2 -1
View File
@@ -344,8 +344,8 @@ export class SessionManagerCore
);
await storage.init();
const mediaManager = new MediaManager(
name,
storage,
this.config.mimetypes,
loggerBuilder.child({ name: 'MediaManager' }),
);
const webhook = new WebhookConductor(loggerBuilder);
@@ -359,6 +359,7 @@ export class SessionManagerCore
proxyConfig: proxyConfig,
sessionConfig: config,
ignore: this.ignoreChatsConfig(config),
media: this.config.mediaConfig,
};
if (this.EngineClass === WhatsappSessionWebJSCore) {
sessionConfig.engineConfig = this.webjsEngineConfigService.getConfig();
+7 -4
View File
@@ -1,17 +1,20 @@
import { IMediaEngineProcessor } from '@waha/core/media/IMediaEngineProcessor';
import { WAMedia } from '@waha/structures/media.dto';
export interface MediaDownloadOptions {
download: boolean;
mimetypes: string[];
}
/**
* General interface for MediaManager - one that handles the logic
* and manipulates MediaStorage and MediaEngineProcessor
*/
interface IMediaManager {
export interface IMediaManager {
processMedia<Message>(
processor: IMediaEngineProcessor<Message>,
message: Message,
session: string,
options: MediaDownloadOptions,
): Promise<WAMedia | null>;
close(): void;
}
export { IMediaManager };
+21 -23
View File
@@ -1,5 +1,8 @@
import { IMediaEngineProcessor } from '@waha/core/media/IMediaEngineProcessor';
import { IMediaManager } from '@waha/core/media/IMediaManager';
import {
IMediaManager,
MediaDownloadOptions,
} from '@waha/core/media/IMediaManager';
import {
IMediaStorage,
MediaData,
@@ -22,52 +25,38 @@ export class MediaManager implements IMediaManager {
};
constructor(
private sessionName: string,
private storage: IMediaStorage,
private mimetypes: string[],
protected log: Logger,
) {
// Log mimetypes
if (this.mimetypes && this.mimetypes.length > 0) {
const mimetypes = this.mimetypes.join(',');
const msg = `Only '${mimetypes}' mimetypes will be downloaded for the session`;
this.log.info(msg);
}
}
) {}
/**
* Check that we need to download files with the mimetype
*/
private shouldProcessMimetype(mimetype: string) {
private shouldProcessMimetype(mimetypes: string[], mimetype: string) {
// No specific mimetypes provided - always download
if (!this.mimetypes || this.mimetypes.length === 0) {
if (!mimetypes || mimetypes.length === 0) {
return true;
}
// Found "right" mimetype in the list of allowed mimetypes - download it
return this.mimetypes.some((type) => mimetype.startsWith(type));
return mimetypes.some((type) => mimetype.startsWith(type));
}
private async processMediaInternal<Message>(
processor: IMediaEngineProcessor<Message>,
message: Message,
session: string,
): Promise<WAMedia | null> {
const messageId = processor.getMessageId(message);
const chatId = processor.getChatId(message);
const mimetype = processor.getMimetype(message);
const filename = processor.getFilename(message);
if (!this.shouldProcessMimetype(mimetype)) {
this.log.info(
`The message '${messageId}' has '${mimetype}' mimetype media, skip it.`,
);
return null;
}
let extension = mime.extension(mimetype);
if (mimetype == 'application/was' && !extension) {
extension = 'zip';
}
const mediaData: MediaData = {
session: session,
session: this.sessionName,
message: {
id: messageId,
chatId: chatId,
@@ -105,7 +94,7 @@ export class MediaManager implements IMediaManager {
async processMedia<Message>(
processor: IMediaEngineProcessor<Message>,
message: Message,
session: string,
options: MediaDownloadOptions,
): Promise<WAMedia | null> {
let messageId: string;
try {
@@ -129,7 +118,16 @@ export class MediaManager implements IMediaManager {
try {
media.filename = processor.getFilename(message);
media.mimetype = processor.getMimetype(message);
const data = await this.processMediaInternal(processor, message, session);
if (!options.download) {
return media;
}
if (!this.shouldProcessMimetype(options.mimetypes, media.mimetype)) {
this.log.info(
`The message '${messageId}' has '${media.mimetype}' mimetype media, skip it.`,
);
return media;
}
const data = await this.processMediaInternal(processor, message);
media = { ...media, ...data };
} catch (err) {
this.log.error(err, `Error processing media for message '${messageId}'`);
@@ -0,0 +1,45 @@
import { CommaSeparatedStrings } from './CommaSeparatedStrings';
describe('CommaSeparatedStrings', () => {
it('should convert a single value to an array', () => {
// ?a=image/jpeg
expect(CommaSeparatedStrings({ value: 'image/jpeg' })).toEqual([
'image/jpeg',
]);
});
it('should keep the repeated form as is', () => {
// ?a=image/jpeg&a=image/png
expect(
CommaSeparatedStrings({ value: ['image/jpeg', 'image/png'] }),
).toEqual(['image/jpeg', 'image/png']);
});
it('should split comma separated values', () => {
// ?a=image/jpeg,image/png
expect(CommaSeparatedStrings({ value: 'image/jpeg,image/png' })).toEqual([
'image/jpeg',
'image/png',
]);
});
it('should split comma separated values in the repeated form', () => {
// ?a=image/jpeg,image/png&a=video/mp4
expect(
CommaSeparatedStrings({ value: ['image/jpeg,image/png', 'video/mp4'] }),
).toEqual(['image/jpeg', 'image/png', 'video/mp4']);
});
it('should trim values and remove empty ones', () => {
// ?a=image/jpeg, image/png,
expect(CommaSeparatedStrings({ value: 'image/jpeg, image/png,' })).toEqual([
'image/jpeg',
'image/png',
]);
});
it('should keep null and undefined as is', () => {
expect(CommaSeparatedStrings({ value: null })).toBeNull();
expect(CommaSeparatedStrings({ value: undefined })).toBeUndefined();
});
});
@@ -0,0 +1,20 @@
/**
* Convert a query param to a list of strings.
* Accepts:
* - single value: ?a=x => ['x']
* - repeated form: ?a=x&a=y => ['x', 'y']
* - comma separated form: ?a=x,y => ['x', 'y']
* Values are trimmed, empty ones are removed.
* @param value
* @constructor
*/
export function CommaSeparatedStrings({ value }: { value: any }) {
if (value == null) {
return value;
}
const values = Array.isArray(value) ? value : [value];
return values
.flatMap((item: string) => String(item).split(','))
.map((item: string) => item.trim())
.filter(Boolean);
}
+13
View File
@@ -1,5 +1,6 @@
import { ApiParam, ApiProperty, getSchemaPath } from '@nestjs/swagger';
import { BooleanString } from '@waha/nestjs/validation/BooleanString';
import { CommaSeparatedStrings } from '@waha/nestjs/validation/CommaSeparatedStrings';
import { BinaryFile, RemoteFile } from '@waha/structures/files.dto';
import { WAMessage } from '@waha/structures/responses.dto';
import { Transform, Type } from 'class-transformer';
@@ -193,6 +194,18 @@ export class PreviewChannelMessages {
@IsBoolean()
downloadMedia: boolean = false;
@ApiProperty({
type: String,
isArray: true,
required: false,
description: 'Download only media with these mimetypes (prefix match)',
})
@Transform(CommaSeparatedStrings)
@IsArray()
@IsString({ each: true })
@IsOptional()
downloadMediaMimetypes?: string[];
@IsNumber()
@Type(() => Number)
limit: number = 10;
+28 -5
View File
@@ -1,6 +1,7 @@
import { BadRequestException } from '@nestjs/common';
import { ApiProperty } from '@nestjs/swagger';
import { BooleanString } from '@waha/nestjs/validation/BooleanString';
import { CommaSeparatedStrings } from '@waha/nestjs/validation/CommaSeparatedStrings';
import { WAMessageAck, WAMessageAckName } from '@waha/structures/enums.dto';
import {
LimitOffsetParams,
@@ -113,14 +114,25 @@ export class GetChatMessagesQuery extends PaginationParams {
sortBy?: string = MessageSortField.TIMESTAMP;
@ApiProperty({
example: false,
required: false,
description: 'Download media for messages',
})
@Transform(BooleanString)
@IsBoolean()
@IsOptional()
downloadMedia: boolean = true;
downloadMedia?: boolean;
@ApiProperty({
type: String,
isArray: true,
required: false,
description: 'Download only media with these mimetypes (prefix match)',
})
@Transform(CommaSeparatedStrings)
@IsArray()
@IsString({ each: true })
@IsOptional()
downloadMediaMimetypes?: string[];
@ApiProperty({
example: true,
@@ -165,14 +177,25 @@ export class ReadChatMessagesResponse {
export class GetChatMessageQuery {
@ApiProperty({
example: true,
required: false,
description: 'Download media for messages',
})
@Transform(BooleanString)
@IsBoolean()
@IsOptional()
downloadMedia: boolean = true;
downloadMedia?: boolean;
@ApiProperty({
type: String,
isArray: true,
required: false,
description: 'Download only media with these mimetypes (prefix match)',
})
@Transform(CommaSeparatedStrings)
@IsArray()
@IsString({ each: true })
@IsOptional()
downloadMediaMimetypes?: string[];
@ApiProperty({
example: true,
@@ -256,7 +279,7 @@ export class OverviewFilter {
@IsOptional()
@IsArray()
@IsString({ each: true })
@Transform(({ value }) => (Array.isArray(value) ? value : [value]))
@Transform(CommaSeparatedStrings)
@ApiProperty({
description: 'Filter by chat ids',
required: false,