namespace_replication_queue.go ×8

Frontier kind: Code frontier

unlabeled · c_e70468fddca9

7 tests · 4282 LOC · 178 files · introduces 0 tests · 161 LOC · 5 files

Introduces — evidence that enters the hierarchy at this concept

Code
24 ranges161 lines · 5 files
Tests
0 tests

Contains — complete concept membership

All code (extent)
788 ranges4282 lines · 178 files · Browse complete extent
All tests (intent)
7 testsBrowse complete intent

Neighbourhood graph

The orange circle is the focus. Violet and green circles are every ancestor and descendant, broader and narrower, at any distance; blue squares and pink diamonds are the introduced files and exact introduced tests of every visible concept, not only the focus's. Arrows point from broader to narrower concepts and bridge only concepts omitted from this view. Undirected links show source or test introduction. Concept and file size follows LOC; exact test nodes use test-count units.

Introduced files, introduced tests, and structurally relevant concept specialization

In the embedded map, ordinary wheel input scrolls the page; use the visible controls to zoom and drag to pan. Open the full-screen map for canvas navigation: wheel pans, Ctrl/Command plus wheel zooms, and arrow keys pan when this region is focused. On touch screens, open the full-screen map to pan or pinch. If JavaScript or WebGL is unavailable, use the native relationship evidence on this page.

Graph controls are ready.

Interactive rendering requires JavaScript and WebGL. Use the native relationship evidence on this page while the interactive map is unavailable.

Native relationship evidence

Every exact file and test below is linked only from the concept that introduces it.

Introduced tests

Every collected test enters the hierarchy at exactly one concept.

No tests are introduced at this concept. Its intent tests are introduced by other concepts.

Introduced code

Every collected source range enters the hierarchy at exactly one concept.

5 files ranked by introduced lines: 161 introduced LOC across 24 ranges. Expand a file to inspect source; the > gutter marks introduced lines.

go.temporal.io/server/common/persistence/persistence-tests/queue_persistence.go 58 introduced LOC · 3 ranges

Open complete file

124
125 // TestNamespaceReplicationDLQ tests namespace DLQ operations
126 > func (s *QueuePersistenceSuite) TestNamespaceReplicationDLQ() { queue_persistence.go
127 > maxMessageID := int64(100)
128 > numMessages := 100
129 > concurrentSenders := 10
130 >
131 > messageChan := make(chan *replicationspb.ReplicationTask)
132 >
133 > taskType := enumsspb.REPLICATION_TASK_TYPE_NAMESPACE_TASK
134 > go func() {
135 > for i := range numMessages {
136 > messageChan <- &replicationspb.ReplicationTask{
137 > TaskType: taskType,
138 > Attributes: &replicationspb.ReplicationTask_NamespaceTaskAttributes{
139 > NamespaceTaskAttributes: &replicationspb.NamespaceTaskAttributes{
140 > Id: fmt.Sprintf("message-%v", i),
141 > },
142 > },
143 > }
144 > }
145 > close(messageChan)
146 }()
147
148 > wg := sync.WaitGroup{} queue_persistence.go
149 > wg.Add(concurrentSenders)
150 >
151 > for i := range concurrentSenders {
152 > go func(senderNum int) {
153 > defer wg.Done()
154 > for message := range messageChan {
155 > err := s.PublishToNamespaceDLQ(s.ctx, message)
156 > id := message.Attributes.(*replicationspb.ReplicationTask_NamespaceTaskAttributes).NamespaceTaskAttributes.Id
157 > s.Nil(err, "Enqueue message failed when sender %d tried to send %s", senderNum, id)
158 > }
159 }(i)
160 }
161
162 > wg.Wait() queue_persistence.go
163 >
164 > result1, token, err := s.GetMessagesFromNamespaceDLQ(s.ctx, persistence.EmptyQueueMessageID, maxMessageID, numMessages/2, nil)
165 > s.Nil(err, "GetReplicationMessages failed.")
166 > s.NotNil(token)
167 > result2, token, err := s.GetMessagesFromNamespaceDLQ(s.ctx, persistence.EmptyQueueMessageID, maxMessageID, numMessages, token)
168 > s.Nil(err, "GetReplicationMessages failed.")
169 > s.Equal(len(token), 0)
170 > s.Equal(len(result1)+len(result2), numMessages)
171 > _, _, err = s.GetMessagesFromNamespaceDLQ(s.ctx, persistence.EmptyQueueMessageID, 1<<63-1, numMessages, nil)
172 > s.NoError(err, "GetReplicationMessages failed.")
173 > s.Equal(len(token), 0)
174 >
175 > lastMessageID := result2[len(result2)-1].SourceTaskId
176 > err = s.DeleteMessageFromNamespaceDLQ(s.ctx, lastMessageID)
177 > s.NoError(err)
178 > result3, token, err := s.GetMessagesFromNamespaceDLQ(s.ctx, persistence.EmptyQueueMessageID, maxMessageID, numMessages, token)
179 > s.Nil(err, "GetReplicationMessages failed.")
180 > s.Equal(len(token), 0)
181 > s.Equal(len(result3), numMessages-1)
182 >
183 > err = s.RangeDeleteMessagesFromNamespaceDLQ(s.ctx, persistence.EmptyQueueMessageID, lastMessageID)
184 > s.NoError(err)
185 > result4, token, err := s.GetMessagesFromNamespaceDLQ(s.ctx, persistence.EmptyQueueMessageID, maxMessageID, numMessages, token)
186 > s.Nil(err, "GetReplicationMessages failed.")
187 > s.Equal(len(token), 0)
188 > s.Equal(len(result4), 0)
189 }
190
go.temporal.io/server/common/persistence/namespace_replication_queue.go 27 introduced LOC · 8 ranges

