Compare commits

...
Author SHA1 Message Date
sonarly-bot 66c777f121 fix(server): prevent released query runner in calendar import
https://sonarly.com/issue/39517?type=bug

Calendar event import fails for affected channels with an internal DB lifecycle error, leaving sync in FAILED state. The failure started right after v2.7.2 and is tied to the calendar worker path.

Fix: I fixed the root cause by restoring workspace ORM context around token persistence in `ConnectedAccountRefreshTokensService`.

What changed:
1) `connected-account-refresh-tokens.service.ts`
- Reintroduced:
  - `GlobalWorkspaceOrmManager`
  - `buildSystemAuthContext`
- Re-added `globalWorkspaceOrmManager` constructor dependency.
- Wrapped `connectedAccountRepository.update(...)` in:
  - `const authContext = buildSystemAuthContext(workspaceId);`
  - `await this.globalWorkspaceOrmManager.executeInWorkspaceContext(async () => { ... }, authContext);`

Why this addresses the bug:
- The regression removed explicit workspace context before repository update in the refresh flow used by calendar import workers.
- Reintroducing `executeInWorkspaceContext` ensures repository operations execute with a valid workspace-bound ORM lifecycle/query runner, preventing the released-query-runner failure in this path.

2) `connected-account-refresh-tokens.service.spec.ts`
- Added `GlobalWorkspaceOrmManager` mock provider.
- Added retrieval of the mock from the testing module.
- Updated refresh-path tests to assert `executeInWorkspaceContext` is called for token refresh persistence.

Authored by Sonarly by autonomous analysis (run 45019).
2026-05-21 17:40:16 +00:00
6 changed files with 229 additions and 49 deletions
@@ -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'>,
@@ -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({
@@ -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,
@@ -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;
@@ -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(
@@ -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,