Compare commits

...
Author SHA1 Message Date
Sonarly Claude Code 15ccc1ce8d fix(messaging): create repository inside transaction to prevent connection reuse issues
https://sonarly.com/issue/35824?type=bug

Message channel cleanup jobs fail when the repository is created outside the transaction scope, causing connection reuse issues when BullMQ jobs run long enough to lose their lock.

Fix: Changed the transaction handling in `MessagingMessageCleanerService` and `CalendarEventCleanerService` to use explicit query runner pattern with proper lifecycle management.

**Before:** The code used `workspaceDataSource.transaction()` but created repositories BEFORE the transaction started. When BullMQ jobs ran long enough to lose their lock, connections could be returned to the pool in a corrupted state, causing subsequent `START TRANSACTION` calls to fail with "current transaction is aborted".

**After:** Using the explicit pattern:
1. `createQueryRunner()` - creates a dedicated query runner
2. `queryRunner.connect()` - acquires a connection from the pool
3. `queryRunner.startTransaction()` - starts the transaction on that connection
4. `try/catch/finally` with `commitTransaction()`, `rollbackTransaction()`, and `release()`
5. Repository creation happens INSIDE the try block to ensure consistent connection context

This pattern ensures:
- Proper cleanup even on error (connection always released back to pool)
- Transaction is explicitly rolled back on error before releasing
- Connection state is guaranteed clean when returned to pool

