118
119
// "passive" means the namespace is in standby mode and only replicates data
120
>
namespaceState := metrics.PassiveNamespaceStateTagValue
dlq_writer.go
121
>
if isNamespaceActive {
122
namespaceState = metrics.ActiveNamespaceStateTagValue
123
}
124
>
taskType := GetTaskTypeTagValue(task, isNamespaceActive, q.chasmRegistry)
dlq_writer.go
125
>
metrics.DLQWrites.With(q.metricsHandler).Record(
126
>
1,
127
>
metrics.TaskCategoryTag(task.GetCategory().Name()),
128
>
metrics.NamespaceStateTag(namespaceState),
129
>
nsMetricTag,
130
>
metrics.TaskTypeTag(taskType),
131
>
metrics.OperationTag(taskType),
132
>
getArchetypeTag(task, q.chasmRegistry),
133
>
)
134
>
q.logger.Warn("Task enqueued to DLQ",
135
>
tag.DLQMessageID(resp.Metadata.ID),
136
>
tag.SourceCluster(sourceCluster),
137
>
tag.TargetCluster(targetCluster),
138
>
tag.TaskType(task.GetType()),
139
>
tag.String("task-category", task.GetCategory().Name()),
140
>
nsLogTag,
141
>
)
142
>
return nil
143
}
144
145
// getQueueMutex returns a per-queue mutex, creating it if it doesn't exist.
146
// This provides process-level locking to serialize concurrent writes to the same queue.
147
>
func (q *DLQWriter) getQueueMutex(queueKey persistence.QueueKey) *sync.Mutex {
dlq_writer.go
148
>
if mu, ok := q.enqueueMutex.Load(queueKey); ok {
149
return mu.(*sync.Mutex) //nolint:revive
150
}
151
153
>
actual, _ := q.enqueueMutex.LoadOrStore(queueKey, newMutex)
154
>
return actual.(*sync.Mutex) //nolint:revive
155
}