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

Method initializeQueueMetadata

common/persistence/sql/queue.go:315–345  ·  view source on GitHub ↗
(
	ctx context.Context,
	blob *commonpb.DataBlob,
)

Source from the content-addressed store, hash-verified

313}
314
315func (q *sqlQueue) initializeQueueMetadata(
316 ctx context.Context,
317 blob *commonpb.DataBlob,
318) error {
319 _, err := q.DB.SelectFromQueueMetadata(ctx, sqlplugin.QueueMetadataFilter{
320 QueueType: q.queueType,
321 })
322 switch err {
323 case nil:
324 return nil
325 case sql.ErrNoRows:
326 result, err := q.DB.InsertIntoQueueMetadata(ctx, &sqlplugin.QueueMetadataRow{
327 QueueType: q.queueType,
328 Data: blob.Data,
329 DataEncoding: blob.EncodingType.String(),
330 })
331 if err != nil {
332 return serviceerror.NewUnavailablef("initializeQueueMetadata operation failed. Error %v", err)
333 }
334 rowsAffected, err := result.RowsAffected()
335 if err != nil {
336 return fmt.Errorf("rowsAffected returned error when initializing queue metadata %v: %v", q.queueType, err)
337 }
338 if rowsAffected != 1 {
339 return fmt.Errorf("rowsAffected returned %v queue metadata instead of one", rowsAffected)
340 }
341 return nil
342 default:
343 return err
344 }
345}
346
347func (q *sqlQueue) initializeDLQMetadata(
348 ctx context.Context,

Callers 1

InitMethod · 0.95

Calls 4

StringMethod · 0.65
ErrorfMethod · 0.45

Tested by

no test coverage detected