(t *testing.T)
| 153 | } |
| 154 | |
| 155 | func TestCluster_LoadAndProbe(t *testing.T) { |
| 156 | ctx := context.Background() |
| 157 | ns := "test-ns" |
| 158 | clusterName := "test-clusterProbe" |
| 159 | cluster, err := store.NewCluster(clusterName, []string{"127.0.0.1:7770", "127.0.0.1:7771"}, 2) |
| 160 | require.NoError(t, err) |
| 161 | |
| 162 | nodes := make([]*store.ClusterNode, 0) |
| 163 | for _, shard := range cluster.Shards { |
| 164 | for _, node := range shard.Nodes { |
| 165 | clusterNode, _ := node.(*store.ClusterNode) |
| 166 | nodes = append(nodes, clusterNode) |
| 167 | } |
| 168 | } |
| 169 | require.NoError(t, cluster.Reset(ctx)) |
| 170 | defer func() { |
| 171 | require.NoError(t, cluster.Reset(ctx)) |
| 172 | }() |
| 173 | |
| 174 | s := NewMockClusterStore() |
| 175 | require.NoError(t, s.CreateCluster(ctx, ns, cluster)) |
| 176 | |
| 177 | clusterProbe := NewClusterChecker(s, ns, clusterName) |
| 178 | clusterProbe.WithPingInterval(100 * time.Millisecond) |
| 179 | clusterProbe.Start() |
| 180 | defer clusterProbe.Close() |
| 181 | |
| 182 | ticker := time.NewTicker(400 * time.Millisecond) |
| 183 | defer ticker.Stop() |
| 184 | <-ticker.C |
| 185 | |
| 186 | for _, node := range nodes { |
| 187 | info, err := node.GetClusterInfo(ctx) |
| 188 | require.NoError(t, err) |
| 189 | require.EqualValues(t, 1, info.CurrentEpoch) |
| 190 | } |
| 191 | require.NoError(t, s.UpdateCluster(ctx, ns, cluster)) |
| 192 | |
| 193 | <-ticker.C |
| 194 | // should sync the clusterName info |
| 195 | for _, node := range nodes { |
| 196 | info, err := node.GetClusterInfo(ctx) |
| 197 | require.NoError(t, err) |
| 198 | require.EqualValues(t, 2, info.CurrentEpoch) |
| 199 | } |
| 200 | require.NoError(t, s.UpdateCluster(ctx, ns, cluster)) |
| 201 | clusterProbe.sendSyncEvent() |
| 202 | <-ticker.C |
| 203 | for _, node := range nodes { |
| 204 | info, err := node.GetClusterInfo(ctx) |
| 205 | require.NoError(t, err) |
| 206 | require.EqualValues(t, 3, info.CurrentEpoch) |
| 207 | } |
| 208 | } |
| 209 | |
| 210 | func TestCluster_MigrateSlot(t *testing.T) { |
| 211 | ctx := context.Background() |
nothing calls this directly
no test coverage detected