go.temporal.io/server/tests/dlq_test.go

610 LOC · 0 covered · 610 uncovered · 0 ranges · 0 concepts · 0 introducers · 0 tests

1 package tests
2
3 import (
4 "bytes"
5 "context"
6 "encoding/base64"
7 "encoding/json"
8 "errors"
9 "io"
10 "os"
11 "strconv"
12 "strings"
13 "sync/atomic"
14 "testing"
15 "time"
16
17 "github.com/google/uuid"
18 "github.com/stretchr/testify/require"
19 enumspb "go.temporal.io/api/enums/v1"
20 sdkclient "go.temporal.io/sdk/client"
21 "go.temporal.io/sdk/workflow"
22 "go.temporal.io/server/api/adminservice/v1"
23 enumsspb "go.temporal.io/server/api/enums/v1"
24 "go.temporal.io/server/api/historyservice/v1"
25 persistencespb "go.temporal.io/server/api/persistence/v1"
26 "go.temporal.io/server/common/codec"
27 "go.temporal.io/server/common/config"
28 "go.temporal.io/server/common/definition"
29 "go.temporal.io/server/common/persistence"
30 "go.temporal.io/server/common/persistence/serialization"
31 "go.temporal.io/server/common/primitives"
32 "go.temporal.io/server/common/testing/await"
33 "go.temporal.io/server/common/testing/parallelsuite"
34 "go.temporal.io/server/common/testing/testhooks"
35 "go.temporal.io/server/service/history/tasks"
36 "go.temporal.io/server/tests/testcore"
37 "go.temporal.io/server/tests/testutils"
38 "go.temporal.io/server/tools/tdbg"
39 "go.temporal.io/server/tools/tdbg/tdbgtest"
40 "google.golang.org/protobuf/encoding/protojson"
41 "google.golang.org/protobuf/proto"
42 )
43
44 type (
45 DLQSuite struct {
46 parallelsuite.Suite[*DLQSuite]
47 }
48 dlqTestEnv struct {
49 *testcore.TestEnv
50
51 dlq persistence.HistoryTaskQueueManager
52 writer bytes.Buffer
53 systemSDKClient sdkclient.Client
54 deleteBlockCh chan any
55
56 failingWorkflowIDPrefix atomic.Pointer[string]
57 }
58 dlqTestCase struct {
59 name string
60 dlqTestParams
61 configure func(*dlqTestParams)
62 }
63 dlqTestParams struct {
64 maxMessageCount string
65 lastMessageID string
66 targetCluster string
67 expectedNumMessages int
68 }
69 )
70
71 func TestDLQSuite(t *testing.T) {
72 parallelsuite.RunLegacySequential(t, &DLQSuite{}) //nolint:staticcheck // SA1019: DLQ tests use dedicated clusters with fault injection and worker-service DLQ jobs.
73 }
74
75 func (s *DLQSuite) newTestEnv(opts ...testcore.TestOption) *dlqTestEnv {
76 w := &dlqTestEnv{}
77 testPrefix := "dlq-test-terminal-wfts-"
78 w.failingWorkflowIDPrefix.Store(&testPrefix)
79
80 baseOpts := []testcore.TestOption{
81 // Return a terminal error that will cause workflow task to be added to the DLQ.
82 testcore.WithPersistenceFaultInjection(&config.FaultInjection{
83 Injector: func(target config.FaultInjectionTarget) error {
84 if target.Store != config.ExecutionStoreName || target.Method != "GetWorkflowExecution" {
85 return nil
86 }
87 request, ok := target.Request.(*persistence.GetWorkflowExecutionRequest)
88 if !ok || !strings.HasPrefix(request.WorkflowID, *w.failingWorkflowIDPrefix.Load()) {
89 return nil
90 }
91 return serialization.NewDeserializationError(enumspb.ENCODING_TYPE_PROTO3, errors.New("test error"))
92 },
93 }),
94 }
95 w.TestEnv = testcore.NewEnv(s.T(), append(baseOpts, opts...)...)
96 w.SdkWorker().RegisterWorkflow(s.myWorkflow)
97
98 var err error
99 w.dlq, err = w.GetTestCluster().TestBase().Factory.NewHistoryTaskQueueManager()
100 s.NoError(err)
101 s.T().Cleanup(w.dlq.Close)
102
103 w.systemSDKClient, err = sdkclient.Dial(sdkclient.Options{
104 HostPort: w.FrontendGRPCAddress(),
105 Namespace: primitives.SystemLocalNamespace,
106 })
107 s.NoError(err)
108 s.T().Cleanup(w.systemSDKClient.Close)
109
110 w.deleteBlockCh = make(chan any)
111 close(w.deleteBlockCh)
112
113 // DeleteDLQTasks is used to block the dlq job workflow until one of them is cancelled in TestCancelRunningMerge.
114 w.InjectHook(testhooks.NewHook(
115 testhooks.HistoryDLQTaskDeleteInterceptor,
116 func(
117 ctx context.Context,
118 request *historyservice.DeleteDLQTasksRequest,
119 deleteTasks func(context.Context, *historyservice.DeleteDLQTasksRequest) (*historyservice.DeleteDLQTasksResponse, error),
120 ) (*historyservice.DeleteDLQTasksResponse, error) {
121 <-w.deleteBlockCh
122 return deleteTasks(ctx, request)
123 },
124 ))
125
126 return w
127 }
128
129 func (env *dlqTestEnv) runTdbg(ctx context.Context, args []string) error {
130 return tdbgtest.NewCliApp(
131 func(params *tdbg.Params) {
132 params.ClientFactory = tdbg.NewClientFactory(tdbg.WithFrontendAddress(env.FrontendGRPCAddress()))
133 params.Writer = &env.writer
134 },
135 ).RunContext(ctx, args)
136 }
137
138 func (s *DLQSuite) myWorkflow(workflow.Context) (string, error) {
139 return "hello", nil
140 }
141
142 func (s *DLQSuite) TestReadArtificialDLQTasks() {
143 env := s.newTestEnv()
144
145 namespaceID := "test-namespace"
146 workflowID := "test-workflow-id"
147 workflowKey := definition.NewWorkflowKey(namespaceID, workflowID, "test-run-id")
148
149 category := tasks.CategoryTransfer
150 sourceCluster := "test-source-cluster-" + s.T().Name()
151
152 // Note: it's ok that this isn't unique across tests because the queue name will still be unique due to the source
153 // cluster name being included in the queue name. We use the current cluster name because that's what the default
154 // is if the target cluster flag isn't specified.
155 targetCluster := "active"
156 queueKey := persistence.QueueKey{
157 QueueType: persistence.QueueTypeHistoryDLQ,
158 Category: category,
159 SourceCluster: sourceCluster,
160 TargetCluster: targetCluster,
161 }
162 _, err := env.dlq.CreateQueue(s.Context(), &persistence.CreateQueueRequest{
163 QueueKey: queueKey,
164 })
165 s.NoError(err)
166 for i := range 4 {
167 task := &tasks.WorkflowTask{
168 WorkflowKey: workflowKey,
169 TaskID: int64(42 + i),
170 }
171 _, err := env.dlq.EnqueueTask(s.Context(), &persistence.EnqueueTaskRequest{
172 QueueType: queueKey.QueueType,
173 SourceCluster: queueKey.SourceCluster,
174 TargetCluster: queueKey.TargetCluster,
175 Task: task,
176 SourceShardID: tasks.GetShardIDForTask(task, int(env.GetTestClusterConfig().HistoryConfig.NumHistoryShards)),
177 })
178 s.NoError(err)
179 }
180 for _, tc := range []dlqTestCase{
181 {
182 name: "max message count exceeded",
183 configure: func(params *dlqTestParams) {
184 params.maxMessageCount = "2"
185 params.lastMessageID = "999"
186 params.expectedNumMessages = 2
187 },
188 },
189 {
190 name: "last message ID exceeded",
191 configure: func(params *dlqTestParams) {
192 params.maxMessageCount = "999"
193 params.lastMessageID = "2" // first message is 0, so this should return 3 messages: 0, 1, 2
194 params.expectedNumMessages = 3
195 },
196 },
197 {
198 name: "target cluster specified",
199 configure: func(params *dlqTestParams) {
200 },
201 },
202 {
203 name: "target cluster not specified",
204 configure: func(params *dlqTestParams) {
205 params.targetCluster = ""
206 },
207 },
208 } {
209 s.Run(tc.name, func(s *DLQSuite) {
210 file := testutils.CreateTemp(s.T(), "", "*")
211 tc.maxMessageCount = "999"
212 tc.lastMessageID = "999"
213 tc.expectedNumMessages = 4
214 tc.targetCluster = targetCluster
215 tc.configure(&tc.dlqTestParams)
216 args := []string{
217 "tdbg",
218 "--" + tdbg.FlagYes,
219 "dlq",
220 "read",
221 "--" + tdbg.FlagDLQType, strconv.Itoa(tasks.CategoryTransfer.ID()),
222 "--" + tdbg.FlagCluster, sourceCluster,
223 "--" + tdbg.FlagPageSize, "1",
224 "--" + tdbg.FlagMaxMessageCount, tc.maxMessageCount,
225 "--" + tdbg.FlagLastMessageID, tc.lastMessageID,
226 "--" + tdbg.FlagOutputFilename, file.Name(),
227 }
228 if tc.targetCluster != "" {
229 args = append(args, "--"+tdbg.FlagTargetCluster, tc.targetCluster)
230 }
231 cmdString := strings.Join(args, " ")
232 s.T().Log("TDBG command:", cmdString)
233 err := env.runTdbg(s.Context(), args)
234 s.NoError(err)
235
236 s.T().Log("TDBG output:")
237 s.T().Log("========================================")
238 output, err := io.ReadAll(file)
239 s.NoError(err)
240 s.T().Log(string(output))
241 _, err = file.Seek(0, io.SeekStart)
242 s.NoError(err)
243 s.T().Log("========================================")
244 s.verifyNumTasks(file, tc.expectedNumMessages)
245 })
246 }
247 }
248
249 // This test executes an actual workflow for which we've set up an executor wrapper to return a terminal error. This
250 // causes the workflow task to be added to the DLQ. This tests the end-to-end functionality of the DLQ, whereas the
251 // above test is more for testing specific CLI flags when reading from the DLQ. After the workflow task is added to the
252 // DLQ, this test then purges the DLQ and verifies that the task was deleted.
253 // This test will then call DescribeDLQJob and CancelDLQJob api to verify.
254 func (s *DLQSuite) TestPurgeRealWorkflow() {
255 env := s.newTestEnv(testcore.WithWorkerService("dlq purge workflow"))
256
257 _, dlqMessageID := s.executeDoomedWorkflow(env)
258
259 // Delete the workflow task from the DLQ.
260 token := s.purgeMessages(env, dlqMessageID)
261
262 // Verify that the workflow task is no longer in the DLQ.
263 dlqTasks := s.readDLQTasks(env)
264 s.Empty(dlqTasks, "expected DLQ to be empty after purge")
265
266 // Run DescribeJob and validate
267 response := s.describeJob(env, token)
268 s.Equal(enumsspb.DLQ_OPERATION_TYPE_PURGE, response.OperationType)
269 s.Equal(enumsspb.DLQ_OPERATION_STATE_COMPLETED, response.OperationState)
270 s.Equal(dlqMessageID, response.MaxMessageId)
271 s.Equal(dlqMessageID, response.LastProcessedMessageId)
272 s.Equal(int64(1), response.MessagesProcessed)
273
274 // Try to cancel completed workflow
275 cancelResponse := s.cancelJob(env, token)
276 s.False(cancelResponse.Canceled)
277 }
278
279 // This test executes actual workflows for which we've set up an executor wrapper to return a terminal error. This
280 // causes the workflow tasks to be added to the DLQ. This tests the end-to-end functionality of the DLQ, whereas the
281 // above test is more for testing specific CLI flags when reading from the DLQ.
282 // This test will then call DescribeDLQJob and CancelDLQJob api to verify.
283 func (s *DLQSuite) TestMergeRealWorkflow() {
284 env := s.newTestEnv(testcore.WithWorkerService("dlq merge workflow"))
285
286 // Verify that we can execute a normal workflow.
287 run := s.executeWorkflow(env, "dlq-test-ok-workflow-id")
288 s.validateWorkflowRun(env, run)
289
290 // Execute several doomed workflows.
291 numWorkflows := 3
292 var dlqMessageID int64
293 var runs []sdkclient.WorkflowRun
294 for range numWorkflows {
295 run, dlqMessageID = s.executeDoomedWorkflow(env)
296 runs = append(runs, run)
297 }
298
299 // Re-enqueue the workflow tasks from the DLQ, but don't fail its WFTs this time.
300 nonExistantID := "some-workflow-id-that-wont-exist"
301 env.failingWorkflowIDPrefix.Store(&nonExistantID)
302 token := s.mergeMessages(env, dlqMessageID)
303
304 // Verify that the workflow task was deleted from the DLQ after merging.
305 dlqTasks := s.readDLQTasks(env)
306 s.Empty(dlqTasks)
307
308 // Verify that the workflows now eventually complete successfully.
309 for i := range numWorkflows {
310 s.validateWorkflowRun(env, runs[i])
311 }
312
313 // Run DescribeJob and validate
314 response := s.describeJob(env, token)
315 s.Equal(enumsspb.DLQ_OPERATION_TYPE_MERGE, response.OperationType)
316 s.Equal(enumsspb.DLQ_OPERATION_STATE_COMPLETED, response.OperationState)
317 s.Equal(dlqMessageID, response.MaxMessageId)
318 s.Equal(dlqMessageID, response.LastProcessedMessageId)
319 s.Equal(int64(numWorkflows), response.MessagesProcessed)
320
321 // Try to cancel completed workflow
322 cancelResponse := s.cancelJob(env, token)
323 s.False(cancelResponse.Canceled)
324 }
325
326 func (s *DLQSuite) TestCancelRunningMerge() {
327 env := s.newTestEnv(testcore.WithWorkerService("dlq merge workflow"))
328 env.deleteBlockCh = make(chan any)
329
330 // Execute several doomed workflows.
331 _, dlqMessageID := s.executeDoomedWorkflow(env)
332
333 token := s.mergeMessagesWithoutBlocking(env, dlqMessageID)
334
335 // Try to cancel running workflow
336 cancelResponse := s.cancelJob(env, token)
337 s.True(cancelResponse.Canceled)
338 // Unblock waiting tests on Delete
339 close(env.deleteBlockCh)
340 // Delete the workflow task from the DLQ.
341 s.purgeMessages(env, dlqMessageID)
342 }
343
344 func (s *DLQSuite) TestListQueues() {
345 env := s.newTestEnv()
346 targetCluster := "active"
347 category := tasks.CategoryTransfer
348 sourceCluster := "test-source-cluster-" + s.T().Name()
349
350 queueKey1 := persistence.QueueKey{
351 QueueType: persistence.QueueTypeHistoryDLQ,
352 Category: category,
353 SourceCluster: sourceCluster + "_1",
354 TargetCluster: targetCluster,
355 }
356 _, err := env.dlq.CreateQueue(s.Context(), &persistence.CreateQueueRequest{
357 QueueKey: queueKey1,
358 })
359 s.NoError(err)
360
361 queueKey2 := persistence.QueueKey{
362 QueueType: persistence.QueueTypeHistoryDLQ,
363 Category: category,
364 SourceCluster: sourceCluster + "_2",
365 TargetCluster: targetCluster,
366 }
367 _, err = env.dlq.CreateQueue(s.Context(), &persistence.CreateQueueRequest{
368 QueueKey: queueKey2,
369 })
370 s.NoError(err)
371
372 // Insert a message to second queue
373 _, err = env.dlq.EnqueueTask(s.Context(), &persistence.EnqueueTaskRequest{
374 QueueType: persistence.QueueTypeHistoryDLQ,
375 SourceCluster: sourceCluster + "_2",
376 TargetCluster: targetCluster,
377 Task: &tasks.WorkflowTask{},
378 SourceShardID: 1,
379 })
380 s.NoError(err)
381
382 queueInfos := s.listQueues(env)
383 qi0 := adminservice.ListQueuesResponse_QueueInfo{
384 QueueName: queueKey1.GetQueueName(),
385 MessageCount: 0,
386 LastMessageId: -1,
387 }
388 qi1 := adminservice.ListQueuesResponse_QueueInfo{
389 QueueName: queueKey2.GetQueueName(),
390 MessageCount: 1,
391 LastMessageId: 0,
392 }
393 var found0, found1 bool
394 for _, qi := range queueInfos {
395 found0 = found0 || proto.Equal(qi, &qi0)
396 found1 = found1 || proto.Equal(qi, &qi1)
397
398 }
399 s.True(found0, "unable to find %v in %v", &qi0, queueInfos)
400 s.True(found1, "unable to find %v in %v", &qi1, queueInfos)
401 }
402
403 func (s *DLQSuite) validateWorkflowRun(env *dlqTestEnv, run sdkclient.WorkflowRun) {
404 var result string
405 err := run.Get(s.Context(), &result)
406 s.NoError(err)
407 s.Equal("hello", result)
408 }
409
410 // executeDoomedWorkflow runs a workflow that is guaranteed to produce a workflow task that will be added to the DLQ. It
411 // then returns the sdk workflow run and the message ID of the DLQ message for the failed workflow task.
412 func (s *DLQSuite) executeDoomedWorkflow(env *dlqTestEnv) (sdkclient.WorkflowRun, int64) {
413 // Execute a workflow.
414 // Use a random workflow ID to ensure that we don't have any collisions with other runs.
415 run := s.executeWorkflow(env, *env.failingWorkflowIDPrefix.Load()+uuid.NewString())
416
417 // Wait for the workflow task to be added to the DLQ.
418 var found *tdbgtest.DLQMessage[*persistencespb.TransferTaskInfo]
419 await.Require(s.Context(), s.T(), func(t *await.T) {
420 dlqTasks := s.readDLQTasks(env)
421 for _, task := range dlqTasks {
422 if task.Payload.RunId == run.GetRunID() {
423 found = &task
424 return
425 }
426 }
427 require.Failf(t, "workflow task not found in DLQ", "run ID: %s", run.GetRunID())
428 }, 10*time.Second, 100*time.Millisecond)
429
430 return run, found.MessageID
431 }
432
433 // executeWorkflow just executes a simple no-op workflow that returns "hello" and returns the sdk workflow run.
434 func (s *DLQSuite) executeWorkflow(env *dlqTestEnv, workflowID string) sdkclient.WorkflowRun {
435 run, err := env.SdkClient().ExecuteWorkflow(s.Context(), sdkclient.StartWorkflowOptions{
436 ID: workflowID,
437 TaskQueue: env.WorkerTaskQueue(),
438 }, s.myWorkflow)
439 s.NoError(err)
440 return run
441 }
442
443 // purgeMessages from the DLQ up to and including the specified message ID, blocking until the purge workflow completes.
444 func (s *DLQSuite) purgeMessages(env *dlqTestEnv, maxMessageIDToDelete int64) string {
445 args := []string{
446 "tdbg",
447 "--" + tdbg.FlagYes,
448 "dlq",
449 "purge",
450 "--" + tdbg.FlagDLQType, strconv.Itoa(tasks.CategoryTransfer.ID()),
451 "--" + tdbg.FlagLastMessageID, strconv.FormatInt(maxMessageIDToDelete, 10),
452 }
453 err := env.runTdbg(s.Context(), args)
454 s.NoError(err)
455 output := env.writer.Bytes()
456 env.writer.Truncate(0)
457 var data map[string]string
458 err = json.Unmarshal(output, &data)
459 s.NoError(err)
460 tokenString := data["jobToken"]
461
462 var response adminservice.PurgeDLQTasksResponse
463 s.NoError(protojson.Unmarshal(output, &response))
464 var token adminservice.DLQJobToken
465 s.NoError(proto.Unmarshal(response.GetJobToken(), &token))
466
467 run := env.systemSDKClient.GetWorkflow(s.Context(), token.WorkflowId, token.RunId)
468 s.NoError(run.Get(s.Context(), nil))
469 return tokenString
470 }
471
472 // mergeMessages from the DLQ up to and including the specified message ID, blocking until the merge workflow completes.
473 func (s *DLQSuite) mergeMessages(env *dlqTestEnv, maxMessageID int64) string {
474 tokenString := s.mergeMessagesWithoutBlocking(env, maxMessageID)
475 tokenBytes, err := base64.StdEncoding.DecodeString(tokenString)
476 s.NoError(err)
477 var token adminservice.DLQJobToken
478 s.NoError(token.Unmarshal(tokenBytes))
479 run := env.systemSDKClient.GetWorkflow(s.Context(), token.WorkflowId, token.RunId)
480 s.NoError(run.Get(s.Context(), nil))
481 return tokenString
482 }
483
484 // mergeMessages from the DLQ up to and including the specified message ID, returns immediately after running tdbg command.
485 func (s *DLQSuite) mergeMessagesWithoutBlocking(env *dlqTestEnv, maxMessageID int64) string {
486 args := []string{
487 "tdbg",
488 "--" + tdbg.FlagYes,
489 "dlq",
490 "merge",
491 "--" + tdbg.FlagDLQType, strconv.Itoa(tasks.CategoryTransfer.ID()),
492 "--" + tdbg.FlagLastMessageID, strconv.FormatInt(maxMessageID, 10),
493 "--" + tdbg.FlagPageSize, "1", // to ensure that we test pagination
494 }
495 err := env.runTdbg(s.Context(), args)
496 s.NoError(err)
497 output := env.writer.Bytes()
498 env.writer.Truncate(0)
499 var data map[string]string
500 err = json.Unmarshal(output, &data)
501 s.NoError(err)
502 tokenString := data["jobToken"]
503 var response adminservice.MergeDLQTasksResponse
504 s.NoError(protojson.Unmarshal(output, &response))
505 return tokenString
506 }
507
508 // readDLQTasks from the transfer task DLQ for this cluster and return them.
509 func (s *DLQSuite) readDLQTasks(env *dlqTestEnv) []tdbgtest.DLQMessage[*persistencespb.TransferTaskInfo] {
510 file := testutils.CreateTemp(s.T(), "", "*")
511 args := []string{
512 "tdbg",
513 "--" + tdbg.FlagYes,
514 "dlq",
515 "read",
516 "--" + tdbg.FlagDLQType, strconv.Itoa(tasks.CategoryTransfer.ID()),
517 "--" + tdbg.FlagOutputFilename, file.Name(),
518 }
519 s.NoError(env.runTdbg(s.Context(), args))
520 dlqTasks := s.readTransferTasks(file)
521 return dlqTasks
522 }
523
524 // Calls describe dlq job and verify the output
525 func (s *DLQSuite) describeJob(env *dlqTestEnv, token string) *adminservice.DescribeDLQJobResponse {
526 args := []string{
527 "tdbg",
528 "dlq",
529 "job",
530 "describe",
531 "--" + tdbg.FlagJobToken, token,
532 }
533 err := env.runTdbg(s.Context(), args)
534 s.NoError(err)
535 output := env.writer.Bytes()
536 s.T().Log(string(output))
537 env.writer.Truncate(0)
538 var response adminservice.DescribeDLQJobResponse
539 s.NoError(protojson.Unmarshal(output, &response))
540 return &response
541 }
542
543 // Calls delete dlq job and verify the output
544 func (s *DLQSuite) cancelJob(env *dlqTestEnv, token string) *adminservice.CancelDLQJobResponse {
545 args := []string{
546 "tdbg",
547 "dlq",
548 "job",
549 "cancel",
550 "--" + tdbg.FlagJobToken, token,
551 "--" + tdbg.FlagReason, "testing cancel",
552 }
553 err := env.runTdbg(s.Context(), args)
554 s.NoError(err)
555 output := env.writer.Bytes()
556 s.T().Log(string(output))
557 env.writer.Truncate(0)
558 var response adminservice.CancelDLQJobResponse
559 s.NoError(protojson.Unmarshal(output, &response))
560 return &response
561 }
562
563 // List all queues
564 func (s *DLQSuite) listQueues(env *dlqTestEnv) []*adminservice.ListQueuesResponse_QueueInfo {
565 args := []string{
566 "tdbg",
567 "dlq",
568 "list",
569 "--" + tdbg.FlagPrintJSON,
570 }
571
572 err := env.runTdbg(s.Context(), args)
573 s.NoError(err)
574 b := env.writer.Bytes()
575 env.writer.Truncate(0)
576 var arr []*adminservice.ListQueuesResponse_QueueInfo
577 jsonpb := codec.NewJSONPBEncoder()
578 err = jsonpb.DecodeSlice(b, func() proto.Message {
579 resp := &adminservice.ListQueuesResponse_QueueInfo{}
580 arr = append(arr, resp)
581 return resp
582 })
583 s.NoError(err)
584 return arr
585 }
586
587 // verifyNumTasks verifies that the specified file contains the expected number of DLQ tasks, and that each task has the
588 // expected metadata and payload.
589 func (s *DLQSuite) verifyNumTasks(file *os.File, expectedNumTasks int) {
590 dlqTasks := s.readTransferTasks(file)
591 s.Len(dlqTasks, expectedNumTasks)
592
593 for i, task := range dlqTasks {
594 s.Equal(int64(persistence.FirstQueueMessageID+i), task.MessageID)
595 taskInfo := task.Payload
596 s.Equal(enumsspb.TASK_TYPE_TRANSFER_WORKFLOW_TASK, taskInfo.TaskType)
597 s.Equal("test-namespace", taskInfo.NamespaceId)
598 s.Equal("test-workflow-id", taskInfo.WorkflowId)
599 s.Equal("test-run-id", taskInfo.RunId)
600 s.Equal(int64(42+i), taskInfo.TaskId)
601 }
602 }
603
604 func (s *DLQSuite) readTransferTasks(file *os.File) []tdbgtest.DLQMessage[*persistencespb.TransferTaskInfo] {
605 dlqTasks, err := tdbgtest.ParseDLQMessages(file, func() *persistencespb.TransferTaskInfo {
606 return new(persistencespb.TransferTaskInfo)
607 })
608 s.NoError(err)
609 return dlqTasks
610 }