MCPcopy Create free account
hub / github.com/apache/kvrocks-controller / TestClusterMigrateData

Function TestClusterMigrateData

server/api/cluster_test.go:215–305  ·  view source on GitHub ↗
(t *testing.T)

Source from the content-addressed store, hash-verified

213}
214
215func TestClusterMigrateData(t *testing.T) {
216 ns := "test-ns"
217 clusterName := "test-cluster"
218 clusterStore := store.NewClusterStore(engine.NewMock())
219 handler := &ClusterHandler{s: clusterStore}
220
221 ctx := context.Background()
222 nodeAddrs := []string{"127.0.0.1:7770", "127.0.0.1:7771"}
223 sourceRedisClient := redis.NewClient(&redis.Options{Addr: nodeAddrs[0]})
224 targetRedisClient := redis.NewClient(&redis.Options{Addr: nodeAddrs[1]})
225
226 cluster, err := store.NewCluster(clusterName, nodeAddrs, 1)
227 require.NoError(t, err)
228 require.NoError(t, cluster.Reset(ctx))
229 defer func() {
230 require.NoError(t, cluster.Reset(ctx))
231 }()
232 require.NoError(t, cluster.SyncToNodes(ctx))
233 clusterStore.CreateCluster(ctx, ns, cluster)
234
235 sendRequest := func(t *testing.T, ns, cluster string, slotRange store.SlotRange) {
236 recorder := httptest.NewRecorder()
237 reqCtx := GetTestContext(recorder)
238 reqCtx.Set(consts.ContextKeyStore, handler.s)
239 reqCtx.Params = []gin.Param{{Key: "namespace", Value: ns}, {Key: "cluster", Value: cluster}}
240 body, err := json.Marshal(&MigrateSlotRequest{Target: 1, Slot: slotRange})
241 require.NoError(t, err)
242 reqCtx.Request.Body = io.NopCloser(bytes.NewBuffer(body))
243
244 middleware.RequiredCluster(reqCtx)
245 handler.MigrateSlot(reqCtx)
246 require.Equal(t, http.StatusOK, recorder.Code)
247 }
248
249 runController := func(t *testing.T) *controller.Controller {
250 ctrl, err := controller.New(clusterStore, &config.ControllerConfig{
251 FailOver: &config.FailOverConfig{
252 PingIntervalSeconds: 1,
253 MaxPingCount: 3,
254 },
255 })
256 require.NoError(t, err)
257 require.NoError(t, ctrl.Start(ctx))
258 ctrl.WaitForReady()
259 return ctrl
260 }
261
262 t.Run("migrate slot(s) from one shard to another", func(t *testing.T) {
263 for _, slotRange := range []store.SlotRange{
264 {Start: 10, Stop: 10},
265 {Start: 11, Stop: 14},
266 } {
267 for i := slotRange.Start; i <= slotRange.Stop; i++ {
268 require.NoError(t, sourceRedisClient.Set(ctx, util.SlotTable[i], "test-value", 0).Err())
269 }
270 sendRequest(t, ns, "test-cluster", slotRange)
271 gotCluster, err := clusterStore.GetCluster(ctx, ns, clusterName)
272 require.NoError(t, err)

Callers

nothing calls this directly

Calls 15

ResetMethod · 0.95
SyncToNodesMethod · 0.95
CreateClusterMethod · 0.95
MigrateSlotMethod · 0.95
GetClusterMethod · 0.95
NewClusterStoreFunction · 0.92
NewMockFunction · 0.92
NewClusterFunction · 0.92
RequiredClusterFunction · 0.92
NewFunction · 0.92
RemoveSlotFromSlotRangesFunction · 0.92
AddSlotToSlotRangesFunction · 0.92

Tested by

no test coverage detected