1 parent
f35fd8d7c6
commit
72a301c579
7 files changed
+79
-37
No files matched your search
@@ -6,6 +6,7 @@ import {
|
||||
NoRetriesJobOptions,
|
||||
} from '@waha/apps/app_sdk/constants';
|
||||
import { MessageCleanupConsumer } from '@waha/apps/chatwoot/consumers/scheduled/message.cleanup';
|
||||
import { CheckTierConsumer } from '@waha/apps/chatwoot/consumers/scheduled/check.tier';
|
||||
import { ChatWootAppService } from '@waha/apps/chatwoot/services/ChatWootAppService';
|
||||
import * as lodash from 'lodash';
|
||||
|
||||
@@ -51,6 +52,10 @@ const IMPORTS = lodash.flatten([
|
||||
name: QueueName.SCHEDULED_CHECK_VERSION,
|
||||
defaultJobOptions: merge(NoRetriesJobOptions, JobRemoveOptions),
|
||||
}),
|
||||
RegisterAppQueue({
|
||||
name: QueueName.SCHEDULED_CHECK_TIER,
|
||||
defaultJobOptions: merge(NoRetriesJobOptions, JobRemoveOptions),
|
||||
}),
|
||||
RegisterAppQueue({
|
||||
name: QueueName.TASK_CONTACTS_PULL,
|
||||
defaultJobOptions: merge(ExponentialRetriesJobOptions, JobRemoveOptions),
|
||||
@@ -145,6 +150,7 @@ const PROVIDERS = [
|
||||
// Scheduled
|
||||
MessageCleanupConsumer,
|
||||
CheckVersionConsumer,
|
||||
CheckTierConsumer,
|
||||
// Services
|
||||
ChatWootWAHAQueueService,
|
||||
ChatWootQueueService,
|
||||
|
||||
@@ -4,6 +4,7 @@ export enum QueueName {
|
||||
//
|
||||
SCHEDULED_MESSAGE_CLEANUP = 'chatwoot.scheduled | message.cleanup',
|
||||
SCHEDULED_CHECK_VERSION = 'chatwoot.scheduled | check.version',
|
||||
SCHEDULED_CHECK_TIER = 'chatwoot.scheduled | check.tier',
|
||||
|
||||
//
|
||||
// Task
|
||||
|
||||
@@ -0,0 +1,42 @@
|
||||
import { Processor } from '@nestjs/bullmq';
|
||||
import { Injectable } from '@nestjs/common';
|
||||
import { JOB_CONCURRENCY } from '@waha/apps/app_sdk/constants';
|
||||
import { QueueName } from '@waha/apps/chatwoot/consumers/QueueName';
|
||||
import { DIContainer } from '@waha/apps/chatwoot/di/DIContainer';
|
||||
import { SessionManager } from '@waha/core/abc/manager.abc';
|
||||
import { SUPPORT_US_URL } from '@waha/core/constants';
|
||||
import { RMutexService } from '@waha/modules/rmutex/rmutex.service';
|
||||
import { VERSION, WAHAVersion } from '@waha/version';
|
||||
import { Job } from 'bullmq';
|
||||
import { PinoLogger } from 'nestjs-pino';
|
||||
|
||||
import { TKey } from '@waha/apps/chatwoot/i18n/templates';
|
||||
import { ChatWootScheduledConsumer } from './base';
|
||||
|
||||
@Processor(QueueName.SCHEDULED_CHECK_TIER, { concurrency: JOB_CONCURRENCY })
|
||||
export class CheckTierConsumer extends ChatWootScheduledConsumer {
|
||||
constructor(manager: SessionManager, log: PinoLogger, rmutex: RMutexService) {
|
||||
super(manager, log, rmutex, CheckTierConsumer.name);
|
||||
}
|
||||
|
||||
protected ErrorHeaderKey(): TKey {
|
||||
return TKey.JOB_SCHEDULED_ERROR_HEADER;
|
||||
}
|
||||
|
||||
protected async Process(container: DIContainer, job: Job): Promise<any> {
|
||||
const logger = container.Logger();
|
||||
const locale = container.Locale();
|
||||
const conversation = await container
|
||||
.ContactConversationService()
|
||||
.InboxNotifications();
|
||||
if (VERSION.tier !== WAHAVersion.CORE) {
|
||||
logger.info('WAHA is not using the CORE version');
|
||||
return;
|
||||
}
|
||||
logger.info('WAHA is using the CORE version');
|
||||
const supportMessage = locale.key(TKey.WAHA_CORE_VERSION_USED).render({
|
||||
supportUrl: SUPPORT_US_URL,
|
||||
});
|
||||
await conversation.incoming(supportMessage);
|
||||
}
|
||||
}
|
||||
@@ -1,24 +1,19 @@
|
||||
import { Processor } from '@nestjs/bullmq';
|
||||
import { Injectable } from '@nestjs/common';
|
||||
import { JOB_CONCURRENCY } from '@waha/apps/app_sdk/constants';
|
||||
import { ILogger } from '@waha/apps/app_sdk/ILogger';
|
||||
import { QueueName } from '@waha/apps/chatwoot/consumers/QueueName';
|
||||
import { DIContainer } from '@waha/apps/chatwoot/di/DIContainer';
|
||||
import { SessionManager } from '@waha/core/abc/manager.abc';
|
||||
import { CHANGELOG_URL, SUPPORT_US_URL } from '@waha/core/constants';
|
||||
import { CHANGELOG_URL } from '@waha/core/constants';
|
||||
import { RMutexService } from '@waha/modules/rmutex/rmutex.service';
|
||||
import { VERSION, WAHAVersion } from '@waha/version';
|
||||
import { VERSION } from '@waha/version';
|
||||
import axios from 'axios';
|
||||
import { Job } from 'bullmq';
|
||||
import { PinoLogger } from 'nestjs-pino';
|
||||
|
||||
import { ChatWootScheduledConsumer } from './base';
|
||||
import { TKey } from '@waha/apps/chatwoot/i18n/templates';
|
||||
import { ChatWootScheduledConsumer } from './base';
|
||||
|
||||
/**
|
||||
* Scheduled consumer for version checking
|
||||
* Periodically checks for new versions and notifies if updates are available
|
||||
*/
|
||||
@Processor(QueueName.SCHEDULED_CHECK_VERSION, { concurrency: JOB_CONCURRENCY })
|
||||
export class CheckVersionConsumer extends ChatWootScheduledConsumer {
|
||||
constructor(manager: SessionManager, log: PinoLogger, rmutex: RMutexService) {
|
||||
@@ -29,34 +24,7 @@ export class CheckVersionConsumer extends ChatWootScheduledConsumer {
|
||||
return TKey.JOB_SCHEDULED_ERROR_HEADER;
|
||||
}
|
||||
|
||||
/**
|
||||
* Process the version check job
|
||||
* Checks for new versions and notifies if updates are available
|
||||
*/
|
||||
protected async Process(container: DIContainer, job: Job): Promise<any> {
|
||||
await this.CheckWAHACoreVersion(container);
|
||||
await this.CheckNewVersionAvailable(container);
|
||||
}
|
||||
|
||||
private async CheckWAHACoreVersion(container: DIContainer) {
|
||||
const logger = container.Logger();
|
||||
const locale = container.Locale();
|
||||
const conversation = await container
|
||||
.ContactConversationService()
|
||||
.InboxNotifications();
|
||||
// If using the CORE version, send a message encouraging support
|
||||
if (VERSION.tier !== WAHAVersion.CORE) {
|
||||
logger.info('WAHA is not using the CORE version');
|
||||
return;
|
||||
}
|
||||
logger.info('WAHA is using the CORE version');
|
||||
const supportMessage = locale.key(TKey.WAHA_CORE_VERSION_USED).render({
|
||||
supportUrl: SUPPORT_US_URL,
|
||||
});
|
||||
await conversation.incoming(supportMessage);
|
||||
}
|
||||
|
||||
private async CheckNewVersionAvailable(container: DIContainer) {
|
||||
const logger = container.Logger();
|
||||
logger.info('Processing version check job');
|
||||
const currentVersion = VERSION.version;
|
||||
@@ -72,7 +40,6 @@ export class CheckVersionConsumer extends ChatWootScheduledConsumer {
|
||||
`New version available: ${latestVersion} (current: ${currentVersion})`,
|
||||
);
|
||||
|
||||
// Send a message to Chatwoot conversation about a new version
|
||||
const locale = container.Locale();
|
||||
const message = locale.key(TKey.WAHA_NEW_VERSION_AVAILABLE).render({
|
||||
currentVersion: currentVersion,
|
||||
|
||||
@@ -54,6 +54,21 @@ export class ChatWootScheduleService {
|
||||
},
|
||||
},
|
||||
);
|
||||
// Check the tier
|
||||
const checkTierQueue = this.queueRegistry.queue(
|
||||
QueueName.SCHEDULED_CHECK_TIER,
|
||||
);
|
||||
await checkTierQueue.upsertJobScheduler(
|
||||
this.JobId(QueueName.SCHEDULED_CHECK_TIER, appId),
|
||||
// Every Monday (1) at 14:00
|
||||
{ pattern: '0 0 14 * * 1' },
|
||||
{
|
||||
data: {
|
||||
app: appId,
|
||||
session: sessionName,
|
||||
},
|
||||
},
|
||||
);
|
||||
}
|
||||
|
||||
async unschedule(appId: string, sessionName: string): Promise<void> {
|
||||
@@ -71,6 +86,13 @@ export class ChatWootScheduleService {
|
||||
await checkVersionQueue.removeJobScheduler(
|
||||
this.JobId(QueueName.SCHEDULED_CHECK_VERSION, appId),
|
||||
);
|
||||
// Check the tier
|
||||
const checkTierQueue = this.queueRegistry.queue(
|
||||
QueueName.SCHEDULED_CHECK_TIER,
|
||||
);
|
||||
await checkTierQueue.removeJobScheduler(
|
||||
this.JobId(QueueName.SCHEDULED_CHECK_TIER, appId),
|
||||
);
|
||||
|
||||
// contacts
|
||||
const contactsPullQueue = this.queueRegistry.queue(
|
||||
|
||||
@@ -18,7 +18,8 @@ export class QueueManager {
|
||||
constructor(private readonly registry: QueueRegistry) {
|
||||
this.queues = {
|
||||
[QueueName.SCHEDULED_MESSAGE_CLEANUP]: Locked,
|
||||
[QueueName.SCHEDULED_CHECK_VERSION]: Locked,
|
||||
[QueueName.SCHEDULED_CHECK_VERSION]: Managable,
|
||||
[QueueName.SCHEDULED_CHECK_TIER]: Locked,
|
||||
[QueueName.TASK_CONTACTS_PULL]: Locked,
|
||||
[QueueName.TASK_MESSAGES_PULL]: Locked,
|
||||
[QueueName.WAHA_SESSION_STATUS]: Locked,
|
||||
|
||||
@@ -16,6 +16,8 @@ export class QueueRegistry {
|
||||
private readonly scheduledMessageCleanupQueue: Queue,
|
||||
@InjectQueue(QueueName.SCHEDULED_CHECK_VERSION)
|
||||
private readonly scheduledCheckVersionQueue: Queue,
|
||||
@InjectQueue(QueueName.SCHEDULED_CHECK_TIER)
|
||||
private readonly scheduledCheckTierQueue: Queue,
|
||||
@InjectQueue(QueueName.TASK_CONTACTS_PULL)
|
||||
private readonly taskContactsPullQueue: Queue,
|
||||
@InjectQueue(QueueName.TASK_MESSAGES_PULL)
|
||||
@@ -55,6 +57,7 @@ export class QueueRegistry {
|
||||
this.queues = {
|
||||
[QueueName.SCHEDULED_MESSAGE_CLEANUP]: this.scheduledMessageCleanupQueue,
|
||||
[QueueName.SCHEDULED_CHECK_VERSION]: this.scheduledCheckVersionQueue,
|
||||
[QueueName.SCHEDULED_CHECK_TIER]: this.scheduledCheckTierQueue,
|
||||
[QueueName.TASK_CONTACTS_PULL]: this.taskContactsPullQueue,
|
||||
[QueueName.TASK_MESSAGES_PULL]: this.taskMessagesPullQueue,
|
||||
[QueueName.WAHA_SESSION_STATUS]: this.wahaSessionStatusQueue,
|
||||
|
||||
Reference in new issue
Block a user