mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-08-16 12:16:36 +00:00
338 lines
14 KiB
Go
338 lines
14 KiB
Go
package shell
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"flag"
|
|
"fmt"
|
|
"io"
|
|
"log"
|
|
"time"
|
|
|
|
"github.com/seaweedfs/seaweedfs/weed/pb"
|
|
"github.com/seaweedfs/seaweedfs/weed/wdclient"
|
|
|
|
"github.com/seaweedfs/seaweedfs/weed/operation"
|
|
"github.com/seaweedfs/seaweedfs/weed/pb/volume_server_pb"
|
|
"github.com/seaweedfs/seaweedfs/weed/storage/needle"
|
|
"github.com/seaweedfs/seaweedfs/weed/util"
|
|
|
|
"google.golang.org/grpc"
|
|
)
|
|
|
|
func init() {
|
|
Commands = append(Commands, &commandVolumeMove{})
|
|
}
|
|
|
|
type commandVolumeMove struct {
|
|
}
|
|
|
|
func (c *commandVolumeMove) Name() string {
|
|
return "volume.move"
|
|
}
|
|
|
|
func (c *commandVolumeMove) Help() string {
|
|
return `move a live volume from one volume server to another volume server
|
|
|
|
volume.move -source <source volume server host:port> -target <target volume server host:port> -volumeId <volume id>
|
|
volume.move -source <source volume server host:port> -target <target volume server host:port> -volumeId <volume id> -disk [hdd|ssd|<tag>]
|
|
|
|
This command move a live volume from one volume server to another volume server. Here are the steps:
|
|
|
|
1. This command marks the source volume as read-only, copies it to the target volume server, and records the last entry timestamp.
|
|
2. This command asks the target volume server to mount the new volume.
|
|
3. This command asks the target volume server to tail the source volume for updates after the timestamp, for 1 minutes to drain any in-flight requests.
|
|
4. This command asks the source volume server to delete the source volume.
|
|
|
|
The option "-disk [hdd|ssd|<tag>]" can be used to change the volume disk type.
|
|
The option "-timeout" fails the whole move if it does not finish in time.
|
|
|
|
`
|
|
}
|
|
|
|
func (c *commandVolumeMove) HasTag(CommandTag) bool {
|
|
return false
|
|
}
|
|
|
|
func (c *commandVolumeMove) Do(args []string, commandEnv *CommandEnv, writer io.Writer) (err error) {
|
|
|
|
volMoveCommand := flag.NewFlagSet(c.Name(), flag.ContinueOnError)
|
|
volumeIdInt := volMoveCommand.Int("volumeId", 0, "the volume id")
|
|
sourceNodeStr := volMoveCommand.String("source", "", "the source volume server <host>:<port>")
|
|
targetNodeStr := volMoveCommand.String("target", "", "the target volume server <host>:<port>")
|
|
diskTypeStr := volMoveCommand.String("disk", "", "[hdd|ssd|<tag>] hard drive or solid state drive or any tag")
|
|
ioBytePerSecond := volMoveCommand.Int64("ioBytePerSecond", 0, "limit the speed of move")
|
|
timeout := volMoveCommand.Duration("timeout", 0, "wall-clock cap on the whole move; 0 = no timeout")
|
|
noLock := volMoveCommand.Bool("noLock", false, "do not lock the admin shell at one's own risk")
|
|
|
|
if err = volMoveCommand.Parse(args); err != nil {
|
|
return nil
|
|
}
|
|
|
|
if *noLock {
|
|
commandEnv.noLock = true
|
|
} else {
|
|
if err = commandEnv.confirmIsLocked(args); err != nil {
|
|
return
|
|
}
|
|
}
|
|
|
|
sourceVolumeServer, targetVolumeServer := pb.ServerAddress(*sourceNodeStr), pb.ServerAddress(*targetNodeStr)
|
|
|
|
volumeId := needle.VolumeId(*volumeIdInt)
|
|
|
|
if sourceVolumeServer == targetVolumeServer {
|
|
return fmt.Errorf("source and target volume servers are the same!")
|
|
}
|
|
|
|
ctx := context.Background()
|
|
if *timeout > 0 {
|
|
var cancel context.CancelFunc
|
|
ctx, cancel = context.WithTimeout(ctx, *timeout)
|
|
defer cancel()
|
|
}
|
|
|
|
return LiveMoveVolume(ctx, commandEnv.option.GrpcDialOption, writer, volumeId, sourceVolumeServer, targetVolumeServer, 5*time.Second, *diskTypeStr, *ioBytePerSecond, false)
|
|
}
|
|
|
|
// LiveMoveVolume moves one volume from one source volume server to one target volume server, with idleTimeout to drain the incoming requests.
|
|
func LiveMoveVolume(ctx context.Context, grpcDialOption grpc.DialOption, writer io.Writer, volumeId needle.VolumeId, sourceVolumeServer, targetVolumeServer pb.ServerAddress, idleTimeout time.Duration, diskType string, ioBytePerSecond int64, skipTailError bool) (err error) {
|
|
|
|
log.Printf("copying volume %d from %s to %s", volumeId, sourceVolumeServer, targetVolumeServer)
|
|
lastAppendAtNs, leftReadonly, err := copyVolume(ctx, grpcDialOption, writer, volumeId, sourceVolumeServer, targetVolumeServer, diskType, ioBytePerSecond, false)
|
|
if err != nil {
|
|
return fmt.Errorf("copy volume %d from %s to %s: %v", volumeId, sourceVolumeServer, targetVolumeServer, err)
|
|
}
|
|
|
|
// A move aborted after the copy must restore the source writability, or
|
|
// the source volume is left permanently readonly.
|
|
var sourceDeleteStarted bool
|
|
defer func() {
|
|
if err == nil || !leftReadonly {
|
|
return
|
|
}
|
|
if !sourceDeleteStarted {
|
|
// The target copy may be missing tailed entries; remove it so the
|
|
// restored source stays the only replica.
|
|
deleteCtx, deleteCancel := context.WithTimeout(context.WithoutCancel(ctx), 30*time.Second)
|
|
defer deleteCancel()
|
|
if dErr := deleteVolume(deleteCtx, grpcDialOption, volumeId, targetVolumeServer, false, true); dErr != nil {
|
|
// Restoring the source while the stale target stays mounted
|
|
// risks divergent replicas; re-running the move recovers both.
|
|
log.Printf("volume %d is left readonly on %s: failed to delete the incomplete copy on %s: %v; re-run volume.move to recover", volumeId, sourceVolumeServer, targetVolumeServer, dErr)
|
|
return
|
|
}
|
|
}
|
|
restoreCtx, restoreCancel := context.WithTimeout(context.WithoutCancel(ctx), 30*time.Second)
|
|
defer restoreCancel()
|
|
if wErr := markVolumeWritable(restoreCtx, grpcDialOption, volumeId, sourceVolumeServer, true, false); wErr != nil {
|
|
log.Printf("failed to restore volume %d writable on %s: %v", volumeId, sourceVolumeServer, wErr)
|
|
}
|
|
}()
|
|
|
|
log.Printf("tailing volume %d from %s to %s", volumeId, sourceVolumeServer, targetVolumeServer)
|
|
if err = tailVolume(ctx, grpcDialOption, volumeId, sourceVolumeServer, targetVolumeServer, lastAppendAtNs, idleTimeout); err != nil {
|
|
if skipTailError {
|
|
fmt.Fprintf(writer, "tail volume %d from %s to %s: %v\n", volumeId, sourceVolumeServer, targetVolumeServer, err)
|
|
} else {
|
|
return fmt.Errorf("tail volume %d from %s to %s: %v", volumeId, sourceVolumeServer, targetVolumeServer, err)
|
|
}
|
|
}
|
|
|
|
log.Printf("deleting volume %d from %s", volumeId, sourceVolumeServer)
|
|
sourceDeleteStarted = true
|
|
if err = deleteVolume(ctx, grpcDialOption, volumeId, sourceVolumeServer, false, true); err != nil {
|
|
return fmt.Errorf("delete volume %d from %s: %v", volumeId, sourceVolumeServer, err)
|
|
}
|
|
|
|
log.Printf("moved volume %d from %s to %s", volumeId, sourceVolumeServer, targetVolumeServer)
|
|
return nil
|
|
}
|
|
|
|
func copyVolume(ctx context.Context, grpcDialOption grpc.DialOption, writer io.Writer, volumeId needle.VolumeId, sourceVolumeServer, targetVolumeServer pb.ServerAddress, diskType string, ioBytePerSecond int64, restoreWritable bool) (lastAppendAtNs uint64, leftReadonly bool, err error) {
|
|
|
|
// check to see if the volume is already read-only and if its not then we need
|
|
// to mark it as read-only and then before we return we need to undo what we
|
|
// did
|
|
var shouldMarkWritable bool
|
|
defer func() {
|
|
if !shouldMarkWritable {
|
|
return
|
|
}
|
|
if !restoreWritable && err == nil {
|
|
leftReadonly = true
|
|
return
|
|
}
|
|
|
|
// Restoring writability must outlive a cancelled copy, or the source
|
|
// is left permanently readonly on abort, as balance_task.go already
|
|
// guards against.
|
|
restoreCtx, restoreCancel := context.WithTimeout(context.WithoutCancel(ctx), 30*time.Second)
|
|
defer restoreCancel()
|
|
clientErr := operation.WithVolumeServerClient(false, sourceVolumeServer, grpcDialOption, func(volumeServerClient volume_server_pb.VolumeServerClient) error {
|
|
_, writableErr := volumeServerClient.VolumeMarkWritable(restoreCtx, &volume_server_pb.VolumeMarkWritableRequest{
|
|
VolumeId: uint32(volumeId),
|
|
})
|
|
return writableErr
|
|
})
|
|
if clientErr != nil {
|
|
log.Printf("failed to mark volume %d as writable after copy from %s: %v", volumeId, sourceVolumeServer, clientErr)
|
|
}
|
|
}()
|
|
|
|
err = operation.WithVolumeServerClient(false, sourceVolumeServer, grpcDialOption, func(volumeServerClient volume_server_pb.VolumeServerClient) error {
|
|
resp, statusErr := volumeServerClient.VolumeStatus(ctx, &volume_server_pb.VolumeStatusRequest{
|
|
VolumeId: uint32(volumeId),
|
|
})
|
|
if statusErr == nil && !resp.IsReadOnly {
|
|
shouldMarkWritable = true
|
|
_, readonlyErr := volumeServerClient.VolumeMarkReadonly(ctx, &volume_server_pb.VolumeMarkReadonlyRequest{
|
|
VolumeId: uint32(volumeId),
|
|
Persist: false,
|
|
})
|
|
return readonlyErr
|
|
}
|
|
return statusErr
|
|
})
|
|
if err != nil {
|
|
return
|
|
}
|
|
|
|
err = operation.WithVolumeServerClient(true, targetVolumeServer, grpcDialOption, func(volumeServerClient volume_server_pb.VolumeServerClient) error {
|
|
stream, replicateErr := volumeServerClient.VolumeCopy(ctx, &volume_server_pb.VolumeCopyRequest{
|
|
VolumeId: uint32(volumeId),
|
|
SourceDataNode: string(sourceVolumeServer),
|
|
DiskType: diskType,
|
|
IoBytePerSecond: ioBytePerSecond,
|
|
})
|
|
if replicateErr != nil {
|
|
return replicateErr
|
|
}
|
|
for {
|
|
resp, recvErr := stream.Recv()
|
|
if recvErr != nil {
|
|
if recvErr == io.EOF {
|
|
break
|
|
} else {
|
|
return recvErr
|
|
}
|
|
}
|
|
if resp.LastAppendAtNs != 0 {
|
|
lastAppendAtNs = resp.LastAppendAtNs
|
|
} else {
|
|
fmt.Fprintf(writer, "%s => %s volume %d processed %s\n", sourceVolumeServer, targetVolumeServer, volumeId, util.BytesToHumanReadable(uint64(resp.ProcessedBytes)))
|
|
}
|
|
}
|
|
|
|
return nil
|
|
})
|
|
|
|
return
|
|
}
|
|
|
|
func tailVolume(ctx context.Context, grpcDialOption grpc.DialOption, volumeId needle.VolumeId, sourceVolumeServer, targetVolumeServer pb.ServerAddress, lastAppendAtNs uint64, idleTimeout time.Duration) (err error) {
|
|
|
|
return operation.WithVolumeServerClient(true, targetVolumeServer, grpcDialOption, func(volumeServerClient volume_server_pb.VolumeServerClient) error {
|
|
_, replicateErr := volumeServerClient.VolumeTailReceiver(ctx, &volume_server_pb.VolumeTailReceiverRequest{
|
|
VolumeId: uint32(volumeId),
|
|
SinceNs: lastAppendAtNs,
|
|
IdleTimeoutSeconds: uint32(idleTimeout.Seconds()),
|
|
SourceVolumeServer: string(sourceVolumeServer),
|
|
})
|
|
return replicateErr
|
|
})
|
|
|
|
}
|
|
|
|
// deleteVolume removes the volume from sourceVolumeServer. When keepRemoteData
|
|
// is true, the cloud-tier object backing the volume is left intact — used on
|
|
// the source side of a move where another server is taking over the same .vif.
|
|
func deleteVolume(ctx context.Context, grpcDialOption grpc.DialOption, volumeId needle.VolumeId, sourceVolumeServer pb.ServerAddress, onlyEmpty bool, keepRemoteData bool) (err error) {
|
|
return operation.WithVolumeServerClient(false, sourceVolumeServer, grpcDialOption, func(volumeServerClient volume_server_pb.VolumeServerClient) error {
|
|
_, deleteErr := volumeServerClient.VolumeDelete(ctx, &volume_server_pb.VolumeDeleteRequest{
|
|
VolumeId: uint32(volumeId),
|
|
OnlyEmpty: onlyEmpty,
|
|
KeepRemoteData: keepRemoteData,
|
|
})
|
|
return deleteErr
|
|
})
|
|
}
|
|
|
|
func markVolumeWritable(ctx context.Context, grpcDialOption grpc.DialOption, volumeId needle.VolumeId, sourceVolumeServer pb.ServerAddress, writable, persist bool) (err error) {
|
|
return operation.WithVolumeServerClient(false, sourceVolumeServer, grpcDialOption, func(volumeServerClient volume_server_pb.VolumeServerClient) error {
|
|
if writable {
|
|
_, err = volumeServerClient.VolumeMarkWritable(ctx, &volume_server_pb.VolumeMarkWritableRequest{
|
|
VolumeId: uint32(volumeId),
|
|
})
|
|
} else {
|
|
_, err = volumeServerClient.VolumeMarkReadonly(ctx, &volume_server_pb.VolumeMarkReadonlyRequest{
|
|
VolumeId: uint32(volumeId),
|
|
Persist: persist,
|
|
})
|
|
}
|
|
return err
|
|
})
|
|
}
|
|
|
|
func markVolumeReplicaWritable(ctx context.Context, grpcDialOption grpc.DialOption, volumeId needle.VolumeId, location wdclient.Location, writable, persist bool) error {
|
|
if writable {
|
|
fmt.Printf("markVolumeWritable %d on %s ...\n", volumeId, location.Url)
|
|
} else {
|
|
fmt.Printf("markVolumeReadonly %d on %s persist=%v ...\n", volumeId, location.Url, persist)
|
|
}
|
|
return markVolumeWritable(ctx, grpcDialOption, volumeId, location.ServerAddress(), writable, persist)
|
|
}
|
|
|
|
func markVolumeReplicasWritable(ctx context.Context, grpcDialOption grpc.DialOption, volumeId needle.VolumeId, locations []wdclient.Location, writable, persist bool) error {
|
|
for _, location := range locations {
|
|
if err := markVolumeReplicaWritable(ctx, grpcDialOption, volumeId, location, writable, persist); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// replicateVolumeToServer copies a volume from sourceAddress to targetAddress via the VolumeCopy gRPC stream.
|
|
func replicateVolumeToServer(ctx context.Context, grpcDialOption grpc.DialOption, writer io.Writer, volumeId needle.VolumeId, sourceAddress, targetAddress pb.ServerAddress, diskType string) error {
|
|
return operation.WithVolumeServerClient(false, targetAddress, grpcDialOption, func(volumeServerClient volume_server_pb.VolumeServerClient) error {
|
|
stream, replicateErr := volumeServerClient.VolumeCopy(ctx, &volume_server_pb.VolumeCopyRequest{
|
|
VolumeId: uint32(volumeId),
|
|
SourceDataNode: string(sourceAddress),
|
|
DiskType: diskType,
|
|
})
|
|
if replicateErr != nil {
|
|
return replicateErr
|
|
}
|
|
for {
|
|
resp, recvErr := stream.Recv()
|
|
if recvErr != nil {
|
|
if recvErr == io.EOF {
|
|
break
|
|
}
|
|
return recvErr
|
|
}
|
|
if resp.ProcessedBytes > 0 {
|
|
fmt.Fprintf(writer, "volume %d processed %s bytes\n", volumeId, util.BytesToHumanReadable(uint64(resp.ProcessedBytes)))
|
|
}
|
|
}
|
|
return nil
|
|
})
|
|
}
|
|
|
|
// configureVolumeReplication sets the replication setting on a volume at the given server.
|
|
func configureVolumeReplication(ctx context.Context, grpcDialOption grpc.DialOption, volumeId needle.VolumeId, targetAddress pb.ServerAddress, replicationString string) error {
|
|
return operation.WithVolumeServerClient(false, targetAddress, grpcDialOption, func(volumeServerClient volume_server_pb.VolumeServerClient) error {
|
|
resp, configureErr := volumeServerClient.VolumeConfigure(ctx, &volume_server_pb.VolumeConfigureRequest{
|
|
VolumeId: uint32(volumeId),
|
|
Replication: replicationString,
|
|
})
|
|
if configureErr != nil {
|
|
return configureErr
|
|
}
|
|
if resp.Error != "" {
|
|
return errors.New(resp.Error)
|
|
}
|
|
return nil
|
|
})
|
|
}
|