aboutsummaryrefslogtreecommitdiff
path: root/weed/server/volume_grpc_sync.go
diff options
context:
space:
mode:
Diffstat (limited to 'weed/server/volume_grpc_sync.go')
-rw-r--r--weed/server/volume_grpc_sync.go46
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
}