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

Method getQueueMetadata

common/persistence/sql/queue_v2.go:318–348  ·  view source on GitHub ↗
(
	ctx context.Context,
	tc sqlplugin.TableCRUD,
	queueType persistence.QueueV2Type,
	queueName string,
)

Source from the content-addressed store, hash-verified

316}
317
318func (q *queueV2) getQueueMetadata(
319 ctx context.Context,
320 tc sqlplugin.TableCRUD,
321 queueType persistence.QueueV2Type,
322 queueName string,
323) (*persistencespb.Queue, error) {
324
325 filter := sqlplugin.QueueV2MetadataFilter{
326 QueueType: queueType,
327 QueueName: queueName,
328 }
329 var (
330 metadata *sqlplugin.QueueV2MetadataRow
331 err error
332 )
333 switch tc.(type) {
334 case sqlplugin.Tx:
335 metadata, err = tc.SelectFromQueueV2MetadataForUpdate(ctx, filter)
336 default:
337 metadata, err = tc.SelectFromQueueV2Metadata(ctx, filter)
338 }
339 if err != nil {
340 if errors.Is(err, sql.ErrNoRows) {
341 return nil, persistence.NewQueueNotFoundError(queueType, queueName)
342 }
343 return nil, serviceerror.NewUnavailablef(
344 "failed to get metadata for queue with type: %v and name: %v. Error: %v", queueType, queueName, err,
345 )
346 }
347 return q.extractQueueMetadata(metadata)
348}
349
350func (q queueV2) extractQueueMetadata(metadataRow *sqlplugin.QueueV2MetadataRow) (*persistencespb.Queue, error) {
351 if metadataRow.MetadataEncoding != enumspb.ENCODING_TYPE_PROTO3.String() {

Callers 3

EnqueueMessageMethod · 0.95
ReadMessagesMethod · 0.95
RangeDeleteMessagesMethod · 0.95

Tested by

no test coverage detected