339
}
340
341
>
func (c *fairBacklogManagerImpl) InternalStatus() []*taskqueuespb.InternalTaskQueueStatus {
fair_backlog_manager.go
342
>
currentTaskIDBlock := c.taskWriter.getCurrentTaskIDBlock()
343
>
344
>
c.subqueueLock.Lock()
345
>
defer c.subqueueLock.Unlock()
346
>
347
>
status := make([]*taskqueuespb.InternalTaskQueueStatus, len(c.subqueues))
348
>
for i, r := range c.subqueues {
349
>
readLevel, ackLevel := r.getLevels()
350
>
count, maxReadLevel := c.db.getApproximateBacklogCountAndMaxReadLevel(subqueueIndex(i))
351
>
status[i] = &taskqueuespb.InternalTaskQueueStatus{
352
>
FairReadLevel: readLevel.toProto(),
353
>
FairAckLevel: ackLevel.toProto(),
354
>
TaskIdBlock: &taskqueuepb.TaskIdBlock{
355
>
StartId: currentTaskIDBlock.start,
356
>
EndId: currentTaskIDBlock.end,
357
>
},
358
>
LoadedTasks: int64(r.getLoadedTasks()),
359
>
FairMaxReadLevel: maxReadLevel.toProto(),
360
>
ApproximateBacklogCount: count,
361
>
BacklogDrained: r.isDrained(),
362
>
}
363
>
}
364
>
return status
365
}
366