( ctx context.Context, tx sqlplugin.Tx, tqHash uint32, tqId []byte, oldRangeID int64, v sqlplugin.MatchingTaskVersion, )
| 185 | } |
| 186 | |
| 187 | func lockTaskQueue( |
| 188 | ctx context.Context, |
| 189 | tx sqlplugin.Tx, |
| 190 | tqHash uint32, |
| 191 | tqId []byte, |
| 192 | oldRangeID int64, |
| 193 | v sqlplugin.MatchingTaskVersion, |
| 194 | ) error { |
| 195 | rangeID, err := tx.LockTaskQueues(ctx, sqlplugin.TaskQueuesFilter{ |
| 196 | RangeHash: tqHash, |
| 197 | TaskQueueID: tqId, |
| 198 | }, v) |
| 199 | switch err { |
| 200 | case nil: |
| 201 | if rangeID != oldRangeID { |
| 202 | return &persistence.ConditionFailedError{ |
| 203 | Msg: fmt.Sprintf("Task queue range ID was %v when it was should have been %v", rangeID, oldRangeID), |
| 204 | } |
| 205 | } |
| 206 | return nil |
| 207 | |
| 208 | case sql.ErrNoRows: |
| 209 | return &persistence.ConditionFailedError{Msg: "Task queue does not exists"} |
| 210 | |
| 211 | default: |
| 212 | return serviceerror.NewUnavailablef("Failed to lock task queue. Error: %v", err) |
| 213 | } |
| 214 | } |
no test coverage detected