431
ctx context.Context,
432
request *p.RangeCompleteHistoryTasksRequest,
434
>
start := p.UnixMilliseconds(request.InclusiveMinTaskKey.FireTime)
435
>
end := p.UnixMilliseconds(request.ExclusiveMaxTaskKey.FireTime)
436
>
query := d.Session.Query(templateRangeCompleteTimerTaskQuery,
437
>
request.ShardID,
438
>
rowTypeTimerTask,
439
>
rowTypeTimerNamespaceID,
440
>
rowTypeTimerWorkflowID,
441
>
rowTypeTimerRunID,
442
>
start,
443
>
end,
444
>
).WithContext(ctx)
445
>
446
>
err := query.Exec()
447
>
return gocql.ConvertError("RangeCompleteTimerTask", err)
448
>
}
449
450
func (d *MutableStateTaskStore) getReplicationTasks(
451
ctx context.Context,
452
request *p.GetHistoryTasksRequest,
454
>
455
>
// Reading replication tasks need to be quorum level consistent, otherwise we could lose task
456
>
query := d.Session.Query(templateGetReplicationTasksQuery,
457
>
request.ShardID,
458
>
rowTypeReplicationTask,
459
>
rowTypeReplicationNamespaceID,
460
>
rowTypeReplicationWorkflowID,
461
>
rowTypeReplicationRunID,
462
>
defaultVisibilityTimestamp,
463
>
request.InclusiveMinTaskKey.TaskID,
464
>
request.ExclusiveMaxTaskKey.TaskID,
465
>
).WithContext(ctx).PageSize(request.BatchSize).PageState(request.NextPageToken)
466
>
467
>
return d.populateGetReplicationTasksResponse(query, "GetReplicationTasks")
468
>
}
469
470
func (d *MutableStateTaskStore) completeReplicationTask(
471
ctx context.Context,
472
request *p.CompleteHistoryTaskRequest,
474
>
query := d.Session.Query(templateCompleteReplicationTaskQuery,
475
>
request.ShardID,
476
>
rowTypeReplicationTask,
477
>
rowTypeReplicationNamespaceID,
478
>
rowTypeReplicationWorkflowID,
479
>
rowTypeReplicationRunID,
480
>
defaultVisibilityTimestamp,
481
>
request.TaskKey.TaskID,
482
>
).WithContext(ctx)
483
>
484
>
err := query.Exec()
485
>
return gocql.ConvertError("CompleteReplicationTask", err)
486
>
}
487
488
func (d *MutableStateTaskStore) rangeCompleteReplicationTasks(
489
ctx context.Context,
490
request *p.RangeCompleteHistoryTasksRequest,
492
>
query := d.Session.Query(templateRangeCompleteReplicationTaskQuery,
493
>
request.ShardID,
494
>
rowTypeReplicationTask,
495
>
rowTypeReplicationNamespaceID,
496
>
rowTypeReplicationWorkflowID,
497
>
rowTypeReplicationRunID,
498
>
defaultVisibilityTimestamp,
499
>
request.InclusiveMinTaskKey.TaskID,
500
>
request.ExclusiveMaxTaskKey.TaskID,
501
>
).WithContext(ctx)
502
>
503
>
err := query.Exec()
504
>
return gocql.ConvertError("RangeCompleteReplicationTask", err)
505
>
}
506
507
func (d *MutableStateTaskStore) PutReplicationTaskToDLQ(