diff options
Diffstat (limited to 'weed/server/volume_grpc_sync.go')
| -rw-r--r-- | weed/server/volume_grpc_sync.go | 46 |
1 files changed, 30 insertions, 16 deletions
diff --git a/weed/server/volume_grpc_sync.go b/weed/server/volume_grpc_sync.go index 084069211..5f56ec17d 100644 --- a/weed/server/volume_grpc_sync.go +++ b/weed/server/volume_grpc_sync.go @@ -25,13 +25,11 @@ func (vs *VolumeServer) VolumeSyncStatus(ctx context.Context, req *volume_server } -func (vs *VolumeServer) VolumeSyncIndex(ctx context.Context, req *volume_server_pb.VolumeSyncIndexRequest) (*volume_server_pb.VolumeSyncIndexResponse, error) { - - resp := &volume_server_pb.VolumeSyncIndexResponse{} +func (vs *VolumeServer) VolumeSyncIndex(req *volume_server_pb.VolumeSyncIndexRequest, stream volume_server_pb.VolumeServer_VolumeSyncIndexServer) error { v := vs.store.GetVolume(storage.VolumeId(req.VolumdId)) if v == nil { - return nil, fmt.Errorf("Not Found Volume Id %d", req.VolumdId) + return fmt.Errorf("Not Found Volume Id %d", req.VolumdId) } content, err := v.IndexFileContent() @@ -42,46 +40,62 @@ func (vs *VolumeServer) VolumeSyncIndex(ctx context.Context, req *volume_server_ glog.V(2).Infof("sync volume %d index", req.VolumdId) } - resp.IndexFileContent = content + const blockSizeLimit = 1024 * 1024 * 2 + for i := 0; i < len(content); i += blockSizeLimit { + blockSize := len(content) - i + if blockSize > blockSizeLimit { + blockSize = blockSizeLimit + } + resp := &volume_server_pb.VolumeSyncIndexResponse{} + resp.IndexFileContent = content[i : i+blockSize] + stream.Send(resp) + } - return resp, nil + return nil } -func (vs *VolumeServer) VolumeSyncData(ctx context.Context, req *volume_server_pb.VolumeSyncDataRequest) (*volume_server_pb.VolumeSyncDataResponse, error) { - - resp := &volume_server_pb.VolumeSyncDataResponse{} +func (vs *VolumeServer) VolumeSyncData(req *volume_server_pb.VolumeSyncDataRequest, stream volume_server_pb.VolumeServer_VolumeSyncDataServer) error { v := vs.store.GetVolume(storage.VolumeId(req.VolumdId)) if v == nil { - return nil, fmt.Errorf("Not Found Volume Id %d", req.VolumdId) + return fmt.Errorf("Not Found Volume Id %d", req.VolumdId) } if uint32(v.SuperBlock.CompactRevision) != req.Revision { - return nil, fmt.Errorf("Requested Volume Revision is %d, but current revision is %d", req.Revision, v.SuperBlock.CompactRevision) + return fmt.Errorf("Requested Volume Revision is %d, but current revision is %d", req.Revision, v.SuperBlock.CompactRevision) } content, err := storage.ReadNeedleBlob(v.DataFile(), int64(req.Offset)*types.NeedlePaddingSize, req.Size, v.Version()) if err != nil { - return nil, fmt.Errorf("read offset:%d size:%d", req.Offset, req.Size) + return fmt.Errorf("read offset:%d size:%d", req.Offset, req.Size) } id, err := types.ParseNeedleId(req.NeedleId) if err != nil { - return nil, fmt.Errorf("parsing needle id %s: %v", req.NeedleId, err) + return fmt.Errorf("parsing needle id %s: %v", req.NeedleId, err) } n := new(storage.Needle) n.ParseNeedleHeader(content) if id != n.Id { - return nil, fmt.Errorf("Expected file entry id %d, but found %d", id, n.Id) + return fmt.Errorf("Expected file entry id %d, but found %d", id, n.Id) } if err != nil { glog.Errorf("sync volume %d data: %v", req.VolumdId, err) } - resp.FileContent = content + const blockSizeLimit = 1024 * 1024 * 2 + for i := 0; i < len(content); i += blockSizeLimit { + blockSize := len(content) - i + if blockSize > blockSizeLimit { + blockSize = blockSizeLimit + } + resp := &volume_server_pb.VolumeSyncDataResponse{} + resp.FileContent = content[i : i+blockSize] + stream.Send(resp) + } - return resp, nil + return nil } |
