( ctx context.Context, blob *commonpb.DataBlob, )
| 161 | } |
| 162 | |
| 163 | func (q *sqlQueue) EnqueueMessageToDLQ( |
| 164 | ctx context.Context, |
| 165 | blob *commonpb.DataBlob, |
| 166 | ) (int64, error) { |
| 167 | var lastMessageID int64 |
| 168 | err := q.txExecute(ctx, "EnqueueMessageToDLQ", func(tx sqlplugin.Tx) error { |
| 169 | var err error |
| 170 | lastMessageID, err = tx.GetLastEnqueuedMessageIDForUpdate(ctx, q.getDLQTypeFromQueueType()) |
| 171 | switch err { |
| 172 | case nil: |
| 173 | _, err = tx.InsertIntoMessages(ctx, []sqlplugin.QueueMessageRow{ |
| 174 | newQueueRow(q.getDLQTypeFromQueueType(), lastMessageID+1, blob), |
| 175 | }) |
| 176 | return err |
| 177 | case sql.ErrNoRows: |
| 178 | _, err = tx.InsertIntoMessages(ctx, []sqlplugin.QueueMessageRow{ |
| 179 | newQueueRow(q.getDLQTypeFromQueueType(), persistence.EmptyQueueMessageID+1, blob), |
| 180 | }) |
| 181 | return err |
| 182 | default: |
| 183 | return fmt.Errorf("failed to get last enqueued message id from DLQ: %v", err) |
| 184 | } |
| 185 | }) |
| 186 | if err != nil { |
| 187 | return persistence.EmptyQueueMessageID, serviceerror.NewUnavailable(err.Error()) |
| 188 | } |
| 189 | return lastMessageID + 1, nil |
| 190 | } |
| 191 | |
| 192 | func (q *sqlQueue) ReadMessagesFromDLQ( |
| 193 | ctx context.Context, |
nothing calls this directly
no test coverage detected