Added messageReplaced to concurrency tracking (when freezing)
This commit is contained in:
@@ -507,6 +507,8 @@ export class MarQS {
|
||||
});
|
||||
|
||||
await this.#callEnqueueMessage(newMessage);
|
||||
|
||||
await this.options.subscriber?.messageReplaced(newMessage);
|
||||
},
|
||||
{
|
||||
kind: SpanKind.CONSUMER,
|
||||
|
||||
@@ -97,6 +97,7 @@ export interface MessageQueueSubscriber {
|
||||
messageDequeued(message: MessagePayload): Promise<void>;
|
||||
messageAcked(message: MessagePayload): Promise<void>;
|
||||
messageNacked(message: MessagePayload): Promise<void>;
|
||||
messageReplaced(message: MessagePayload): Promise<void>;
|
||||
}
|
||||
|
||||
export interface VisibilityTimeoutStrategy {
|
||||
|
||||
@@ -97,6 +97,30 @@ class TaskRunConcurrencyTracker implements MessageQueueSubscriber {
|
||||
});
|
||||
}
|
||||
|
||||
async messageReplaced(message: MessagePayload): Promise<void> {
|
||||
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) {
|
||||
|
||||
Reference in New Issue
Block a user