Improvement on messaging (#16351)
In this PR: - change messaging / calendar stale duration check to 30minutes (cron is running every 1h, duration check was 1h, so evaluation was flaky) - when temporary error (throttling), preserve syncStageStartedAt as this is necessary to assess exponential throttling
This commit is contained in:
+4
@@ -35,6 +35,7 @@ export class MessageChannelSyncStatusService {
|
||||
public async scheduleMessageListFetch(
|
||||
messageChannelIds: string[],
|
||||
workspaceId: string,
|
||||
preserveSyncStageStartedAt: boolean = false,
|
||||
) {
|
||||
if (!messageChannelIds.length) {
|
||||
return;
|
||||
@@ -48,12 +49,14 @@ export class MessageChannelSyncStatusService {
|
||||
|
||||
await messageChannelRepository.update(messageChannelIds, {
|
||||
syncStage: MessageChannelSyncStage.MESSAGE_LIST_FETCH_PENDING,
|
||||
...(!preserveSyncStageStartedAt ? { syncStageStartedAt: null } : {}),
|
||||
});
|
||||
}
|
||||
|
||||
public async scheduleMessagesImport(
|
||||
messageChannelIds: string[],
|
||||
workspaceId: string,
|
||||
preserveSyncStageStartedAt: boolean = false,
|
||||
) {
|
||||
if (!messageChannelIds.length) {
|
||||
return;
|
||||
@@ -67,6 +70,7 @@ export class MessageChannelSyncStatusService {
|
||||
|
||||
await messageChannelRepository.update(messageChannelIds, {
|
||||
syncStage: MessageChannelSyncStage.MESSAGES_IMPORT_PENDING,
|
||||
...(!preserveSyncStageStartedAt ? { syncStageStartedAt: null } : {}),
|
||||
});
|
||||
}
|
||||
|
||||
|
||||
+1
-1
@@ -1 +1 @@
|
||||
export const MESSAGING_IMPORT_ONGOING_SYNC_TIMEOUT = 1000 * 60 * 60; // 1 hour
|
||||
export const MESSAGING_IMPORT_ONGOING_SYNC_TIMEOUT = 1000 * 60 * 30; // 30 minutes
|
||||
|
||||
+1
-1
@@ -51,7 +51,7 @@ export class MessagingMessageListFetchCronJob {
|
||||
const now = new Date().toISOString();
|
||||
|
||||
const [messageChannels] = await this.coreDataSource.query(
|
||||
`UPDATE ${schemaName}."messageChannel" SET "syncStage" = '${MessageChannelSyncStage.MESSAGE_LIST_FETCH_SCHEDULED}', "syncStageStartedAt" = '${now}'
|
||||
`UPDATE ${schemaName}."messageChannel" SET "syncStage" = '${MessageChannelSyncStage.MESSAGE_LIST_FETCH_SCHEDULED}', "syncStageStartedAt" = COALESCE("syncStageStartedAt", '${now}')
|
||||
WHERE "isSyncEnabled" = true AND "syncStage" = '${MessageChannelSyncStage.MESSAGE_LIST_FETCH_PENDING}' RETURNING *`,
|
||||
);
|
||||
|
||||
|
||||
+1
-1
@@ -56,7 +56,7 @@ export class MessagingMessagesImportCronJob {
|
||||
const now = new Date().toISOString();
|
||||
|
||||
const [messageChannels] = await this.coreDataSource.query(
|
||||
`UPDATE ${schemaName}."messageChannel" SET "syncStage" = '${MessageChannelSyncStage.MESSAGES_IMPORT_SCHEDULED}', "syncStageStartedAt" = '${now}'
|
||||
`UPDATE ${schemaName}."messageChannel" SET "syncStage" = '${MessageChannelSyncStage.MESSAGES_IMPORT_SCHEDULED}', "syncStageStartedAt" = COALESCE("syncStageStartedAt", '${now}')
|
||||
WHERE "isSyncEnabled" = true AND "syncStage" = '${MessageChannelSyncStage.MESSAGES_IMPORT_PENDING}' RETURNING *`,
|
||||
);
|
||||
|
||||
|
||||
+31
-24
@@ -5,6 +5,7 @@ import { Processor } from 'src/engine/core-modules/message-queue/decorators/proc
|
||||
import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants';
|
||||
import { TwentyORMManager } from 'src/engine/twenty-orm/twenty-orm.manager';
|
||||
import { isThrottled } from 'src/modules/connected-account/utils/is-throttled';
|
||||
import { MessageChannelSyncStatusService } from 'src/modules/messaging/common/services/message-channel-sync-status.service';
|
||||
import {
|
||||
MessageChannelSyncStage,
|
||||
type MessageChannelWorkspaceEntity,
|
||||
@@ -31,6 +32,7 @@ export class MessagingMessageListFetchJob {
|
||||
private readonly messagingMonitoringService: MessagingMonitoringService,
|
||||
private readonly twentyORMManager: TwentyORMManager,
|
||||
private readonly messageImportErrorHandlerService: MessageImportExceptionHandlerService,
|
||||
private readonly messageChannelSyncStatusService: MessageChannelSyncStatusService,
|
||||
) {}
|
||||
|
||||
@Process(MessagingMessageListFetchJob.name)
|
||||
@@ -65,6 +67,13 @@ export class MessagingMessageListFetchJob {
|
||||
return;
|
||||
}
|
||||
|
||||
if (
|
||||
messageChannel.syncStage !==
|
||||
MessageChannelSyncStage.MESSAGE_LIST_FETCH_SCHEDULED
|
||||
) {
|
||||
return;
|
||||
}
|
||||
|
||||
try {
|
||||
if (
|
||||
isThrottled(
|
||||
@@ -72,35 +81,33 @@ export class MessagingMessageListFetchJob {
|
||||
messageChannel.throttleFailureCount,
|
||||
)
|
||||
) {
|
||||
await this.messageChannelSyncStatusService.scheduleMessageListFetch(
|
||||
[messageChannel.id],
|
||||
workspaceId,
|
||||
true,
|
||||
);
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
switch (messageChannel.syncStage) {
|
||||
case MessageChannelSyncStage.MESSAGE_LIST_FETCH_SCHEDULED:
|
||||
await this.messagingMonitoringService.track({
|
||||
eventName: 'message_list_fetch.started',
|
||||
workspaceId,
|
||||
connectedAccountId: messageChannel.connectedAccount.id,
|
||||
messageChannelId: messageChannel.id,
|
||||
});
|
||||
await this.messagingMonitoringService.track({
|
||||
eventName: 'message_list_fetch.started',
|
||||
workspaceId,
|
||||
connectedAccountId: messageChannel.connectedAccount.id,
|
||||
messageChannelId: messageChannel.id,
|
||||
});
|
||||
|
||||
await this.messagingMessageListFetchService.processMessageListFetch(
|
||||
messageChannel,
|
||||
workspaceId,
|
||||
);
|
||||
await this.messagingMessageListFetchService.processMessageListFetch(
|
||||
messageChannel,
|
||||
workspaceId,
|
||||
);
|
||||
|
||||
await this.messagingMonitoringService.track({
|
||||
eventName: 'message_list_fetch.completed',
|
||||
workspaceId,
|
||||
connectedAccountId: messageChannel.connectedAccount.id,
|
||||
messageChannelId: messageChannel.id,
|
||||
});
|
||||
|
||||
break;
|
||||
|
||||
default:
|
||||
break;
|
||||
}
|
||||
await this.messagingMonitoringService.track({
|
||||
eventName: 'message_list_fetch.completed',
|
||||
workspaceId,
|
||||
connectedAccountId: messageChannel.connectedAccount.id,
|
||||
messageChannelId: messageChannel.id,
|
||||
});
|
||||
} catch (error) {
|
||||
await this.messageImportErrorHandlerService.handleDriverException(
|
||||
error,
|
||||
|
||||
+14
-6
@@ -5,6 +5,7 @@ import { Processor } from 'src/engine/core-modules/message-queue/decorators/proc
|
||||
import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants';
|
||||
import { TwentyORMManager } from 'src/engine/twenty-orm/twenty-orm.manager';
|
||||
import { isThrottled } from 'src/modules/connected-account/utils/is-throttled';
|
||||
import { MessageChannelSyncStatusService } from 'src/modules/messaging/common/services/message-channel-sync-status.service';
|
||||
import {
|
||||
MessageChannelSyncStage,
|
||||
type MessageChannelWorkspaceEntity,
|
||||
@@ -24,6 +25,7 @@ export class MessagingMessagesImportJob {
|
||||
constructor(
|
||||
private readonly messagingMessagesImportService: MessagingMessagesImportService,
|
||||
private readonly messagingMonitoringService: MessagingMonitoringService,
|
||||
private readonly messageChannelSyncStatusService: MessageChannelSyncStatusService,
|
||||
private readonly twentyORMManager: TwentyORMManager,
|
||||
) {}
|
||||
|
||||
@@ -64,18 +66,24 @@ export class MessagingMessagesImportJob {
|
||||
}
|
||||
|
||||
if (
|
||||
isThrottled(
|
||||
messageChannel.syncStageStartedAt,
|
||||
messageChannel.throttleFailureCount,
|
||||
)
|
||||
messageChannel.syncStage !==
|
||||
MessageChannelSyncStage.MESSAGES_IMPORT_SCHEDULED
|
||||
) {
|
||||
return;
|
||||
}
|
||||
|
||||
if (
|
||||
messageChannel.syncStage !==
|
||||
MessageChannelSyncStage.MESSAGES_IMPORT_SCHEDULED
|
||||
isThrottled(
|
||||
messageChannel.syncStageStartedAt,
|
||||
messageChannel.throttleFailureCount,
|
||||
)
|
||||
) {
|
||||
await this.messageChannelSyncStatusService.scheduleMessagesImport(
|
||||
[messageChannel.id],
|
||||
workspaceId,
|
||||
true,
|
||||
);
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
|
||||
+6
-4
@@ -53,10 +53,6 @@ export class MessagingOngoingStaleJob {
|
||||
messageChannel.syncStageStartedAt &&
|
||||
isSyncStale(messageChannel.syncStageStartedAt)
|
||||
) {
|
||||
this.logger.log(
|
||||
`Sync for message channel ${messageChannel.id} and workspace ${workspaceId} is stale. Setting sync stage to MESSAGES_IMPORT_PENDING`,
|
||||
);
|
||||
|
||||
await this.messageChannelSyncStatusService.resetSyncStageStartedAt(
|
||||
[messageChannel.id],
|
||||
workspaceId,
|
||||
@@ -65,6 +61,9 @@ export class MessagingOngoingStaleJob {
|
||||
switch (messageChannel.syncStage) {
|
||||
case MessageChannelSyncStage.MESSAGE_LIST_FETCH_ONGOING:
|
||||
case MessageChannelSyncStage.MESSAGE_LIST_FETCH_SCHEDULED:
|
||||
this.logger.log(
|
||||
`Sync for message channel ${messageChannel.id} and workspace ${workspaceId} is stale. Setting sync stage to MESSAGE_LIST_FETCH_PENDING`,
|
||||
);
|
||||
await this.messageChannelSyncStatusService.scheduleMessageListFetch(
|
||||
[messageChannel.id],
|
||||
workspaceId,
|
||||
@@ -72,6 +71,9 @@ export class MessagingOngoingStaleJob {
|
||||
break;
|
||||
case MessageChannelSyncStage.MESSAGES_IMPORT_ONGOING:
|
||||
case MessageChannelSyncStage.MESSAGES_IMPORT_SCHEDULED:
|
||||
this.logger.log(
|
||||
`Sync for message channel ${messageChannel.id} and workspace ${workspaceId} is stale. Setting sync stage to MESSAGES_IMPORT_PENDING`,
|
||||
);
|
||||
await this.messageChannelSyncStatusService.scheduleMessagesImport(
|
||||
[messageChannel.id],
|
||||
workspaceId,
|
||||
|
||||
+2
@@ -166,6 +166,7 @@ export class MessageImportExceptionHandlerService {
|
||||
await this.messageChannelSyncStatusService.scheduleMessageListFetch(
|
||||
[messageChannel.id],
|
||||
workspaceId,
|
||||
true,
|
||||
);
|
||||
break;
|
||||
|
||||
@@ -174,6 +175,7 @@ export class MessageImportExceptionHandlerService {
|
||||
await this.messageChannelSyncStatusService.scheduleMessagesImport(
|
||||
[messageChannel.id],
|
||||
workspaceId,
|
||||
true,
|
||||
);
|
||||
break;
|
||||
|
||||
|
||||
Reference in New Issue
Block a user