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
}