369
}
370
371
>
func (m *StreamReceiverMonitorImpl) generateStatusMap(inboundKeys map[ClusterShardKeyPair]struct{}) map[ClusterShardKeyPair]*streamStatus {
stream_receiver_monitor.go
372
>
serverToClients := make(map[ClusterShardKey][]ClusterShardKey)
373
>
for keyPair := range inboundKeys {
374
>
serverToClients[keyPair.Server] = append(serverToClients[keyPair.Server], keyPair.Client)
375
>
}
376
>
statusMap := make(map[ClusterShardKeyPair]*streamStatus)
377
>
for serverKey, clientKeys := range serverToClients {
378
>
m.fillStatusMap(statusMap, serverKey, clientKeys)
379
>
}
380
>
return statusMap
381
}
382
383
>
func (m *StreamReceiverMonitorImpl) fillStatusMap(statusMap map[ClusterShardKeyPair]*streamStatus, serverKey ClusterShardKey, clientsKeys []ClusterShardKey) {
stream_receiver_monitor.go
384
>
shardContext, err := m.ShardController.GetShardByID(serverKey.ShardID)
385
>
if err != nil {
386
m.Logger.Error("Failed to get shardContext.", tag.Error(err))
387
return
388
}
390
>
defer cancel()
391
>
engine, err := shardContext.GetEngine(ctx)
392
>
if err != nil {
393
m.Logger.Error("Failed to get engine.", tag.Error(err))
394
return
395
}
397
>
queueState, ok := shardContext.GetQueueState(tasks.CategoryReplication)
398
>
if !ok {
399
m.Logger.Error("Failed to get queue state.")
400
return
401
}
403
>
for _, clientKey := range clientsKeys {
404
>
readerID := shard.ReplicationReaderIDFromClusterShardID(
405
>
int64(clientKey.ClusterID),
406
>
clientKey.ShardID,
407
>
)
408
>
readerState, ok := readerStates[readerID]
409
>
if !ok {
410
m.Logger.Error("Failed to get reader state.")
411
statusMap[ClusterShardKeyPair{Client: clientKey, Server: serverKey}] = &streamStatus{