Open complete file

101 }
102
103 > func (q *namespaceReplicationQueueImpl) PublishToDLQ(ctx context.Context, task *replicationspb.ReplicationTask) error { namespace_replication_queue.go
104 > blob, err := q.serializer.ReplicationTaskToBlob(task)
105 > if err != nil {
106 return fmt.Errorf("failed to encode message: %v", err)
107 }
108 > messageID, err := q.queue.EnqueueMessageToDLQ(ctx, blob) namespace_replication_queue.go
109 > if err != nil {
110 return err
111 }
112
113 > metrics.NamespaceReplicationDLQMaxLevelGauge.With(q.metricsHandler). namespace_replication_queue.go
114 > Record(float64(messageID), metrics.OperationTag(metrics.PersistenceNamespaceReplicationQueueScope))
115 > return nil
116 }
117
262 pageSize int,
263 pageToken []byte,
264 > ) ([]*replicationspb.ReplicationTask, []byte, error) { namespace_replication_queue.go
265 >
266 > messages, token, err := q.queue.ReadMessagesFromDLQ(ctx, firstMessageID, lastMessageID, pageSize, pageToken)
267 > if err != nil {
268 return nil, nil, err
269 }
270
271 > var replicationTasks []*replicationspb.ReplicationTask namespace_replication_queue.go
272 > for _, message := range messages {
273 > replicationTask, err := q.serializer.ReplicationTaskFromBlob(NewDataBlob(message.Data, message.Encoding))
274 > if err != nil {
275 return nil, nil, fmt.Errorf("failed to decode dlq task: %v", err)
276 }
277
278 // Overwrite to local cluster message id
279 > replicationTask.SourceTaskId = message.ID namespace_replication_queue.go
280 > replicationTasks = append(replicationTasks, replicationTask)
281 }
282
283 > return replicationTasks, token, nil namespace_replication_queue.go
284 }
285
314 firstMessageID int64,
315 lastMessageID int64,
317 >
318 > return q.queue.RangeDeleteMessagesFromDLQ(
319 > ctx,
320 > firstMessageID,
321 > lastMessageID,
322 > )
323 > }
324
325 func (q *namespaceReplicationQueueImpl) DeleteMessageFromDLQ(
go.temporal.io/server/common/persistence/persistence_metric_clients.go 27 introduced LOC · 3 ranges

Open complete file

1132 ctx context.Context,
1133 blob *commonpb.DataBlob,
1134 > ) (_ int64, retErr error) { persistence_metric_clients.go
1135 > caller := headers.GetCallerInfo(ctx).CallerName
1136 > startTime := time.Now().UTC()
1137 > defer func() {
1138 > p.healthSignals.Record(CallerSegmentMissing, time.Since(startTime), retErr)
1139 > p.recordRequestMetrics(metrics.PersistenceEnqueueMessageToDLQScope, caller, time.Since(startTime), retErr)
1140 > p.recordDataLossMetrics(metrics.PersistenceEnqueueMessageToDLQScope, caller, retErr, "", "")
1141 > }()
1142 > return p.persistence.EnqueueMessageToDLQ(ctx, blob)
1143 }
1144
1149 pageSize int,
1150 pageToken []byte,
1151 > ) (_ []*QueueMessage, _ []byte, retErr error) { persistence_metric_clients.go
1152 > caller := headers.GetCallerInfo(ctx).CallerName
1153 > startTime := time.Now().UTC()
1154 > defer func() {
1155 > p.healthSignals.Record(CallerSegmentMissing, time.Since(startTime), retErr)
1156 > p.recordRequestMetrics(metrics.PersistenceReadMessagesFromDLQScope, caller, time.Since(startTime), retErr)
1157 > p.recordDataLossMetrics(metrics.PersistenceReadMessagesFromDLQScope, caller, retErr, "", "")
1158 > }()
1159 > return p.persistence.ReadMessagesFromDLQ(ctx, firstMessageID, lastMessageID, pageSize, pageToken)
1160 }
1161
1178 firstMessageID int64,
1179 lastMessageID int64,
1180 > ) (retErr error) { persistence_metric_clients.go
1181 > caller := headers.GetCallerInfo(ctx).CallerName
1182 > startTime := time.Now().UTC()
1183 > defer func() {
1184 > p.healthSignals.Record(CallerSegmentMissing, time.Since(startTime), retErr)
1185 > p.recordRequestMetrics(metrics.PersistenceRangeDeleteMessagesFromDLQScope, caller, time.Since(startTime), retErr)
1186 > p.recordDataLossMetrics(metrics.PersistenceRangeDeleteMessagesFromDLQScope, caller, retErr, "", "")
1187 > }()
1188 > return p.persistence.RangeDeleteMessagesFromDLQ(ctx, firstMessageID, lastMessageID)
1189 }
1190
go.temporal.io/server/common/persistence/persistence-tests/persistence_test_base.go 25 introduced LOC · 4 ranges

