From 6cf022d3425d50174c22e08f5bb024b3744c0ed1 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Mon, 21 Sep 2026 15:36:14 +0000 Subject: [PATCH] Add move-tables start and stop commands Wire `pscale branch vtctld move-tables start` and `stop` to the existing vtctld workflow start/stop endpoints, using the same --workflow and --target-keyspace flags as the other move-tables commands. Co-authored-by: Gomez --- AGENTS.md | 8 +- internal/cmd/branch/vtctld/move_tables.go | 112 +++++++++++- .../cmd/branch/vtctld/move_tables_test.go | 167 ++++++++++++++++++ internal/cmd/branch/vtctld/progress_test.go | 50 ++++++ .../cmd/branch/vtctld/workflow_next_steps.go | 26 +++ .../branch/vtctld/workflow_next_steps_test.go | 6 + 6 files changed, 365 insertions(+), 4 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index 4fc08236..bbe850f0 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -479,14 +479,18 @@ External create required flags: `--host`, `--source-database`, `--username`, `-- ## Vitess MoveTables -Copy tables between keyspaces with `pscale branch vtctld move-tables`. `pscale workflow` will be deprecated soon; prefer `move-tables` for new work. JSON output includes `next_steps` — follow those commands. Typical order: create the target keyspace (`keyspace create` or `keyspace create-external`), create the workflow, poll `status`, switch replica traffic, then primary traffic (ask the user first), then `complete --dry-run` and `complete` after approval. +Copy tables between keyspaces with `pscale branch vtctld move-tables`. `pscale workflow` will be deprecated soon; prefer `move-tables` for new work. JSON output includes `next_steps` — follow those commands. Typical order: create the target keyspace (`keyspace create` or `keyspace create-external`), create the workflow, `start` if you used `--auto-start=false` or after `stop`, poll `status`, switch replica traffic, then primary traffic (ask the user first), then `complete --dry-run` and `complete` after approval. -`--workflow` is the workflow name you choose. `--source-keyspace` and `--target-keyspace` are required on create. Pass `--tables t1,t2` or `--all-tables` (mutually exclusive). +`--workflow` is the workflow name you choose. `--source-keyspace` and `--target-keyspace` are required on create. Pass `--tables t1,t2` or `--all-tables` (mutually exclusive). `start` and `stop` take `--workflow` and `--target-keyspace`. ```bash pscale branch vtctld move-tables list --org --format json pscale branch vtctld move-tables create --org --format json \ --workflow --source-keyspace --target-keyspace --tables +pscale branch vtctld move-tables start --org --format json \ + --workflow --target-keyspace +pscale branch vtctld move-tables stop --org --format json \ + --workflow --target-keyspace pscale branch vtctld move-tables status --org --format json \ --workflow --target-keyspace pscale branch vtctld move-tables switch-traffic --org --format json \ diff --git a/internal/cmd/branch/vtctld/move_tables.go b/internal/cmd/branch/vtctld/move_tables.go index e05c2d49..86ff8bc3 100644 --- a/internal/cmd/branch/vtctld/move_tables.go +++ b/internal/cmd/branch/vtctld/move_tables.go @@ -29,6 +29,8 @@ func MoveTablesCmd(ch *cmdutil.Helper) *cobra.Command { cmd.AddCommand(MoveTablesCreateCmd(ch)) cmd.AddCommand(MoveTablesShowCmd(ch)) cmd.AddCommand(MoveTablesStatusCmd(ch)) + cmd.AddCommand(MoveTablesStartCmd(ch)) + cmd.AddCommand(MoveTablesStopCmd(ch)) cmd.AddCommand(MoveTablesSwitchTrafficCmd(ch)) cmd.AddCommand(MoveTablesReverseTrafficCmd(ch)) cmd.AddCommand(MoveTablesCancelCmd(ch)) @@ -121,9 +123,15 @@ func MoveTablesCreateCmd(ch *cmdutil.Helper) *cobra.Command { } end() - return printWorkflowJSON(ch.Printer, data, []workflowNextStep{ + nextSteps := []workflowNextStep{ moveTablesStatusStep(ch.Config.Organization, database, branch, flags.workflow, flags.targetKeyspace, "Monitor copy and replication progress"), - }) + } + if cmd.Flags().Changed("auto-start") && !flags.autoStart { + nextSteps = []workflowNextStep{ + moveTablesStartStep(ch.Config.Organization, database, branch, flags.workflow, flags.targetKeyspace, "Start the workflow after creating it with --auto-start=false"), + } + } + return printWorkflowJSON(ch.Printer, data, nextSteps) }, } @@ -300,6 +308,106 @@ func MoveTablesStatusCmd(ch *cmdutil.Helper) *cobra.Command { return cmd } +func MoveTablesStartCmd(ch *cmdutil.Helper) *cobra.Command { + var flags struct { + workflow string + targetKeyspace string + } + + cmd := &cobra.Command{ + Use: "start ", + Short: "Start a MoveTables workflow", + Args: cmdutil.RequiredArgs("database", "branch"), + RunE: func(cmd *cobra.Command, args []string) error { + ctx := cmd.Context() + database, branch := args[0], args[1] + + client, err := ch.Client() + if err != nil { + return err + } + + end := ch.Printer.PrintProgress( + fmt.Sprintf("Starting MoveTables workflow %s on %s\u2026", + printer.BoldBlue(flags.workflow), progressTarget(ch.Config.Organization, database, branch))) + defer end() + + data, err := client.Vtctld.StartWorkflow(ctx, &ps.VtctldStartWorkflowRequest{ + Organization: ch.Config.Organization, + Database: database, + Branch: branch, + Workflow: flags.workflow, + Keyspace: flags.targetKeyspace, + }) + if err != nil { + return cmdutil.HandleError(err) + } + + end() + return printWorkflowJSON(ch.Printer, data, []workflowNextStep{ + moveTablesStatusStep(ch.Config.Organization, database, branch, flags.workflow, flags.targetKeyspace, "Monitor copy and replication progress"), + }) + }, + } + + cmd.Flags().StringVar(&flags.workflow, "workflow", "", "Name of the workflow") + cmd.Flags().StringVar(&flags.targetKeyspace, "target-keyspace", "", "Target keyspace") + cmd.MarkFlagRequired("workflow") // nolint:errcheck + cmd.MarkFlagRequired("target-keyspace") // nolint:errcheck + + return cmd +} + +func MoveTablesStopCmd(ch *cmdutil.Helper) *cobra.Command { + var flags struct { + workflow string + targetKeyspace string + } + + cmd := &cobra.Command{ + Use: "stop ", + Short: "Stop a MoveTables workflow", + Args: cmdutil.RequiredArgs("database", "branch"), + RunE: func(cmd *cobra.Command, args []string) error { + ctx := cmd.Context() + database, branch := args[0], args[1] + + client, err := ch.Client() + if err != nil { + return err + } + + end := ch.Printer.PrintProgress( + fmt.Sprintf("Stopping MoveTables workflow %s on %s\u2026", + printer.BoldBlue(flags.workflow), progressTarget(ch.Config.Organization, database, branch))) + defer end() + + data, err := client.Vtctld.StopWorkflow(ctx, &ps.VtctldStopWorkflowRequest{ + Organization: ch.Config.Organization, + Database: database, + Branch: branch, + Workflow: flags.workflow, + Keyspace: flags.targetKeyspace, + }) + if err != nil { + return cmdutil.HandleError(err) + } + + end() + return printWorkflowJSON(ch.Printer, data, []workflowNextStep{ + moveTablesStartStep(ch.Config.Organization, database, branch, flags.workflow, flags.targetKeyspace, "Resume the workflow when you are ready to continue"), + }) + }, + } + + cmd.Flags().StringVar(&flags.workflow, "workflow", "", "Name of the workflow") + cmd.Flags().StringVar(&flags.targetKeyspace, "target-keyspace", "", "Target keyspace") + cmd.MarkFlagRequired("workflow") // nolint:errcheck + cmd.MarkFlagRequired("target-keyspace") // nolint:errcheck + + return cmd +} + func MoveTablesSwitchTrafficCmd(ch *cmdutil.Helper) *cobra.Command { var flags struct { workflow string diff --git a/internal/cmd/branch/vtctld/move_tables_test.go b/internal/cmd/branch/vtctld/move_tables_test.go index b1b75663..935eebe2 100644 --- a/internal/cmd/branch/vtctld/move_tables_test.go +++ b/internal/cmd/branch/vtctld/move_tables_test.go @@ -241,6 +241,55 @@ func TestMoveTablesCreateWithAllFlags(t *testing.T) { }) } +func TestMoveTablesCreateWithAutoStartFalse(t *testing.T) { + c := qt.New(t) + setMoveTablesPollInterval(t, 0) + + org := "my-org" + db := "my-db" + branch := "my-branch" + + svc := &mock.MoveTablesService{ + CreateFn: func(ctx context.Context, req *ps.MoveTablesCreateRequest) (*ps.VtctldOperationReference, error) { + c.Assert(req.AutoStart, qt.IsNotNil) + c.Assert(*req.AutoStart, qt.IsFalse) + return &ps.VtctldOperationReference{ID: "create-op"}, nil + }, + } + + vtctldSvc := &mock.VtctldService{ + GetOperationFn: func(ctx context.Context, req *ps.GetVtctldOperationRequest) (*ps.VtctldOperation, error) { + return &ps.VtctldOperation{ + ID: "create-op", + State: "completed", + Completed: true, + Result: json.RawMessage(`{"summary":"created"}`), + }, nil + }, + } + + var buf bytes.Buffer + ch := moveTablesTestHelper(org, svc, vtctldSvc, &buf) + + cmd := MoveTablesCmd(ch) + cmd.SetArgs([]string{"create", db, branch, + "--workflow", "my-workflow", + "--target-keyspace", "target-ks", + "--source-keyspace", "source-ks", + "--auto-start=false", + }) + err := cmd.Execute() + c.Assert(err, qt.IsNil) + c.Assert(svc.CreateFnInvoked, qt.IsTrue) + c.Assert(buf.String(), qt.JSONEquals, map[string]any{ + "summary": "created", + "next_steps": []any{map[string]any{ + "command": "pscale branch vtctld move-tables start my-db my-branch --org my-org --workflow my-workflow --target-keyspace target-ks --format json", + "reason": "Start the workflow after creating it with --auto-start=false", + }}, + }) +} + func TestMoveTablesSwitchTrafficWithMaxLag(t *testing.T) { c := qt.New(t) setMoveTablesPollInterval(t, 0) @@ -764,3 +813,121 @@ func TestMoveTablesStatusAddsNextSteps(t *testing.T) { }}, }) } + +func TestMoveTablesStart(t *testing.T) { + c := qt.New(t) + + org := "my-org" + db := "my-db" + branch := "my-branch" + + vtctldSvc := &mock.VtctldService{ + StartWorkflowFn: func(ctx context.Context, req *ps.VtctldStartWorkflowRequest) (json.RawMessage, error) { + c.Assert(req.Organization, qt.Equals, org) + c.Assert(req.Database, qt.Equals, db) + c.Assert(req.Branch, qt.Equals, branch) + c.Assert(req.Workflow, qt.Equals, "my-workflow") + c.Assert(req.Keyspace, qt.Equals, "target-ks") + return json.RawMessage(`{"summary":"started"}`), nil + }, + } + + var buf bytes.Buffer + ch := moveTablesTestHelper(org, nil, vtctldSvc, &buf) + + cmd := MoveTablesCmd(ch) + cmd.SetArgs([]string{"start", db, branch, + "--workflow", "my-workflow", + "--target-keyspace", "target-ks", + }) + err := cmd.Execute() + c.Assert(err, qt.IsNil) + c.Assert(vtctldSvc.StartWorkflowFnInvoked, qt.IsTrue) + c.Assert(buf.String(), qt.JSONEquals, map[string]any{ + "summary": "started", + "next_steps": []any{map[string]any{ + "command": "pscale branch vtctld move-tables status my-db my-branch --org my-org --workflow my-workflow --target-keyspace target-ks --format json", + "reason": "Monitor copy and replication progress", + }}, + }) +} + +func TestMoveTablesStartRequiresFlags(t *testing.T) { + c := qt.New(t) + + org := "my-org" + db := "my-db" + branch := "my-branch" + + vtctldSvc := &mock.VtctldService{} + var buf bytes.Buffer + ch := moveTablesTestHelper(org, nil, vtctldSvc, &buf) + + cmd := MoveTablesCmd(ch) + cmd.SetArgs([]string{"start", db, branch}) + err := cmd.Execute() + + c.Assert(err, qt.IsNotNil) + c.Assert(err.Error(), qt.Contains, "required flag") + c.Assert(vtctldSvc.StartWorkflowFnInvoked, qt.IsFalse) + c.Assert(buf.String(), qt.Equals, "") +} + +func TestMoveTablesStop(t *testing.T) { + c := qt.New(t) + + org := "my-org" + db := "my-db" + branch := "my-branch" + + vtctldSvc := &mock.VtctldService{ + StopWorkflowFn: func(ctx context.Context, req *ps.VtctldStopWorkflowRequest) (json.RawMessage, error) { + c.Assert(req.Organization, qt.Equals, org) + c.Assert(req.Database, qt.Equals, db) + c.Assert(req.Branch, qt.Equals, branch) + c.Assert(req.Workflow, qt.Equals, "my-workflow") + c.Assert(req.Keyspace, qt.Equals, "target-ks") + return json.RawMessage(`{"summary":"stopped"}`), nil + }, + } + + var buf bytes.Buffer + ch := moveTablesTestHelper(org, nil, vtctldSvc, &buf) + + cmd := MoveTablesCmd(ch) + cmd.SetArgs([]string{"stop", db, branch, + "--workflow", "my-workflow", + "--target-keyspace", "target-ks", + }) + err := cmd.Execute() + c.Assert(err, qt.IsNil) + c.Assert(vtctldSvc.StopWorkflowFnInvoked, qt.IsTrue) + c.Assert(buf.String(), qt.JSONEquals, map[string]any{ + "summary": "stopped", + "next_steps": []any{map[string]any{ + "command": "pscale branch vtctld move-tables start my-db my-branch --org my-org --workflow my-workflow --target-keyspace target-ks --format json", + "reason": "Resume the workflow when you are ready to continue", + }}, + }) +} + +func TestMoveTablesStopRequiresFlags(t *testing.T) { + c := qt.New(t) + + org := "my-org" + db := "my-db" + branch := "my-branch" + + vtctldSvc := &mock.VtctldService{} + var buf bytes.Buffer + ch := moveTablesTestHelper(org, nil, vtctldSvc, &buf) + + cmd := MoveTablesCmd(ch) + cmd.SetArgs([]string{"stop", db, branch, "--workflow", "my-workflow"}) + err := cmd.Execute() + + c.Assert(err, qt.IsNotNil) + c.Assert(err.Error(), qt.Contains, "required flag") + c.Assert(vtctldSvc.StopWorkflowFnInvoked, qt.IsFalse) + c.Assert(buf.String(), qt.Equals, "") +} diff --git a/internal/cmd/branch/vtctld/progress_test.go b/internal/cmd/branch/vtctld/progress_test.go index c9074818..66768283 100644 --- a/internal/cmd/branch/vtctld/progress_test.go +++ b/internal/cmd/branch/vtctld/progress_test.go @@ -117,6 +117,56 @@ func TestMoveTablesShowProgressIncludesOrganization(t *testing.T) { c.Assert(progress.String(), qt.Contains, "Fetching MoveTables workflow my-workflow on my-org/my-db/my-branch…") } +func TestMoveTablesStartProgressIncludesOrganization(t *testing.T) { + c := qt.New(t) + setNonTTYProgress(t) + + org := "my-org" + db := "my-db" + branch := "my-branch" + + svc := &mock.VtctldService{ + StartWorkflowFn: func(ctx context.Context, req *ps.VtctldStartWorkflowRequest) (json.RawMessage, error) { + return json.RawMessage(`{"summary":"started"}`), nil + }, + } + + var progress bytes.Buffer + ch := newHumanProgressHelper(org, &progress, &ps.Client{Vtctld: svc}) + + cmd := MoveTablesCmd(ch) + cmd.SetArgs([]string{"start", db, branch, "--workflow", "my-workflow", "--target-keyspace", "target-ks"}) + + err := cmd.Execute() + c.Assert(err, qt.IsNil) + c.Assert(progress.String(), qt.Contains, "Starting MoveTables workflow my-workflow on my-org/my-db/my-branch…") +} + +func TestMoveTablesStopProgressIncludesOrganization(t *testing.T) { + c := qt.New(t) + setNonTTYProgress(t) + + org := "my-org" + db := "my-db" + branch := "my-branch" + + svc := &mock.VtctldService{ + StopWorkflowFn: func(ctx context.Context, req *ps.VtctldStopWorkflowRequest) (json.RawMessage, error) { + return json.RawMessage(`{"summary":"stopped"}`), nil + }, + } + + var progress bytes.Buffer + ch := newHumanProgressHelper(org, &progress, &ps.Client{Vtctld: svc}) + + cmd := MoveTablesCmd(ch) + cmd.SetArgs([]string{"stop", db, branch, "--workflow", "my-workflow", "--target-keyspace", "target-ks"}) + + err := cmd.Execute() + c.Assert(err, qt.IsNil) + c.Assert(progress.String(), qt.Contains, "Stopping MoveTables workflow my-workflow on my-org/my-db/my-branch…") +} + func TestVDiffListProgressIncludesOrganization(t *testing.T) { c := qt.New(t) setNonTTYProgress(t) diff --git a/internal/cmd/branch/vtctld/workflow_next_steps.go b/internal/cmd/branch/vtctld/workflow_next_steps.go index bad3c7f4..54ef5b66 100644 --- a/internal/cmd/branch/vtctld/workflow_next_steps.go +++ b/internal/cmd/branch/vtctld/workflow_next_steps.go @@ -182,6 +182,13 @@ func moveTablesStatusStep(org, database, branch, workflow, targetKeyspace, reaso } } +func moveTablesStartStep(org, database, branch, workflow, targetKeyspace, reason string) workflowNextStep { + return workflowNextStep{ + Command: moveTablesCommand(org, "start", database, branch, workflow, targetKeyspace), + Reason: reason, + } +} + func moveTablesVDiffCreateStep(org, database, branch, workflow, targetKeyspace string) workflowNextStep { return workflowNextStep{ Command: fmt.Sprintf( @@ -222,6 +229,12 @@ func moveTablesStatusNextSteps(data json.RawMessage, org, database, branch, work return nil } + if moveTablesStreamsAllStopped(status) { + return []workflowNextStep{ + moveTablesStartStep(org, database, branch, workflow, targetKeyspace, "Resume the stopped workflow"), + } + } + hasStreams, streamsNeedMonitoring := moveTablesStreamState(status) if len(status.TableCopyState) > 0 || streamsNeedMonitoring { return []workflowNextStep{ @@ -287,6 +300,19 @@ func moveTablesStreamState(status moveTablesStatus) (bool, bool) { return hasStreams, false } +func moveTablesStreamsAllStopped(status moveTablesStatus) bool { + count := 0 + for _, shard := range status.ShardStreams { + for _, stream := range shard.Streams { + count++ + if !strings.EqualFold(stream.Status, "Stopped") { + return false + } + } + } + return count > 0 +} + func vdiffCreateNextSteps(data json.RawMessage, org, database, branch, workflow, targetKeyspace string) []workflowNextStep { var result vdiffResult if err := json.Unmarshal(data, &result); err != nil || result.UUID == "" { diff --git a/internal/cmd/branch/vtctld/workflow_next_steps_test.go b/internal/cmd/branch/vtctld/workflow_next_steps_test.go index 1a234766..9d71581c 100644 --- a/internal/cmd/branch/vtctld/workflow_next_steps_test.go +++ b/internal/cmd/branch/vtctld/workflow_next_steps_test.go @@ -51,6 +51,12 @@ func TestMoveTablesStatusNextSteps(t *testing.T) { wantCommand: "pscale branch vtctld move-tables status my-db my-branch --org my-org --workflow my-workflow --target-keyspace target-ks --format json", wantSteps: 1, }, + { + name: "stopped streams", + data: `{"table_copy_state":{},"shard_streams":{"target/-":{"streams":[{"status":"Stopped"}]}},"traffic_state":"Reads Not Switched. Writes Not Switched"}`, + wantCommand: "pscale branch vtctld move-tables start my-db my-branch --org my-org --workflow my-workflow --target-keyspace target-ks --format json", + wantSteps: 1, + }, { name: "unrecognized traffic state", data: `{"traffic_state":"Something Vitess Has Not Told Us About"}`,