794
ctx context.Context,
795
request *p.GetHistoryTasksRequest,
797
>
// execution manager should already validated the request
798
>
// Reading history tasks need to be quorum level consistent, otherwise we could lose task
799
>
800
>
minTimestamp := p.UnixMilliseconds(request.InclusiveMinTaskKey.FireTime)
801
>
maxTimestamp := p.UnixMilliseconds(request.ExclusiveMaxTaskKey.FireTime)
802
>
query := d.Session.Query(templateGetHistoryScheduledTasksQuery,
803
>
request.ShardID,
804
>
request.TaskCategory.ID(),
805
>
rowTypeHistoryTaskNamespaceID,
806
>
rowTypeHistoryTaskWorkflowID,
807
>
rowTypeHistoryTaskRunID,
808
>
minTimestamp,
809
>
maxTimestamp,
810
>
).WithContext(ctx)
811
>
812
>
iter := query.PageSize(request.BatchSize).PageState(request.NextPageToken).Iter()
813
>
814
>
response := &p.InternalGetHistoryTasksResponse{}
815
>
var timestamp time.Time
816
>
var taskID int64
817
>
var data []byte
818
>
var encoding string
819
>
820
>
for iter.Scan(×tamp, &taskID, &data, &encoding) {
821
>
response.Tasks = append(response.Tasks, p.InternalHistoryTask{
822
>
Key: tasks.NewKey(timestamp, taskID),
823
>
Blob: p.NewDataBlob(data, encoding),
824
>
})
825
>
826
>
timestamp = time.Time{}
827
>
taskID = 0
828
>
data = nil
829
>
encoding = ""
830
>
}
831
>
if len(iter.PageState()) > 0 {
832
response.NextPageToken = iter.PageState()
833
}
834
836
return nil, gocql.ConvertError("GetHistoryScheduledTasks", err)
837
}
838
840
}
841