307
needWrite := db.lastChange.After(db.lastWrite) || time.Since(db.lastWrite) > ttl/2
308
if !needWrite {
309
>
// If we don't write, though, we wouldn't know if someone else has stolen ownership
db.go
310
>
// momentarily (this could happen due to eventual consistency of membership updates).
311
>
// So instead, do a (cheaper) read to just check the range id.
312
>
return db.verifyOwnershipLocked(ctx)
313
>
}
314
315
return db.updateTaskQueueLocked(ctx, false)
316
}
317
318
>
func (db *taskQueueDB) verifyOwnershipLocked(ctx context.Context) error {
db.go
319
>
response, err := db.store.GetTaskQueue(ctx, &persistence.GetTaskQueueRequest{
320
>
NamespaceID: db.queue.NamespaceId(),
321
>
TaskQueue: db.queue.PersistenceName(),
322
>
TaskType: db.queue.TaskType(),
323
>
})
324
>
if err != nil {
325
return err
326
}
327
>
if response.RangeID != db.rangeID {
db.go
328
return &persistence.ConditionFailedError{
329
Msg: fmt.Sprintf("task queue ownership lost: stored rangeID %d, in-memory rangeID %d",