171
}
172
173
>
func (m *StreamReceiverMonitorImpl) generateInboundStreamKeys() map[ClusterShardKeyPair]struct{} {
stream_receiver_monitor.go
174
>
allClusterInfo := m.ClusterMetadata.GetAllClusterInfo()
175
>
176
>
clientClusterIDs := make(map[int32]struct{})
177
>
serverClusterID := int32(m.ClusterMetadata.GetClusterID())
178
>
clusterIDToShardCount := make(map[int32]int32)
179
>
for _, clusterInfo := range allClusterInfo {
180
>
clusterIDToShardCount[int32(clusterInfo.InitialFailoverVersion)] = clusterInfo.ShardCount
181
>
182
>
if !cluster.IsReplicationEnabledForCluster(clusterInfo, m.Config.EnableSeparateReplicationEnableFlag()) || int32(clusterInfo.InitialFailoverVersion) == serverClusterID {
183
>
continue
184
}
185
clientClusterIDs[int32(clusterInfo.InitialFailoverVersion)] = struct{}{}
186
}
188
>
for _, shardID := range m.ShardController.ShardIDs() {
189
for clientClusterID := range clientClusterIDs {
190
serverShardID := shardID