Open complete file

372
373 // PublishToNamespaceDLQ is a utility method to add messages to the namespace DLQ
374 > func (s *TestBase) PublishToNamespaceDLQ(ctx context.Context, task *replicationspb.ReplicationTask) error { persistence_test_base.go
375 > retryPolicy := backoff.NewExponentialRetryPolicy(100 * time.Millisecond).
376 > WithBackoffCoefficient(1.5).
377 > WithMaximumAttempts(20)
378 >
379 > return backoff.ThrottleRetryContext(
380 > ctx,
381 > func(ctx context.Context) error {
382 > return s.NamespaceReplicationQueue.PublishToDLQ(ctx, task)
383 > },
384 retryPolicy,
385 func(e error) bool {
395 pageSize int,
396 pageToken []byte,
397 > ) ([]*replicationspb.ReplicationTask, []byte, error) { persistence_test_base.go
398 > return s.NamespaceReplicationQueue.GetMessagesFromDLQ(
399 > ctx,
400 > firstMessageID,
401 > lastMessageID,
402 > pageSize,
403 > pageToken,
404 > )
405 > }
406
407 // UpdateNamespaceDLQAckLevel updates namespace dlq ack level
424 ctx context.Context,
425 messageID int64,
426 > ) error { persistence_test_base.go
427 > return s.NamespaceReplicationQueue.DeleteMessageFromDLQ(ctx, messageID)
428 > }
429
430 // RangeDeleteMessagesFromNamespaceDLQ deletes messages from namespace DLQ
433 firstMessageID int64,
434 lastMessageID int64,
435 > ) error { persistence_test_base.go
436 > return s.NamespaceReplicationQueue.RangeDeleteMessagesFromDLQ(ctx, firstMessageID, lastMessageID)
437 > }
438
439 func GenerateRandomDBName() string {
go.temporal.io/server/common/persistence/persistence_retryable_clients.go 24 introduced LOC · 6 ranges

Open complete file

1120 ctx context.Context,
1121 blob *commonpb.DataBlob,
1122 > ) (int64, error) { persistence_retryable_clients.go
1123 > var response int64
1124 > op := func(ctx context.Context) error {
1125 > var err error
1126 > response, err = p.persistence.EnqueueMessageToDLQ(ctx, blob)
1127 > return err
1128 > }
1129
1130 > err := backoff.ThrottleRetryContext(ctx, op, p.policy, p.isRetryable) persistence_retryable_clients.go
1131 > return response, err
1132 }
1133
1138 pageSize int,
1139 pageToken []byte,
1140 > ) ([]*QueueMessage, []byte, error) { persistence_retryable_clients.go
1141 > var messages []*QueueMessage
1142 > var nextPageToken []byte
1143 > op := func(ctx context.Context) error {
1144 > var err error
1145 > messages, nextPageToken, err = p.persistence.ReadMessagesFromDLQ(ctx, firstMessageID, lastMessageID, pageSize, pageToken)
1146 > return err
1147 > }
1148
1149 > err := backoff.ThrottleRetryContext(ctx, op, p.policy, p.isRetryable) persistence_retryable_clients.go
1150 > return messages, nextPageToken, err
1151 }
1152
1155 firstMessageID int64,
1156 lastMessageID int64,
1158 > op := func(ctx context.Context) error {
1159 > return p.persistence.RangeDeleteMessagesFromDLQ(ctx, firstMessageID, lastMessageID)
1160 > }
1161
1162 > return backoff.ThrottleRetryContext(ctx, op, p.policy, p.isRetryable) persistence_retryable_clients.go
1163 }
1164 func (p *queueRetryablePersistenceClient) UpdateDLQAckLevel(