138
}
139
141
>
state, err := w.renewLeaseWithRetry(foreverRetryPolicy, common.IsPersistenceTransientError)
142
>
if err != nil {
143
w.backlogMgr.initState(taskQueueState{}, err)
144
return err
145
}
146
>
w.taskIDBlock = rangeIDToTaskIDBlock(state.rangeID, w.config.RangeSize)
pri_task_writer.go
147
>
w.currentTaskIDBlock = w.taskIDBlock
148
>
w.backlogMgr.initState(state, nil)
149
>
return nil
150
}
151
153
>
if w.initState() != nil {
154
return
155
}
156
158
>
for {
159
>
atomic.StoreInt64(&w.currentTaskIDBlock.start, w.taskIDBlock.start)
160
>
atomic.StoreInt64(&w.currentTaskIDBlock.end, w.taskIDBlock.end)
161
>
162
>
select {
163
case request := <-w.appendCh:
164
// read a batch of requests from the channel