208
}
209
210
>
func (m *StreamReceiverMonitorImpl) generateOutboundStreamKeys() map[ClusterShardKeyPair]struct{} {
stream_receiver_monitor.go
211
>
allClusterInfo := m.ClusterMetadata.GetAllClusterInfo()
212
>
213
>
clientClusterID := int32(m.ClusterMetadata.GetClusterID())
214
>
serverClusterIDs := make(map[int32]struct{})
215
>
clusterIDToShardCount := make(map[int32]int32)
216
>
for _, clusterInfo := range allClusterInfo {
217
>
clusterIDToShardCount[int32(clusterInfo.InitialFailoverVersion)] = clusterInfo.ShardCount
218
>
219
>
if !clusterInfo.Enabled || !cluster.IsReplicationEnabledForCluster(clusterInfo, m.Config.EnableSeparateReplicationEnableFlag()) || int32(clusterInfo.InitialFailoverVersion) == clientClusterID {
220
>
continue
221
}
222
serverClusterIDs[int32(clusterInfo.InitialFailoverVersion)] = struct{}{}
223
}
225
>
for _, shardID := range m.ShardController.ShardIDs() {
226
for serverClusterID := range serverClusterIDs {
227
clientShardID := shardID