From f1ffc2e6c820f230a08424ad3a0cdc989433da3d Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Tue, 13 Aug 2024 16:36:39 +0100 Subject: [PATCH] Added messageReplaced to concurrency tracking (when freezing) --- apps/webapp/app/v3/marqs/index.server.ts | 2 ++ apps/webapp/app/v3/marqs/types.ts | 1 + .../taskRunConcurrencyTracker.server.ts | 24 +++++++++++++++++++ 3 files changed, 27 insertions(+) diff --git a/apps/webapp/app/v3/marqs/index.server.ts b/apps/webapp/app/v3/marqs/index.server.ts index bd8247b5a..252723f90 100644 --- a/apps/webapp/app/v3/marqs/index.server.ts +++ b/apps/webapp/app/v3/marqs/index.server.ts @@ -507,6 +507,8 @@ export class MarQS { }); await this.#callEnqueueMessage(newMessage); + + await this.options.subscriber?.messageReplaced(newMessage); }, { kind: SpanKind.CONSUMER, diff --git a/apps/webapp/app/v3/marqs/types.ts b/apps/webapp/app/v3/marqs/types.ts index 559b4b23f..afe6c8fef 100644 --- a/apps/webapp/app/v3/marqs/types.ts +++ b/apps/webapp/app/v3/marqs/types.ts @@ -97,6 +97,7 @@ export interface MessageQueueSubscriber { messageDequeued(message: MessagePayload): Promise; messageAcked(message: MessagePayload): Promise; messageNacked(message: MessagePayload): Promise; + messageReplaced(message: MessagePayload): Promise; } export interface VisibilityTimeoutStrategy { diff --git a/apps/webapp/app/v3/services/taskRunConcurrencyTracker.server.ts b/apps/webapp/app/v3/services/taskRunConcurrencyTracker.server.ts index dba230419..451a5e20b 100644 --- a/apps/webapp/app/v3/services/taskRunConcurrencyTracker.server.ts +++ b/apps/webapp/app/v3/services/taskRunConcurrencyTracker.server.ts @@ -97,6 +97,30 @@ class TaskRunConcurrencyTracker implements MessageQueueSubscriber { }); } + async messageReplaced(message: MessagePayload): Promise { + logger.debug("TaskRunConcurrencyTracker.messageReplaced()", { + data: message.data, + messageId: message.messageId, + }); + + const data = this.getMessageData(message); + if (!data) { + logger.info( + `TaskRunConcurrencyTracker.messageReplaced(): could not parse message data`, + message + ); + return; + } + + await this.executionFinished({ + projectId: data.projectId, + taskId: data.taskIdentifier, + runId: message.messageId, + environmentId: data.environmentId, + deployed: data.environmentType !== "DEVELOPMENT", + }); + } + private getMessageData(message: MessagePayload) { const result = ConcurrentMessageData.safeParse(message.data); if (result.success) {