MCPcopy Create free account
hub / github.com/temporalio/temporal / ListTaskQueue

Method ListTaskQueue

common/persistence/sql/task_queues.go:130–230  ·  view source on GitHub ↗
(
	ctx context.Context,
	request *persistence.ListTaskQueueRequest,
)

Source from the content-addressed store, hash-verified

128}
129
130func (m *taskQueueStore) ListTaskQueue(
131 ctx context.Context,
132 request *persistence.ListTaskQueueRequest,
133) (*persistence.InternalListTaskQueueResponse, error) {
134 pageToken := taskQueuePageToken{MinTaskQueueId: minTaskQueueId}
135 if request.PageToken != nil {
136 if err := gobDeserialize(request.PageToken, &pageToken); err != nil {
137 return nil, serviceerror.NewInternalf("error deserializing page token: %v", err)
138 }
139 }
140 var err error
141 var rows []sqlplugin.TaskQueuesRow
142 var shardGreaterThan uint32
143 var shardLessThan uint32
144
145 i := uint32(0)
146 if pageToken.MinRangeHash > 0 {
147 // Resume partition position from page token, if exists, before entering loop
148 i = getPartitionForRangeHash(pageToken.MinRangeHash, m.taskScanPartitions)
149 }
150
151 lastPageFull := !bytes.Equal(pageToken.MinTaskQueueId, minTaskQueueId)
152 for ; i < m.taskScanPartitions; i++ {
153 // Get start/end boundaries for partition
154 shardGreaterThan, shardLessThan = getBoundariesForPartition(i, m.taskScanPartitions)
155
156 // If page token hash is greater than the boundaries for this partition, use the pageToken hash for resume point
157 if pageToken.MinRangeHash > shardGreaterThan {
158 shardGreaterThan = pageToken.MinRangeHash
159 }
160
161 filter := sqlplugin.TaskQueuesFilter{
162 RangeHashGreaterThanEqualTo: shardGreaterThan,
163 RangeHashLessThanEqualTo: shardLessThan,
164 TaskQueueIDGreaterThan: minTaskQueueId,
165 PageSize: &request.PageSize,
166 }
167
168 if lastPageFull {
169 // Use page token TaskQueueID filter for this query and set this to false
170 // in order for the next partition so we don't miss any results.
171 filter.TaskQueueIDGreaterThan = pageToken.MinTaskQueueId
172 lastPageFull = false
173 }
174
175 rows, err = m.DB.SelectFromTaskQueues(ctx, filter, m.version)
176 if err != nil {
177 return nil, serviceerror.NewUnavailable(err.Error())
178 }
179
180 if len(rows) > 0 {
181 break
182 }
183 }
184
185 maxRangeHash := uint32(0)
186 resp := &persistence.InternalListTaskQueueResponse{
187 Items: make([]*persistence.InternalListTaskQueueItem, len(rows)),

Callers

nothing calls this directly

Calls 8

NewDataBlobFunction · 0.92
gobDeserializeFunction · 0.85
getPartitionForRangeHashFunction · 0.85
gobSerializeFunction · 0.85
EqualMethod · 0.65
SelectFromTaskQueuesMethod · 0.65
ErrorMethod · 0.65

Tested by

no test coverage detected