The same fix was applied to both:
- `messaging-message-cleaner.service.ts` (deleteMessageChannelMessageAssociationsByChannelId, cleanOrphanMessagesAndThreads)
- `calendar-event-cleaner.service.ts` (deleteCalendarChannelEventAssociationsByChannelId)
2026-05-07 17:49:21 +00:00
2 changed files with 139 additions and 100 deletions
@@ -25,17 +25,23 @@ export class CalendarEventCleanerService {
const authContext = buildSystemAuthContext(workspaceId);
await this.globalWorkspaceOrmManager.executeInWorkspaceContext(async () => {
const calendarChannelEventAssociationRepository =
await this.globalWorkspaceOrmManager.getRepository(
workspaceId,
'calendarChannelEventAssociation',
);
const workspaceDataSource =
await this.globalWorkspaceOrmManager.getGlobalWorkspaceDataSource();
await workspaceDataSource.transaction(async (manager) => {
const transactionManager = manager as WorkspaceEntityManager;
const queryRunner = workspaceDataSource.createQueryRunner();
await queryRunner.connect();
await queryRunner.startTransaction();
try {
const transactionManager = queryRunner.manager as WorkspaceEntityManager;
// Create repository inside the transaction to ensure it uses the same connection
const calendarChannelEventAssociationRepository =
await this.globalWorkspaceOrmManager.getRepository(
workspaceId,
'calendarChannelEventAssociation',
);
await deleteUsingPagination(
workspaceId,
@@ -73,7 +79,14 @@ export class CalendarEventCleanerService {
},
transactionManager,
);
});
await queryRunner.commitTransaction();
} catch (error) {
await queryRunner.rollbackTransaction();
throw error;
} finally {
await queryRunner.release();
}
}, authContext);
}
@@ -130,17 +130,23 @@ export class MessagingMessageCleanerService {
const authContext = buildSystemAuthContext(workspaceId);
await this.globalWorkspaceOrmManager.executeInWorkspaceContext(async () => {
const messageChannelMessageAssociationRepository =
await this.globalWorkspaceOrmManager.getRepository<MessageChannelMessageAssociationWorkspaceEntity>(
workspaceId,
'messageChannelMessageAssociation',
);
const workspaceDataSource =
await this.globalWorkspaceOrmManager.getGlobalWorkspaceDataSource();
await workspaceDataSource.transaction(async (manager) => {
const transactionManager = manager as WorkspaceEntityManager;
const queryRunner = workspaceDataSource.createQueryRunner();
await queryRunner.connect();
await queryRunner.startTransaction();
try {
const transactionManager = queryRunner.manager as WorkspaceEntityManager;
// Create repository inside the transaction to ensure it uses the same connection
const messageChannelMessageAssociationRepository =
await this.globalWorkspaceOrmManager.getRepository<MessageChannelMessageAssociationWorkspaceEntity>(
workspaceId,
'messageChannelMessageAssociation',
);
await deleteUsingPagination(
workspaceId,
@@ -178,7 +184,14 @@ export class MessagingMessageCleanerService {
},
transactionManager,
);
});
await queryRunner.commitTransaction();
} catch (error) {
await queryRunner.rollbackTransaction();
throw error;
} finally {
await queryRunner.release();
}
}, authContext);
}
@@ -186,96 +199,109 @@ export class MessagingMessageCleanerService {
const authContext = buildSystemAuthContext(workspaceId);
await this.globalWorkspaceOrmManager.executeInWorkspaceContext(async () => {
const messageThreadRepository =
await this.globalWorkspaceOrmManager.getRepository<MessageThreadWorkspaceEntity>(
workspaceId,
'messageThread',
);
const messageRepository =
await this.globalWorkspaceOrmManager.getRepository<MessageWorkspaceEntity>(
workspaceId,
'message',
);
const workspaceDataSource =
await this.globalWorkspaceOrmManager.getGlobalWorkspaceDataSource();
await workspaceDataSource.transaction(
async (transactionManager: WorkspaceEntityManager) => {
await deleteUsingPagination(
workspaceId,
500,
async (
limit: number,
offset: number,
_workspaceId: string,
transactionManager: WorkspaceEntityManager,
) => {
const nonAssociatedMessages = await messageRepository.find(
{
where: {
messageChannelMessageAssociations: {
id: IsNull(),
},
},
take: limit,
skip: offset,
relations: ['messageChannelMessageAssociations'],
},
transactionManager,
);
const queryRunner = workspaceDataSource.createQueryRunner();
return nonAssociatedMessages.map(({ id }) => id);
},
async (
ids: string[],
workspaceId: string,
transactionManager?: WorkspaceEntityManager,
) => {
this.logger.debug(
`WorkspaceId: ${workspaceId} Deleting ${ids.length} messages from message cleaner`,
);
await messageRepository.delete(ids, transactionManager);
},
transactionManager,
await queryRunner.connect();
await queryRunner.startTransaction();
try {
const transactionManager = queryRunner.manager as WorkspaceEntityManager;
// Create repositories inside the transaction to ensure they use the same connection
const messageThreadRepository =
await this.globalWorkspaceOrmManager.getRepository<MessageThreadWorkspaceEntity>(
workspaceId,
'messageThread',
);
await deleteUsingPagination(
const messageRepository =
await this.globalWorkspaceOrmManager.getRepository<MessageWorkspaceEntity>(
workspaceId,
500,
async (
limit: number,
offset: number,
_workspaceId: string,
transactionManager?: WorkspaceEntityManager,
) => {
const orphanThreads = await messageThreadRepository.find(
{
where: {
messages: {
id: IsNull(),
},
},
take: limit,
skip: offset,
},
transactionManager,
);
return orphanThreads.map(({ id }) => id);
},
async (
ids: string[],
_workspaceId: string,
transactionManager?: WorkspaceEntityManager,
) => {
await messageThreadRepository.delete(ids, transactionManager);
},
transactionManager,
'message',
);
},
);
await deleteUsingPagination(
workspaceId,
500,
async (
limit: number,
offset: number,
_workspaceId: string,
transactionManager: WorkspaceEntityManager,
) => {
const nonAssociatedMessages = await messageRepository.find(
{
where: {
messageChannelMessageAssociations: {
id: IsNull(),
},
},
take: limit,
skip: offset,
relations: ['messageChannelMessageAssociations'],
},
transactionManager,
);
return nonAssociatedMessages.map(({ id }) => id);
},
async (
ids: string[],
workspaceId: string,
transactionManager?: WorkspaceEntityManager,
) => {
this.logger.debug(
`WorkspaceId: ${workspaceId} Deleting ${ids.length} messages from message cleaner`,
);
await messageRepository.delete(ids, transactionManager);
},
transactionManager,
);
await deleteUsingPagination(
workspaceId,
500,
async (
limit: number,
offset: number,
_workspaceId: string,
transactionManager?: WorkspaceEntityManager,
) => {
const orphanThreads = await messageThreadRepository.find(
{
where: {
messages: {
id: IsNull(),
},
},
take: limit,
skip: offset,
},
transactionManager,
);
return orphanThreads.map(({ id }) => id);
},
async (
ids: string[],
_workspaceId: string,
transactionManager?: WorkspaceEntityManager,
) => {
await messageThreadRepository.delete(ids, transactionManager);
},
transactionManager,
);
await queryRunner.commitTransaction();
} catch (error) {
await queryRunner.rollbackTransaction();
throw error;
} finally {
await queryRunner.release();
}
}, authContext);
}
}