333
}
334
335
>
func (c *priBacklogManagerImpl) InternalStatus() []*taskqueuespb.InternalTaskQueueStatus {
pri_backlog_manager.go
336
>
currentTaskIDBlock := c.taskWriter.getCurrentTaskIDBlock()
337
>
338
>
c.subqueueLock.Lock()
339
>
defer c.subqueueLock.Unlock()
340
>
341
>
status := make([]*taskqueuespb.InternalTaskQueueStatus, len(c.subqueues))
342
>
backlogCountsBySubqueue := c.db.getApproximateBacklogCountsBySubqueue()
343
>
for i, r := range c.subqueues {
344
>
readLevel, ackLevel := r.getLevels()
345
>
status[i] = &taskqueuespb.InternalTaskQueueStatus{
346
>
ReadLevel: readLevel,
347
>
AckLevel: ackLevel,
348
>
TaskIdBlock: &taskqueuepb.TaskIdBlock{
349
>
StartId: currentTaskIDBlock.start,
350
>
EndId: currentTaskIDBlock.end,
351
>
},
352
>
LoadedTasks: int64(r.getLoadedTasks()),
353
>
MaxReadLevel: c.db.GetMaxReadLevel(subqueueIndex(i)),
354
>
ApproximateBacklogCount: backlogCountsBySubqueue[i],
355
>
BacklogDrained: r.isDrained(),
356
>
}
357
>
}
358
>
return status
359
}
360