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

Method initializeDLQMetadata

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

Source from the content-addressed store, hash-verified

345}
346
347func (q *sqlQueue) initializeDLQMetadata(
348 ctx context.Context,
349 blob *commonpb.DataBlob,
350) error {
351 _, err := q.DB.SelectFromQueueMetadata(ctx, sqlplugin.QueueMetadataFilter{
352 QueueType: q.getDLQTypeFromQueueType(),
353 })
354 switch err {
355 case nil:
356 return nil
357 case sql.ErrNoRows:
358 result, err := q.DB.InsertIntoQueueMetadata(ctx, &sqlplugin.QueueMetadataRow{
359 QueueType: q.getDLQTypeFromQueueType(),
360 Data: blob.Data,
361 DataEncoding: blob.EncodingType.String(),
362 })
363 if err != nil {
364 return serviceerror.NewUnavailablef("initializeDLQMetadata operation failed. Error %v", err)
365 }
366 rowsAffected, err := result.RowsAffected()
367 if err != nil {
368 return fmt.Errorf("rowsAffected returned error when initializing DLQ metadata %v: %v", q.queueType, err)
369 }
370 if rowsAffected != 1 {
371 return fmt.Errorf("rowsAffected returned %v DLQ metadata instead of one", rowsAffected)
372 }
373 return nil
374 default:
375 return err
376 }
377}
378
379func newQueueRow(
380 queueType persistence.QueueType,

Callers 1

InitMethod · 0.95

Calls 5

StringMethod · 0.65
ErrorfMethod · 0.45

Tested by

no test coverage detected