( ctx context.Context, blob *commonpb.DataBlob, )
| 313 | } |
| 314 | |
| 315 | func (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 | |
| 347 | func (q *sqlQueue) initializeDLQMetadata( |
| 348 | ctx context.Context, |
no test coverage detected