Throw 400 for webhook triggering deleted workspaces + enqueue workflows by batches (#17154)
Fixes https://github.com/twentyhq/twenty/issues/17027 Fixes https://github.com/twentyhq/twenty/issues/16824
This commit is contained in:
+37
-30
@@ -1,6 +1,6 @@
|
||||
import { Injectable, Logger } from '@nestjs/common';
|
||||
|
||||
import { Not } from 'typeorm';
|
||||
import { QUERY_MAX_RECORDS } from 'twenty-shared/constants';
|
||||
|
||||
import { InjectMessageQueue } from 'src/engine/core-modules/message-queue/decorators/message-queue.decorator';
|
||||
import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants';
|
||||
@@ -90,15 +90,17 @@ export class WorkflowRunEnqueueWorkspaceService {
|
||||
workspaceId,
|
||||
);
|
||||
|
||||
const workflowRunIdsToEnqueue: string[] = [];
|
||||
let totalEnqueuedCount = 0;
|
||||
|
||||
if (remainingWorkflowRunToEnqueueCount > 0) {
|
||||
const additionalRunsToEnqueue = await workflowRunRepository.find({
|
||||
while (remainingWorkflowRunToEnqueueCount > 0) {
|
||||
const batchSize = Math.min(
|
||||
remainingWorkflowRunToEnqueueCount,
|
||||
QUERY_MAX_RECORDS,
|
||||
);
|
||||
|
||||
const batchRuns = await workflowRunRepository.find({
|
||||
where: {
|
||||
status: WorkflowRunStatus.NOT_STARTED,
|
||||
...(workflowRunIdsToEnqueue.length > 0
|
||||
? { id: Not(workflowRunIdsToEnqueue[0]) }
|
||||
: {}),
|
||||
},
|
||||
select: {
|
||||
id: true,
|
||||
@@ -106,17 +108,37 @@ export class WorkflowRunEnqueueWorkspaceService {
|
||||
order: {
|
||||
createdAt: 'ASC',
|
||||
},
|
||||
take: remainingWorkflowRunToEnqueueCount,
|
||||
take: batchSize,
|
||||
});
|
||||
|
||||
workflowRunIdsToEnqueue.push(
|
||||
...additionalRunsToEnqueue.map(
|
||||
(workflowRun: WorkflowRunWorkspaceEntity) => workflowRun.id,
|
||||
),
|
||||
if (batchRuns.length === 0) {
|
||||
break;
|
||||
}
|
||||
|
||||
const batchIds = batchRuns.map(
|
||||
(workflowRun: WorkflowRunWorkspaceEntity) => workflowRun.id,
|
||||
);
|
||||
|
||||
await workflowRunRepository.update(batchIds, {
|
||||
enqueuedAt: new Date().toISOString(),
|
||||
status: WorkflowRunStatus.ENQUEUED,
|
||||
});
|
||||
|
||||
for (const workflowRunId of batchIds) {
|
||||
await this.messageQueueService.add<RunWorkflowJobData>(
|
||||
RunWorkflowJob.name,
|
||||
{
|
||||
workflowRunId,
|
||||
workspaceId,
|
||||
},
|
||||
);
|
||||
}
|
||||
|
||||
totalEnqueuedCount += batchRuns.length;
|
||||
remainingWorkflowRunToEnqueueCount -= batchRuns.length;
|
||||
}
|
||||
|
||||
if (workflowRunIdsToEnqueue.length <= 0) {
|
||||
if (totalEnqueuedCount === 0) {
|
||||
if (!isCacheMode) {
|
||||
await this.workflowThrottlingWorkspaceService.recomputeWorkflowRunNotStartedCount(
|
||||
workspaceId,
|
||||
@@ -126,30 +148,15 @@ export class WorkflowRunEnqueueWorkspaceService {
|
||||
return;
|
||||
}
|
||||
|
||||
await workflowRunRepository.update(workflowRunIdsToEnqueue, {
|
||||
enqueuedAt: new Date().toISOString(),
|
||||
status: WorkflowRunStatus.ENQUEUED,
|
||||
});
|
||||
|
||||
await this.workflowThrottlingWorkspaceService.consumeRemainingRunsToEnqueueCount(
|
||||
workspaceId,
|
||||
workflowRunIdsToEnqueue.length,
|
||||
totalEnqueuedCount,
|
||||
);
|
||||
|
||||
for (const workflowRunId of workflowRunIdsToEnqueue) {
|
||||
await this.messageQueueService.add<RunWorkflowJobData>(
|
||||
RunWorkflowJob.name,
|
||||
{
|
||||
workflowRunId,
|
||||
workspaceId,
|
||||
},
|
||||
);
|
||||
}
|
||||
|
||||
if (isCacheMode) {
|
||||
await this.workflowThrottlingWorkspaceService.decreaseWorkflowRunNotStartedCount(
|
||||
workspaceId,
|
||||
workflowRunIdsToEnqueue.length,
|
||||
totalEnqueuedCount,
|
||||
);
|
||||
} else {
|
||||
await this.workflowThrottlingWorkspaceService.recomputeWorkflowRunNotStartedCount(
|
||||
|
||||
Reference in New Issue
Block a user