Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
66c777f121 |
+53
-14
@@ -37,11 +37,22 @@ export class CalendarEventImportErrorHandlerService {
|
||||
) {}
|
||||
|
||||
public async handleDriverException(
|
||||
exception: CalendarEventImportDriverException | TwentyORMException,
|
||||
exception: CalendarEventImportDriverException | TwentyORMException | Error,
|
||||
syncStep: CalendarEventImportSyncStep,
|
||||
calendarChannel: Pick<CalendarChannelEntity, 'id' | 'throttleFailureCount'>,
|
||||
workspaceId: string,
|
||||
): Promise<void> {
|
||||
if (!('code' in exception)) {
|
||||
await this.handleUnknownException(
|
||||
exception,
|
||||
syncStep,
|
||||
calendarChannel,
|
||||
workspaceId,
|
||||
);
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
switch (exception.code) {
|
||||
case CalendarEventImportDriverExceptionCode.NOT_FOUND:
|
||||
await this.handleNotFoundException(
|
||||
@@ -68,11 +79,19 @@ export class CalendarEventImportErrorHandlerService {
|
||||
await this.handleSyncCursorErrorException(calendarChannel, workspaceId);
|
||||
break;
|
||||
case CalendarEventImportDriverExceptionCode.CHANNEL_MISCONFIGURED:
|
||||
await this.handleChannelMisconfiguredException(
|
||||
exception,
|
||||
syncStep,
|
||||
calendarChannel,
|
||||
workspaceId,
|
||||
);
|
||||
break;
|
||||
case CalendarEventImportDriverExceptionCode.UNKNOWN:
|
||||
case CalendarEventImportDriverExceptionCode.UNKNOWN_NETWORK_ERROR:
|
||||
default:
|
||||
await this.handleUnknownException(
|
||||
exception,
|
||||
syncStep,
|
||||
calendarChannel,
|
||||
workspaceId,
|
||||
);
|
||||
@@ -168,7 +187,8 @@ export class CalendarEventImportErrorHandlerService {
|
||||
}
|
||||
|
||||
private async handleUnknownException(
|
||||
exception: { message: string },
|
||||
exception: Error,
|
||||
syncStep: CalendarEventImportSyncStep,
|
||||
calendarChannel: Pick<CalendarChannelEntity, 'id'>,
|
||||
workspaceId: string,
|
||||
): Promise<void> {
|
||||
@@ -182,23 +202,42 @@ export class CalendarEventImportErrorHandlerService {
|
||||
CalendarEventImportExceptionCode.UNKNOWN,
|
||||
);
|
||||
|
||||
this.logger.error(exception);
|
||||
this.exceptionHandlerService.captureExceptions(
|
||||
[calendarEventImportException],
|
||||
{
|
||||
additionalData: {
|
||||
calendarChannelId: calendarChannel.id,
|
||||
exceptionMessage: exception.message,
|
||||
},
|
||||
workspace: {
|
||||
id: workspaceId,
|
||||
},
|
||||
},
|
||||
this.logger.error(
|
||||
exception.message,
|
||||
exception.stack,
|
||||
`CalendarChannelId: ${calendarChannel.id}, WorkspaceId: ${workspaceId}, SyncStep: ${syncStep}`,
|
||||
);
|
||||
|
||||
this.exceptionHandlerService.captureExceptions([exception], {
|
||||
additionalData: {
|
||||
calendarChannelId: calendarChannel.id,
|
||||
exceptionMessage: exception.message,
|
||||
syncStep,
|
||||
},
|
||||
workspace: {
|
||||
id: workspaceId,
|
||||
},
|
||||
});
|
||||
|
||||
throw calendarEventImportException;
|
||||
}
|
||||
|
||||
private async handleChannelMisconfiguredException(
|
||||
exception: CalendarEventImportDriverException,
|
||||
syncStep: CalendarEventImportSyncStep,
|
||||
calendarChannel: Pick<CalendarChannelEntity, 'id'>,
|
||||
workspaceId: string,
|
||||
): Promise<void> {
|
||||
await this.calendarChannelSyncStatusService.markAsFailedUnknownAndFlushCalendarEventsToImport(
|
||||
[calendarChannel.id],
|
||||
workspaceId,
|
||||
);
|
||||
|
||||
this.logger.warn(
|
||||
`Calendar channel ${calendarChannel.id} in workspace ${workspaceId} is misconfigured during ${syncStep}: ${exception.message}`,
|
||||
);
|
||||
}
|
||||
|
||||
private async handleNotFoundException(
|
||||
syncStep: CalendarEventImportSyncStep,
|
||||
calendarChannel: Pick<CalendarChannelEntity, 'id'>,
|
||||
|
||||
+20
@@ -12,6 +12,7 @@ import {
|
||||
ConnectedAccountRefreshAccessTokenExceptionCode,
|
||||
} from 'src/engine/metadata-modules/connected-account/exceptions/connected-account-refresh-tokens.exception';
|
||||
import { ConnectedAccountTokenEncryptionService } from 'src/engine/metadata-modules/connected-account/services/connected-account-token-encryption.service';
|
||||
import { GlobalWorkspaceOrmManager } from 'src/engine/twenty-orm/global-workspace-datasource/global-workspace-orm.manager';
|
||||
import { GoogleAPIRefreshAccessTokenService } from 'src/modules/connected-account/refresh-tokens-manager/drivers/google/services/google-api-refresh-tokens.service';
|
||||
import { MicrosoftAPIRefreshAccessTokenService } from 'src/modules/connected-account/refresh-tokens-manager/drivers/microsoft/services/microsoft-api-refresh-tokens.service';
|
||||
|
||||
@@ -24,6 +25,7 @@ describe('ConnectedAccountRefreshTokensService', () => {
|
||||
let googleAPIRefreshAccessTokenService: GoogleAPIRefreshAccessTokenService;
|
||||
let microsoftAPIRefreshAccessTokenService: MicrosoftAPIRefreshAccessTokenService;
|
||||
let connectedAccountRepository: { update: jest.Mock };
|
||||
let globalWorkspaceOrmManager: { executeInWorkspaceContext: jest.Mock };
|
||||
let connectedAccountTokenEncryptionService: {
|
||||
decrypt: jest.Mock;
|
||||
encryptTokenPair: jest.Mock;
|
||||
@@ -115,6 +117,14 @@ describe('ConnectedAccountRefreshTokensService', () => {
|
||||
provide: ConnectedAccountTokenEncryptionService,
|
||||
useValue: connectedAccountTokenEncryptionService,
|
||||
},
|
||||
{
|
||||
provide: GlobalWorkspaceOrmManager,
|
||||
useValue: {
|
||||
executeInWorkspaceContext: jest.fn(
|
||||
async (operation: () => Promise<void>) => await operation(),
|
||||
),
|
||||
},
|
||||
},
|
||||
],
|
||||
}).compile();
|
||||
|
||||
@@ -132,6 +142,7 @@ describe('ConnectedAccountRefreshTokensService', () => {
|
||||
connectedAccountRepository = module.get(
|
||||
getRepositoryToken(ConnectedAccountEntity),
|
||||
);
|
||||
globalWorkspaceOrmManager = module.get(GlobalWorkspaceOrmManager);
|
||||
});
|
||||
|
||||
afterEach(() => {
|
||||
@@ -199,6 +210,9 @@ describe('ConnectedAccountRefreshTokensService', () => {
|
||||
expect(
|
||||
microsoftAPIRefreshAccessTokenService.refreshTokens,
|
||||
).toHaveBeenCalledWith(mockRefreshTokenPlaintext);
|
||||
expect(globalWorkspaceOrmManager.executeInWorkspaceContext).toHaveBeenCalledTimes(
|
||||
1,
|
||||
);
|
||||
expect(connectedAccountRepository.update).toHaveBeenCalledWith(
|
||||
{ id: mockConnectedAccountId, workspaceId: mockWorkspaceId },
|
||||
expect.objectContaining({
|
||||
@@ -242,6 +256,9 @@ describe('ConnectedAccountRefreshTokensService', () => {
|
||||
expect(
|
||||
googleAPIRefreshAccessTokenService.refreshTokens,
|
||||
).toHaveBeenCalledWith(mockRefreshTokenPlaintext);
|
||||
expect(globalWorkspaceOrmManager.executeInWorkspaceContext).toHaveBeenCalledTimes(
|
||||
1,
|
||||
);
|
||||
expect(connectedAccountRepository.update).toHaveBeenCalledWith(
|
||||
{ id: mockConnectedAccountId, workspaceId: mockWorkspaceId },
|
||||
expect.objectContaining({
|
||||
@@ -285,6 +302,9 @@ describe('ConnectedAccountRefreshTokensService', () => {
|
||||
expect(
|
||||
microsoftAPIRefreshAccessTokenService.refreshTokens,
|
||||
).toHaveBeenCalledWith(mockRefreshTokenPlaintext);
|
||||
expect(globalWorkspaceOrmManager.executeInWorkspaceContext).toHaveBeenCalledTimes(
|
||||
1,
|
||||
);
|
||||
expect(connectedAccountRepository.update).toHaveBeenCalledWith(
|
||||
{ id: mockConnectedAccountId, workspaceId: mockWorkspaceId },
|
||||
expect.objectContaining({
|
||||
|
||||
+15
-8
@@ -12,6 +12,8 @@ import {
|
||||
ConnectedAccountRefreshAccessTokenExceptionCode,
|
||||
} from 'src/engine/metadata-modules/connected-account/exceptions/connected-account-refresh-tokens.exception';
|
||||
import { ConnectedAccountTokenEncryptionService } from 'src/engine/metadata-modules/connected-account/services/connected-account-token-encryption.service';
|
||||
import { GlobalWorkspaceOrmManager } from 'src/engine/twenty-orm/global-workspace-datasource/global-workspace-orm.manager';
|
||||
import { buildSystemAuthContext } from 'src/engine/twenty-orm/utils/build-system-auth-context.util';
|
||||
import { GoogleAPIRefreshAccessTokenService } from 'src/modules/connected-account/refresh-tokens-manager/drivers/google/services/google-api-refresh-tokens.service';
|
||||
import { MicrosoftAPIRefreshAccessTokenService } from 'src/modules/connected-account/refresh-tokens-manager/drivers/microsoft/services/microsoft-api-refresh-tokens.service';
|
||||
|
||||
@@ -32,6 +34,7 @@ export class ConnectedAccountRefreshTokensService {
|
||||
private readonly googleAPIRefreshAccessTokenService: GoogleAPIRefreshAccessTokenService,
|
||||
private readonly microsoftAPIRefreshAccessTokenService: MicrosoftAPIRefreshAccessTokenService,
|
||||
private readonly appOAuthRefreshAccessTokenService: AppOAuthRefreshAccessTokenService,
|
||||
private readonly globalWorkspaceOrmManager: GlobalWorkspaceOrmManager,
|
||||
private readonly connectedAccountTokenEncryptionService: ConnectedAccountTokenEncryptionService,
|
||||
@InjectRepository(ConnectedAccountEntity)
|
||||
private readonly connectedAccountRepository: Repository<ConnectedAccountEntity>,
|
||||
@@ -115,14 +118,18 @@ export class ConnectedAccountRefreshTokensService {
|
||||
workspaceId,
|
||||
});
|
||||
|
||||
await this.connectedAccountRepository.update(
|
||||
{ id: connectedAccount.id, workspaceId },
|
||||
{
|
||||
accessToken: encryptedAccessToken,
|
||||
refreshToken: reEncryptedRefreshToken,
|
||||
lastCredentialsRefreshedAt: new Date(),
|
||||
},
|
||||
);
|
||||
const authContext = buildSystemAuthContext(workspaceId);
|
||||
|
||||
await this.globalWorkspaceOrmManager.executeInWorkspaceContext(async () => {
|
||||
await this.connectedAccountRepository.update(
|
||||
{ id: connectedAccount.id, workspaceId },
|
||||
{
|
||||
accessToken: encryptedAccessToken,
|
||||
refreshToken: reEncryptedRefreshToken,
|
||||
lastCredentialsRefreshedAt: new Date(),
|
||||
},
|
||||
);
|
||||
}, authContext);
|
||||
|
||||
return {
|
||||
accessToken: encryptedAccessToken,
|
||||
|
||||
+43
-7
@@ -27,6 +27,7 @@ import { MessagingMonitoringService } from 'src/modules/messaging/monitoring/ser
|
||||
import { MessageChannelEntity } from 'src/engine/metadata-modules/message-channel/entities/message-channel.entity';
|
||||
import { UserWorkspaceEntity } from 'src/engine/core-modules/user-workspace/user-workspace.entity';
|
||||
import { WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity';
|
||||
import { QueryFailedError } from 'typeorm';
|
||||
|
||||
describe('MessagingMessagesImportService', () => {
|
||||
let service: MessagingMessagesImportService;
|
||||
@@ -48,8 +49,21 @@ describe('MessagingMessagesImportService', () => {
|
||||
>;
|
||||
let mockConnectedAccount: ConnectedAccountEntity;
|
||||
let providersBase: Provider[];
|
||||
const mockWorkspaceMemberRepository = {
|
||||
update: jest.fn().mockResolvedValue(undefined),
|
||||
findOne: jest.fn().mockResolvedValue({
|
||||
id: 'workspace-member-id',
|
||||
userId: 'user-id',
|
||||
}),
|
||||
};
|
||||
|
||||
beforeEach(async () => {
|
||||
mockWorkspaceMemberRepository.findOne.mockReset();
|
||||
mockWorkspaceMemberRepository.findOne.mockResolvedValue({
|
||||
id: 'workspace-member-id',
|
||||
userId: 'user-id',
|
||||
});
|
||||
|
||||
mockConnectedAccount = {
|
||||
id: 'connected-account-id',
|
||||
provider: ConnectedAccountProvider.GOOGLE,
|
||||
@@ -125,13 +139,9 @@ describe('MessagingMessagesImportService', () => {
|
||||
{
|
||||
provide: GlobalWorkspaceOrmManager,
|
||||
useValue: {
|
||||
getRepository: jest.fn().mockResolvedValue({
|
||||
update: jest.fn().mockResolvedValue(undefined),
|
||||
findOne: jest.fn().mockResolvedValue({
|
||||
id: 'workspace-member-id',
|
||||
userId: 'user-id',
|
||||
}),
|
||||
}),
|
||||
getRepository: jest
|
||||
.fn()
|
||||
.mockResolvedValue(mockWorkspaceMemberRepository),
|
||||
executeInWorkspaceContext: jest
|
||||
.fn()
|
||||
.mockImplementation((fn: () => any, _authContext?: any) => fn()),
|
||||
@@ -242,6 +252,32 @@ describe('MessagingMessagesImportService', () => {
|
||||
);
|
||||
});
|
||||
|
||||
it('should retry workspace member lookup when first query fails', async () => {
|
||||
mockWorkspaceMemberRepository.findOne
|
||||
.mockRejectedValueOnce(
|
||||
new QueryFailedError(
|
||||
'SELECT * FROM workspaceMember',
|
||||
[],
|
||||
new Error('Connection terminated unexpectedly'),
|
||||
),
|
||||
)
|
||||
.mockResolvedValueOnce({
|
||||
id: 'workspace-member-id',
|
||||
userId: 'user-id',
|
||||
});
|
||||
|
||||
await service.processMessageBatchImport(
|
||||
mockMessageChannel as MessageChannelEntity,
|
||||
mockConnectedAccount,
|
||||
workspaceId,
|
||||
);
|
||||
|
||||
expect(mockWorkspaceMemberRepository.findOne).toHaveBeenCalledTimes(2);
|
||||
expect(
|
||||
saveMessagesService.saveMessagesAndEnqueueContactCreation,
|
||||
).toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it('should fails if SyncStage is not MESSAGES_IMPORT_SCHEDULED', async () => {
|
||||
mockMessageChannel.syncStage =
|
||||
MessageChannelSyncStage.MESSAGES_IMPORT_PENDING;
|
||||
|
||||
+72
-16
@@ -2,7 +2,7 @@ import { Injectable } from '@nestjs/common';
|
||||
import { InjectRepository } from '@nestjs/typeorm';
|
||||
|
||||
import { isDefined } from 'twenty-shared/utils';
|
||||
import { Repository } from 'typeorm';
|
||||
import { QueryFailedError, Repository } from 'typeorm';
|
||||
|
||||
import { MessageChannelSyncStatus } from 'twenty-shared/types';
|
||||
import { ExceptionHandlerService } from 'src/engine/core-modules/exception-handler/exception-handler.service';
|
||||
@@ -25,6 +25,19 @@ export enum MessageImportSyncStep {
|
||||
MESSAGES_IMPORT_ONGOING = 'MESSAGES_IMPORT_ONGOING',
|
||||
}
|
||||
|
||||
const TRANSIENT_POSTGRES_CONNECTION_ERROR_CODES = new Set<string>([
|
||||
'08000',
|
||||
'08001',
|
||||
'08003',
|
||||
'08004',
|
||||
'08006',
|
||||
'08007',
|
||||
'08P01',
|
||||
'57P01',
|
||||
'57P02',
|
||||
'57P03',
|
||||
]);
|
||||
|
||||
@Injectable()
|
||||
export class MessageImportExceptionHandlerService {
|
||||
constructor(
|
||||
@@ -49,6 +62,17 @@ export class MessageImportExceptionHandlerService {
|
||||
};
|
||||
}
|
||||
|
||||
if (this.isTemporaryException(exception)) {
|
||||
await this.handleTemporaryException(
|
||||
syncStep,
|
||||
messageChannel,
|
||||
workspaceId,
|
||||
exception,
|
||||
);
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
if ('code' in exception) {
|
||||
switch (exception.code) {
|
||||
case MessageImportDriverExceptionCode.NOT_FOUND:
|
||||
@@ -58,20 +82,6 @@ export class MessageImportExceptionHandlerService {
|
||||
workspaceId,
|
||||
);
|
||||
break;
|
||||
case TwentyORMExceptionCode.QUERY_READ_TIMEOUT:
|
||||
case MessageImportDriverExceptionCode.TEMPORARY_ERROR:
|
||||
case MessageNetworkExceptionCode.ECONNABORTED:
|
||||
case MessageNetworkExceptionCode.ENOTFOUND:
|
||||
case MessageNetworkExceptionCode.ECONNRESET:
|
||||
case MessageNetworkExceptionCode.ETIMEDOUT:
|
||||
case MessageNetworkExceptionCode.ERR_NETWORK:
|
||||
await this.handleTemporaryException(
|
||||
syncStep,
|
||||
messageChannel,
|
||||
workspaceId,
|
||||
exception,
|
||||
);
|
||||
break;
|
||||
case MessageImportDriverExceptionCode.INSUFFICIENT_PERMISSIONS:
|
||||
await this.handleInsufficientPermissionsException(
|
||||
messageChannel,
|
||||
@@ -93,14 +103,49 @@ export class MessageImportExceptionHandlerService {
|
||||
exception,
|
||||
messageChannel,
|
||||
workspaceId,
|
||||
syncStep,
|
||||
);
|
||||
break;
|
||||
}
|
||||
} else {
|
||||
await this.handleUnknownException(exception, messageChannel, workspaceId);
|
||||
await this.handleUnknownException(
|
||||
exception,
|
||||
messageChannel,
|
||||
workspaceId,
|
||||
syncStep,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
public isTemporaryException(
|
||||
exception: MessageImportDriverException | Error | TwentyORMException,
|
||||
): boolean {
|
||||
if (exception instanceof QueryFailedError) {
|
||||
const queryFailedError = exception as QueryFailedError & {
|
||||
code?: string;
|
||||
};
|
||||
|
||||
return (
|
||||
isDefined(queryFailedError.code) &&
|
||||
TRANSIENT_POSTGRES_CONNECTION_ERROR_CODES.has(queryFailedError.code)
|
||||
);
|
||||
}
|
||||
|
||||
if (!('code' in exception)) {
|
||||
return false;
|
||||
}
|
||||
|
||||
return [
|
||||
TwentyORMExceptionCode.QUERY_READ_TIMEOUT,
|
||||
MessageImportDriverExceptionCode.TEMPORARY_ERROR,
|
||||
MessageNetworkExceptionCode.ECONNABORTED,
|
||||
MessageNetworkExceptionCode.ENOTFOUND,
|
||||
MessageNetworkExceptionCode.ECONNRESET,
|
||||
MessageNetworkExceptionCode.ETIMEDOUT,
|
||||
MessageNetworkExceptionCode.ERR_NETWORK,
|
||||
].includes(exception.code);
|
||||
}
|
||||
|
||||
private async handleSyncCursorErrorException(
|
||||
messageChannel: Pick<MessageChannelEntity, 'id'>,
|
||||
workspaceId: string,
|
||||
@@ -203,8 +248,19 @@ export class MessageImportExceptionHandlerService {
|
||||
exception: Error,
|
||||
messageChannel: Pick<MessageChannelEntity, 'id'>,
|
||||
workspaceId: string,
|
||||
syncStep: MessageImportSyncStep,
|
||||
): Promise<void> {
|
||||
const exceptionCode =
|
||||
'code' in exception && typeof exception.code === 'string'
|
||||
? exception.code
|
||||
: undefined;
|
||||
|
||||
this.exceptionHandlerService.captureExceptions([exception], {
|
||||
additionalData: {
|
||||
messageChannelId: messageChannel.id,
|
||||
syncStep,
|
||||
exceptionCode,
|
||||
},
|
||||
workspace: { id: workspaceId },
|
||||
});
|
||||
await this.messageChannelSyncStatusService.markAsFailed(
|
||||
|
||||
+26
-4
@@ -2,7 +2,7 @@ import { Injectable, Logger } from '@nestjs/common';
|
||||
import { InjectRepository } from '@nestjs/typeorm';
|
||||
|
||||
import { isDefined } from 'twenty-shared/utils';
|
||||
import { Repository } from 'typeorm';
|
||||
import { QueryFailedError, Repository } from 'typeorm';
|
||||
|
||||
import { InjectCacheStorage } from 'src/engine/core-modules/cache-storage/decorators/cache-storage.decorator';
|
||||
import { CacheStorageService } from 'src/engine/core-modules/cache-storage/services/cache-storage.service';
|
||||
@@ -179,9 +179,10 @@ export class MessagingMessagesImportService {
|
||||
);
|
||||
|
||||
const workspaceMember = userWorkspace
|
||||
? await workspaceMemberRepository.findOne({
|
||||
where: { userId: userWorkspace.userId },
|
||||
})
|
||||
? await this.findWorkspaceMemberWithRetry(
|
||||
workspaceMemberRepository,
|
||||
userWorkspace.userId,
|
||||
)
|
||||
: null;
|
||||
|
||||
const blocklist = workspaceMember
|
||||
@@ -282,6 +283,27 @@ export class MessagingMessagesImportService {
|
||||
);
|
||||
}
|
||||
|
||||
private async findWorkspaceMemberWithRetry(
|
||||
workspaceMemberRepository: Repository<WorkspaceMemberWorkspaceEntity>,
|
||||
userId: string,
|
||||
): Promise<WorkspaceMemberWorkspaceEntity | null> {
|
||||
for (let attempt = 0; attempt < 2; attempt += 1) {
|
||||
try {
|
||||
return await workspaceMemberRepository.findOne({
|
||||
where: { userId },
|
||||
});
|
||||
} catch (error) {
|
||||
const isLastAttempt = attempt === 1;
|
||||
|
||||
if (!(error instanceof QueryFailedError) || isLastAttempt) {
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return null;
|
||||
}
|
||||
|
||||
private async trackMessageImportCompleted(
|
||||
messageChannel: MessageChannelEntity,
|
||||
workspaceId: string,
|
||||
|
||||
Reference in New Issue
Block a user