( ctx context.Context, request *persistence.InternalUpdateTaskQueueRequest, )
| 84 | } |
| 85 | |
| 86 | func (m *taskQueueStore) UpdateTaskQueue( |
| 87 | ctx context.Context, |
| 88 | request *persistence.InternalUpdateTaskQueueRequest, |
| 89 | ) (*persistence.UpdateTaskQueueResponse, error) { |
| 90 | nidBytes, err := primitives.ParseUUID(request.NamespaceID) |
| 91 | if err != nil { |
| 92 | return nil, serviceerror.NewInternal(err.Error()) |
| 93 | } |
| 94 | |
| 95 | tqId, tqHash := taskQueueIdAndHash(nidBytes, request.TaskQueue, request.TaskType, persistence.SubqueueZero) |
| 96 | var resp *persistence.UpdateTaskQueueResponse |
| 97 | err = m.txExecute(ctx, "UpdateTaskQueue", func(tx sqlplugin.Tx) error { |
| 98 | if err := lockTaskQueue(ctx, |
| 99 | tx, |
| 100 | tqHash, |
| 101 | tqId, |
| 102 | request.PrevRangeID, |
| 103 | m.version, |
| 104 | ); err != nil { |
| 105 | return err |
| 106 | } |
| 107 | result, err := tx.UpdateTaskQueues(ctx, &sqlplugin.TaskQueuesRow{ |
| 108 | RangeHash: tqHash, |
| 109 | TaskQueueID: tqId, |
| 110 | RangeID: request.RangeID, |
| 111 | Data: request.TaskQueueInfo.Data, |
| 112 | DataEncoding: request.TaskQueueInfo.EncodingType.String(), |
| 113 | }, m.version) |
| 114 | if err != nil { |
| 115 | return err |
| 116 | } |
| 117 | rowsAffected, err := result.RowsAffected() |
| 118 | if err != nil { |
| 119 | return err |
| 120 | } |
| 121 | if rowsAffected != 1 { |
| 122 | return fmt.Errorf("%v rows were affected instead of 1", rowsAffected) |
| 123 | } |
| 124 | resp = &persistence.UpdateTaskQueueResponse{} |
| 125 | return nil |
| 126 | }) |
| 127 | return resp, err |
| 128 | } |
| 129 | |
| 130 | func (m *taskQueueStore) ListTaskQueue( |
| 131 | ctx context.Context, |
nothing calls this directly
no test coverage detected