messaging fix relaunch cron jobs (#19492)

This commit is contained in:
neo773
2026-04-09 17:55:12 +05:30
committed by GitHub
parent 1577a1933f
commit 7bd84a6029
2 changed files with 70 additions and 42 deletions
@@ -1,7 +1,7 @@
import { InjectDataSource, InjectRepository } from '@nestjs/typeorm';
import { InjectRepository } from '@nestjs/typeorm';
import { WorkspaceActivationStatus } from 'twenty-shared/workspace';
import { DataSource, Repository } from 'typeorm';
import { In, Repository } from 'typeorm';
import {
CalendarChannelSyncStage,
@@ -14,8 +14,8 @@ import { Process } from 'src/engine/core-modules/message-queue/decorators/proces
import { Processor } from 'src/engine/core-modules/message-queue/decorators/processor.decorator';
import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants';
import { MessageQueueService } from 'src/engine/core-modules/message-queue/services/message-queue.service';
import { CalendarChannelEntity } from 'src/engine/metadata-modules/calendar-channel/entities/calendar-channel.entity';
import { WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity';
import { getWorkspaceSchemaName } from 'src/engine/workspace-datasource/utils/get-workspace-schema-name.util';
import {
CalendarRelaunchFailedCalendarChannelJob,
type CalendarRelaunchFailedCalendarChannelJobData,
@@ -29,10 +29,10 @@ export class CalendarRelaunchFailedCalendarChannelsCronJob {
constructor(
@InjectRepository(WorkspaceEntity)
private readonly workspaceRepository: Repository<WorkspaceEntity>,
@InjectRepository(CalendarChannelEntity)
private readonly calendarChannelRepository: Repository<CalendarChannelEntity>,
@InjectMessageQueue(MessageQueue.calendarQueue)
private readonly messageQueueService: MessageQueueService,
@InjectDataSource()
private readonly coreDataSource: DataSource,
private readonly exceptionHandlerService: ExceptionHandlerService,
) {}
@@ -48,27 +48,41 @@ export class CalendarRelaunchFailedCalendarChannelsCronJob {
},
});
for (const activeWorkspace of activeWorkspaces) {
const activeWorkspaceIds = activeWorkspaces.map(
(workspace) => workspace.id,
);
if (activeWorkspaceIds.length === 0) {
return;
}
const failedCalendarChannels = await this.calendarChannelRepository
.find({
where: {
syncStage: CalendarChannelSyncStage.FAILED,
syncStatus: CalendarChannelSyncStatus.FAILED_UNKNOWN,
workspaceId: In(activeWorkspaceIds),
},
})
.catch((error) => {
this.exceptionHandlerService.captureExceptions([error]);
return [];
});
for (const calendarChannel of failedCalendarChannels) {
try {
const schemaName = getWorkspaceSchemaName(activeWorkspace.id);
const failedCalendarChannels = await this.coreDataSource.query(
`SELECT * FROM ${schemaName}."calendarChannel" WHERE "syncStage" = '${CalendarChannelSyncStage.FAILED}' AND "syncStatus" = '${CalendarChannelSyncStatus.FAILED_UNKNOWN}'`,
await this.messageQueueService.add<CalendarRelaunchFailedCalendarChannelJobData>(
CalendarRelaunchFailedCalendarChannelJob.name,
{
workspaceId: calendarChannel.workspaceId,
calendarChannelId: calendarChannel.id,
},
);
for (const calendarChannel of failedCalendarChannels) {
await this.messageQueueService.add<CalendarRelaunchFailedCalendarChannelJobData>(
CalendarRelaunchFailedCalendarChannelJob.name,
{
workspaceId: activeWorkspace.id,
calendarChannelId: calendarChannel.id,
},
);
}
} catch (error) {
this.exceptionHandlerService.captureExceptions([error], {
workspace: {
id: activeWorkspace.id,
id: calendarChannel.workspaceId,
},
});
}
@@ -1,7 +1,7 @@
import { InjectDataSource, InjectRepository } from '@nestjs/typeorm';
import { InjectRepository } from '@nestjs/typeorm';
import { WorkspaceActivationStatus } from 'twenty-shared/workspace';
import { DataSource, Repository } from 'typeorm';
import { In, Repository } from 'typeorm';
import {
MessageChannelSyncStage,
@@ -14,8 +14,8 @@ import { Process } from 'src/engine/core-modules/message-queue/decorators/proces
import { Processor } from 'src/engine/core-modules/message-queue/decorators/processor.decorator';
import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants';
import { MessageQueueService } from 'src/engine/core-modules/message-queue/services/message-queue.service';
import { MessageChannelEntity } from 'src/engine/metadata-modules/message-channel/entities/message-channel.entity';
import { WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity';
import { getWorkspaceSchemaName } from 'src/engine/workspace-datasource/utils/get-workspace-schema-name.util';
import {
MessagingRelaunchFailedMessageChannelJob,
type MessagingRelaunchFailedMessageChannelJobData,
@@ -29,10 +29,10 @@ export class MessagingRelaunchFailedMessageChannelsCronJob {
constructor(
@InjectRepository(WorkspaceEntity)
private readonly workspaceRepository: Repository<WorkspaceEntity>,
@InjectRepository(MessageChannelEntity)
private readonly messageChannelRepository: Repository<MessageChannelEntity>,
@InjectMessageQueue(MessageQueue.messagingQueue)
private readonly messageQueueService: MessageQueueService,
@InjectDataSource()
private readonly coreDataSource: DataSource,
private readonly exceptionHandlerService: ExceptionHandlerService,
) {}
@@ -48,27 +48,41 @@ export class MessagingRelaunchFailedMessageChannelsCronJob {
},
});
for (const activeWorkspace of activeWorkspaces) {
const activeWorkspaceIds = activeWorkspaces.map(
(workspace) => workspace.id,
);
if (activeWorkspaceIds.length === 0) {
return;
}
const failedMessageChannels = await this.messageChannelRepository
.find({
where: {
syncStage: MessageChannelSyncStage.FAILED,
syncStatus: MessageChannelSyncStatus.FAILED_UNKNOWN,
workspaceId: In(activeWorkspaceIds),
},
})
.catch((error) => {
this.exceptionHandlerService.captureExceptions([error]);
return [];
});
for (const messageChannel of failedMessageChannels) {
try {
const schemaName = getWorkspaceSchemaName(activeWorkspace.id);
const failedMessageChannels = await this.coreDataSource.query(
`SELECT * FROM ${schemaName}."messageChannel" WHERE "syncStage" = '${MessageChannelSyncStage.FAILED}' AND "syncStatus" = '${MessageChannelSyncStatus.FAILED_UNKNOWN}'`,
await this.messageQueueService.add<MessagingRelaunchFailedMessageChannelJobData>(
MessagingRelaunchFailedMessageChannelJob.name,
{
workspaceId: messageChannel.workspaceId,
messageChannelId: messageChannel.id,
},
);
for (const messageChannel of failedMessageChannels) {
await this.messageQueueService.add<MessagingRelaunchFailedMessageChannelJobData>(
MessagingRelaunchFailedMessageChannelJob.name,
{
workspaceId: activeWorkspace.id,
messageChannelId: messageChannel.id,
},
);
}
} catch (error) {
this.exceptionHandlerService.captureExceptions([error], {
workspace: {
id: activeWorkspace.id,
id: messageChannel.workspaceId,
},
});
}