volume: EC decode onto a clean peer via staged-new-volume adopt (Go+Rust) (#10463)

Decoding EC shards back to a normal volume in place reconstructs <vid>.dat
in the shards' own directory, so the vid is momentarily registered as both
an EC and a normal volume in one location — the load/scan path then sees it
as both, risking mount ambiguity and needle loss. VolumeEcShardsToVolume
still supports that in-place path; this adds the primitives to decode onto
a *clean* peer instead:

  - ReceiveFile gains a staged-new-volume mode: when the volume does not
    exist here and ReceiveFileInfo.disk_type is set, pick a free-slot disk
    of that medium and write <base><ext>.copying (not a valid volume name,
    so the scanner never half-loads a partial push).
  - VolumeEcShardsToVolume gains from_staged: adopt the pushed .dat/.idx/
    .vif — rename .copying into place under a .note in-progress marker,
    then mount — so <vid> lands on the peer only as a normal volume.

The caller decodes the shards off-box and streams the finished volume to a
peer holding no shard of the vid on the target medium. Go and Rust volume
servers get identical handlers. Proto: ReceiveFileInfo.disk_type (12; 8-11
reserved for versioned-EC), VolumeEcShardsToVolumeRequest.from_staged (3) +
disk_type (4).

Claude-Session: https://claude.ai/code/session_01Ks16jnt4S7gdDk8cheQ3xu
This commit is contained in:
Chris Lu
2026-07-27 17:56:17 -07:00
committed by GitHub
parent 2bea4dd610
commit 9c37e52c9b
6 changed files with 294 additions and 29 deletions
+14
View File
@@ -343,6 +343,13 @@ message ReceiveFileInfo {
uint32 shard_id = 5;
uint64 file_size = 6;
uint32 disk_id = 7; // EC shard disk; 0 = auto-select (see VolumeEcShardsCopyRequest.disk_id)
// Field numbers 8-11 are reserved for versioned-EC; disk_type stays at 12.
// Staged-new-volume mode (EC decode onto a clean peer): set on a non-EC push
// whose volume does not yet exist on this server. The server picks a disk
// location of this medium with a free slot and writes <base><ext>.copying,
// finalized by VolumeEcShardsToVolume(from_staged).
string disk_type = 12;
}
message ReceiveFileResponse {
@@ -503,6 +510,13 @@ message VolumeEcBlobDeleteResponse {
message VolumeEcShardsToVolumeRequest {
uint32 volume_id = 1;
string collection = 2;
// Staged mode: the caller already decoded the EC shards off-box and pushed
// .dat/.idx/.vif as <base><ext>.copying to this server (ReceiveFile
// staged-new-volume mode). Adopt them as a normal volume instead of decoding
// local EC shards in place, so <vid> is never registered as both EC and
// normal on one disk.
bool from_staged = 3;
string disk_type = 4; // target medium's disk location for the normal volume (staged mode)
}
message VolumeEcShardsToVolumeResponse {
}
+124 -10
View File
@@ -1753,16 +1753,47 @@ impl VolumeServer for VolumeGrpcService {
format!("{}/{}{}", dir, ec_base, info.ext)
} else {
let store = self.state.store.read().unwrap();
let (_, v) =
store.find_volume(VolumeId(info.volume_id)).ok_or_else(|| {
Status::not_found(format!(
"volume {} not found",
info.volume_id
))
})?;
let p = v.file_name(&info.ext);
drop(store);
p
let existing = store
.find_volume(VolumeId(info.volume_id))
.map(|(_, v)| v.file_name(&info.ext));
if let Some(p) = existing {
drop(store);
p
} else if !info.disk_type.is_empty() {
// Staged-new-volume mode (EC decode onto a clean peer):
// the volume does not exist here yet. Pick a free-slot
// disk location of the requested medium and stage the
// file as <base><ext>.copying, to be renamed into place
// and mounted by VolumeEcShardsToVolume(from_staged).
// .idx.copying/.vif.copying are not valid volume names,
// so the scanner never half-loads a partial push.
let want = DiskType::from_string(&info.disk_type);
match store.find_free_location_predicate(|l| l.disk_type == want) {
Some(i) => {
let dir = store.locations[i].directory.clone();
drop(store);
let vfile = if info.collection.is_empty() {
format!("{}/{}", dir, info.volume_id)
} else {
format!("{}/{}_{}", dir, info.collection, info.volume_id)
};
format!("{}{}.copying", vfile, info.ext)
}
None => {
drop(store);
return Err(Status::internal(format!(
"no {} disk location with a free slot for volume {}",
info.disk_type, info.volume_id
)));
}
}
} else {
drop(store);
return Err(Status::not_found(format!(
"volume {} not found",
info.volume_id
)));
}
};
target_file = Some(std::fs::File::create(&path).map_err(|e| {
@@ -3095,6 +3126,89 @@ impl VolumeServer for VolumeGrpcService {
let req = request.into_inner();
let vid = VolumeId(req.volume_id);
// Staged mode: the caller decoded off-box and streamed the normal volume
// here as <base><ext>.copying (ReceiveFile staged-new-volume mode). Adopt
// those files as a normal volume — rename .copying into place under a .note
// marker, then mount. This server holds no EC shards for the vid, so there
// is no in-place decode to run.
if req.from_staged {
let want = DiskType::from_string(&req.disk_type);
let base = {
let store = self.state.store.read().unwrap();
if store.has_volume(vid) {
return Err(Status::internal(format!(
"staged volume {} already exists on this server",
req.volume_id
)));
}
let mut base: Option<String> = None;
for loc in store.locations.iter() {
if loc.disk_type != want {
continue;
}
let candidate = if req.collection.is_empty() {
format!("{}/{}", loc.directory, req.volume_id)
} else {
format!("{}/{}_{}", loc.directory, req.collection, req.volume_id)
};
if std::path::Path::new(&format!("{}.dat.copying", candidate)).exists() {
base = Some(candidate);
break;
}
}
base
}
.ok_or_else(|| {
Status::not_found(format!(
"staged volume {}: no .dat.copying found on a {} disk",
req.volume_id, req.disk_type
))
})?;
for ext in [".dat", ".idx", ".vif"] {
if !std::path::Path::new(&format!("{}{}.copying", base, ext)).exists() {
return Err(Status::not_found(format!(
"staged volume {} missing {}.copying",
req.volume_id, ext
)));
}
}
// .note in-progress marker (VolumeCopy discipline): a crash mid-rename
// leaves a .note that fails the load and sweeps the partial volume.
let note = format!("{}.note", base);
std::fs::write(&note, format!("adopting decoded volume {}", req.volume_id))
.map_err(|e| Status::internal(format!("write .note: {}", e)))?;
// Rename staged files into place, then drop the .note before mounting —
// a volume that still carries a .note is swept by the load scan.
for ext in [".vif", ".dat", ".idx"] {
let src = format!("{}{}.copying", base, ext);
let dst = format!("{}{}", base, ext);
if let Err(e) = std::fs::rename(&src, &dst) {
let _ = std::fs::remove_file(&note);
return Err(Status::internal(format!("rename staged {}: {}", ext, e)));
}
}
let _ = std::fs::remove_file(&note);
{
let mut store = self.state.store.write().unwrap();
store.mount_volume_by_id(vid).map_err(|e| {
Status::internal(format!("mount staged volume {}: {}", req.volume_id, e))
})?;
}
self.state.volume_state_notify.notify_one();
tracing::info!(
"volume_ec_shards_to_volume: adopted decoded volume {} from staging ({})",
req.volume_id,
base
);
return Ok(Response::new(
volume_server_pb::VolumeEcShardsToVolumeResponse {},
));
}
let store = self.state.store.read().unwrap();
// Aggregate per-shard data dirs across all locations so the
// shard-presence check + decoder both see the union for
+14
View File
@@ -351,6 +351,13 @@ message ReceiveFileInfo {
uint32 shard_id = 5;
uint64 file_size = 6;
uint32 disk_id = 7; // EC shard disk; 0 = auto-select (see VolumeEcShardsCopyRequest.disk_id)
// Field numbers 8-11 are reserved for versioned-EC; disk_type stays at 12.
// Staged-new-volume mode (EC decode onto a clean peer): set on a non-EC push
// whose volume does not yet exist on this server. The server picks a disk
// location of this medium with a free slot and writes <base><ext>.copying,
// finalized by VolumeEcShardsToVolume(from_staged).
string disk_type = 12;
}
message ReceiveFileResponse {
@@ -512,6 +519,13 @@ message VolumeEcBlobDeleteResponse {
message VolumeEcShardsToVolumeRequest {
uint32 volume_id = 1;
string collection = 2;
// Staged mode: the caller already decoded the EC shards off-box and pushed
// .dat/.idx/.vif as <base><ext>.copying to this server (ReceiveFile
// staged-new-volume mode). Adopt them as a normal volume instead of decoding
// local EC shards in place, so <vid> is never registered as both EC and
// normal on one disk.
bool from_staged = 3;
string disk_type = 4; // target medium's disk location for the normal volume (staged mode)
}
message VolumeEcShardsToVolumeResponse {
}
+53 -15
View File
@@ -2465,14 +2465,20 @@ func (*ReceiveFileRequest_Info) isReceiveFileRequest_Data() {}
func (*ReceiveFileRequest_FileContent) isReceiveFileRequest_Data() {}
type ReceiveFileInfo struct {
state protoimpl.MessageState `protogen:"open.v1"`
VolumeId uint32 `protobuf:"varint,1,opt,name=volume_id,json=volumeId,proto3" json:"volume_id,omitempty"`
Ext string `protobuf:"bytes,2,opt,name=ext,proto3" json:"ext,omitempty"`
Collection string `protobuf:"bytes,3,opt,name=collection,proto3" json:"collection,omitempty"`
IsEcVolume bool `protobuf:"varint,4,opt,name=is_ec_volume,json=isEcVolume,proto3" json:"is_ec_volume,omitempty"`
ShardId uint32 `protobuf:"varint,5,opt,name=shard_id,json=shardId,proto3" json:"shard_id,omitempty"`
FileSize uint64 `protobuf:"varint,6,opt,name=file_size,json=fileSize,proto3" json:"file_size,omitempty"`
DiskId uint32 `protobuf:"varint,7,opt,name=disk_id,json=diskId,proto3" json:"disk_id,omitempty"` // EC shard disk; 0 = auto-select (see VolumeEcShardsCopyRequest.disk_id)
state protoimpl.MessageState `protogen:"open.v1"`
VolumeId uint32 `protobuf:"varint,1,opt,name=volume_id,json=volumeId,proto3" json:"volume_id,omitempty"`
Ext string `protobuf:"bytes,2,opt,name=ext,proto3" json:"ext,omitempty"`
Collection string `protobuf:"bytes,3,opt,name=collection,proto3" json:"collection,omitempty"`
IsEcVolume bool `protobuf:"varint,4,opt,name=is_ec_volume,json=isEcVolume,proto3" json:"is_ec_volume,omitempty"`
ShardId uint32 `protobuf:"varint,5,opt,name=shard_id,json=shardId,proto3" json:"shard_id,omitempty"`
FileSize uint64 `protobuf:"varint,6,opt,name=file_size,json=fileSize,proto3" json:"file_size,omitempty"`
DiskId uint32 `protobuf:"varint,7,opt,name=disk_id,json=diskId,proto3" json:"disk_id,omitempty"` // EC shard disk; 0 = auto-select (see VolumeEcShardsCopyRequest.disk_id)
// Field numbers 8-11 are reserved for versioned-EC; disk_type stays at 12.
// Staged-new-volume mode (EC decode onto a clean peer): set on a non-EC push
// whose volume does not yet exist on this server. The server picks a disk
// location of this medium with a free slot and writes <base><ext>.copying,
// finalized by VolumeEcShardsToVolume(from_staged).
DiskType string `protobuf:"bytes,12,opt,name=disk_type,json=diskType,proto3" json:"disk_type,omitempty"`
unknownFields protoimpl.UnknownFields
sizeCache protoimpl.SizeCache
}
@@ -2556,6 +2562,13 @@ func (x *ReceiveFileInfo) GetDiskId() uint32 {
return 0
}
func (x *ReceiveFileInfo) GetDiskType() string {
if x != nil {
return x.DiskType
}
return ""
}
type ReceiveFileResponse struct {
state protoimpl.MessageState `protogen:"open.v1"`
BytesWritten uint64 `protobuf:"varint,1,opt,name=bytes_written,json=bytesWritten,proto3" json:"bytes_written,omitempty"`
@@ -4257,9 +4270,16 @@ func (*VolumeEcBlobDeleteResponse) Descriptor() ([]byte, []int) {
}
type VolumeEcShardsToVolumeRequest struct {
state protoimpl.MessageState `protogen:"open.v1"`
VolumeId uint32 `protobuf:"varint,1,opt,name=volume_id,json=volumeId,proto3" json:"volume_id,omitempty"`
Collection string `protobuf:"bytes,2,opt,name=collection,proto3" json:"collection,omitempty"`
state protoimpl.MessageState `protogen:"open.v1"`
VolumeId uint32 `protobuf:"varint,1,opt,name=volume_id,json=volumeId,proto3" json:"volume_id,omitempty"`
Collection string `protobuf:"bytes,2,opt,name=collection,proto3" json:"collection,omitempty"`
// Staged mode: the caller already decoded the EC shards off-box and pushed
// .dat/.idx/.vif as <base><ext>.copying to this server (ReceiveFile
// staged-new-volume mode). Adopt them as a normal volume instead of decoding
// local EC shards in place, so <vid> is never registered as both EC and
// normal on one disk.
FromStaged bool `protobuf:"varint,3,opt,name=from_staged,json=fromStaged,proto3" json:"from_staged,omitempty"`
DiskType string `protobuf:"bytes,4,opt,name=disk_type,json=diskType,proto3" json:"disk_type,omitempty"` // target medium's disk location for the normal volume (staged mode)
unknownFields protoimpl.UnknownFields
sizeCache protoimpl.SizeCache
}
@@ -4308,6 +4328,20 @@ func (x *VolumeEcShardsToVolumeRequest) GetCollection() string {
return ""
}
func (x *VolumeEcShardsToVolumeRequest) GetFromStaged() bool {
if x != nil {
return x.FromStaged
}
return false
}
func (x *VolumeEcShardsToVolumeRequest) GetDiskType() string {
if x != nil {
return x.DiskType
}
return ""
}
type VolumeEcShardsToVolumeResponse struct {
state protoimpl.MessageState `protogen:"open.v1"`
unknownFields protoimpl.UnknownFields
@@ -7285,7 +7319,7 @@ const file_volume_server_proto_rawDesc = "" +
"\x12ReceiveFileRequest\x127\n" +
"\x04info\x18\x01 \x01(\v2!.volume_server_pb.ReceiveFileInfoH\x00R\x04info\x12#\n" +
"\ffile_content\x18\x02 \x01(\fH\x00R\vfileContentB\x06\n" +
"\x04data\"\xd3\x01\n" +
"\x04data\"\xf0\x01\n" +
"\x0fReceiveFileInfo\x12\x1b\n" +
"\tvolume_id\x18\x01 \x01(\rR\bvolumeId\x12\x10\n" +
"\x03ext\x18\x02 \x01(\tR\x03ext\x12\x1e\n" +
@@ -7296,7 +7330,8 @@ const file_volume_server_proto_rawDesc = "" +
"isEcVolume\x12\x19\n" +
"\bshard_id\x18\x05 \x01(\rR\ashardId\x12\x1b\n" +
"\tfile_size\x18\x06 \x01(\x04R\bfileSize\x12\x17\n" +
"\adisk_id\x18\a \x01(\rR\x06diskId\"P\n" +
"\adisk_id\x18\a \x01(\rR\x06diskId\x12\x1b\n" +
"\tdisk_type\x18\f \x01(\tR\bdiskType\"P\n" +
"\x13ReceiveFileResponse\x12#\n" +
"\rbytes_written\x18\x01 \x01(\x04R\fbytesWritten\x12\x14\n" +
"\x05error\x18\x02 \x01(\tR\x05error\"`\n" +
@@ -7431,12 +7466,15 @@ const file_volume_server_proto_rawDesc = "" +
"collection\x12\x19\n" +
"\bfile_key\x18\x03 \x01(\x04R\afileKey\x12\x18\n" +
"\aversion\x18\x04 \x01(\rR\aversion\"\x1c\n" +
"\x1aVolumeEcBlobDeleteResponse\"\\\n" +
"\x1aVolumeEcBlobDeleteResponse\"\x9a\x01\n" +
"\x1dVolumeEcShardsToVolumeRequest\x12\x1b\n" +
"\tvolume_id\x18\x01 \x01(\rR\bvolumeId\x12\x1e\n" +
"\n" +
"collection\x18\x02 \x01(\tR\n" +
"collection\" \n" +
"collection\x12\x1f\n" +
"\vfrom_staged\x18\x03 \x01(\bR\n" +
"fromStaged\x12\x1b\n" +
"\tdisk_type\x18\x04 \x01(\tR\bdiskType\" \n" +
"\x1eVolumeEcShardsToVolumeResponse\"8\n" +
"\x19VolumeEcShardsInfoRequest\x12\x1b\n" +
"\tvolume_id\x18\x01 \x01(\rR\bvolumeId\"\xcf\x01\n" +
+23 -4
View File
@@ -687,12 +687,31 @@ func (vs *VolumeServer) ReceiveFile(stream volume_server_pb.VolumeServer_Receive
// Regular volume file
v := vs.store.GetVolume(needle.VolumeId(fileInfo.VolumeId))
if v == nil {
glog.Errorf("ReceiveFile: volume %d not found", fileInfo.VolumeId)
return stream.SendAndClose(&volume_server_pb.ReceiveFileResponse{
Error: fmt.Sprintf("volume %d not found", fileInfo.VolumeId),
if fileInfo.DiskType == "" {
glog.Errorf("ReceiveFile: volume %d not found", fileInfo.VolumeId)
return stream.SendAndClose(&volume_server_pb.ReceiveFileResponse{
Error: fmt.Sprintf("volume %d not found", fileInfo.VolumeId),
})
}
// Staged-new-volume mode (EC decode onto a clean peer): the
// volume does not exist here yet. Pick a free-slot disk location
// of the requested medium and stage the file as
// <base><ext>.copying, to be renamed into place and mounted by
// VolumeEcShardsToVolume(from_staged). .idx.copying/.vif.copying
// are not valid volume names, so the scanner never half-loads.
want := types.ToDiskType(fileInfo.DiskType)
loc := vs.store.FindFreeLocation(func(l *storage.DiskLocation) bool {
return l.DiskType == want
})
if loc == nil {
return stream.SendAndClose(&volume_server_pb.ReceiveFileResponse{
Error: fmt.Sprintf("no %s disk location with a free slot for volume %d", fileInfo.DiskType, fileInfo.VolumeId),
})
}
filePath = storage.VolumeFileName(loc.Directory, fileInfo.Collection, int(fileInfo.VolumeId)) + fileInfo.Ext + ".copying"
} else {
filePath = v.FileName(fileInfo.Ext)
}
filePath = v.FileName(fileInfo.Ext)
}
// Create target file
+66
View File
@@ -921,6 +921,14 @@ func (vs *VolumeServer) VolumeEcShardsToVolume(ctx context.Context, req *volume_
glog.V(0).Infof("VolumeEcShardsToVolume: %v", req)
// Staged mode: the caller decoded the shards off-box and streamed the normal
// volume here as <base><ext>.copying (ReceiveFile staged-new-volume). Adopt
// those files as a normal volume; this server holds no EC shards for the vid,
// so there is no local EC decode to run.
if req.FromStaged {
return vs.adoptStagedVolume(req)
}
// Collect all EC shards (NewEcVolume will load EC config from .vif into v.ECContext)
// Use MaxShardCount (32) to support custom EC ratios up to 32 total shards
tempShards := make([]string, erasure_coding.MaxShardCount)
@@ -1044,6 +1052,64 @@ func (vs *VolumeServer) VolumeEcShardsToVolume(ctx context.Context, req *volume_
return &volume_server_pb.VolumeEcShardsToVolumeResponse{}, nil
}
// adoptStagedVolume finalizes a normal volume the caller decoded off-box and
// streamed here as <base><ext>.copying (ReceiveFile staged-new-volume mode). It
// renames the staged files into place under a .note in-progress marker and
// mounts the volume, so <vid> is registered here only as a normal volume — never
// as an EC/normal twin in one directory.
func (vs *VolumeServer) adoptStagedVolume(req *volume_server_pb.VolumeEcShardsToVolumeRequest) (*volume_server_pb.VolumeEcShardsToVolumeResponse, error) {
vid := needle.VolumeId(req.VolumeId)
if vs.store.GetVolume(vid) != nil {
return nil, fmt.Errorf("staged volume %d already exists on this server", req.VolumeId)
}
want := types.ToDiskType(req.DiskType)
// Locate the disk the ReceiveFile push staged onto: the disk_type location
// whose <base>.dat.copying exists.
var base string
for _, l := range vs.store.Locations {
if l.DiskType != want {
continue
}
candidate := storage.VolumeFileName(l.Directory, req.Collection, int(req.VolumeId))
if util.FileExists(candidate + ".dat.copying") {
base = candidate
break
}
}
if base == "" {
return nil, fmt.Errorf("staged volume %d: no .dat.copying found on a %s disk", req.VolumeId, req.DiskType)
}
for _, ext := range []string{".dat", ".idx", ".vif"} {
if !util.FileExists(base + ext + ".copying") {
return nil, fmt.Errorf("staged volume %d missing %s.copying", req.VolumeId, ext)
}
}
// .note in-progress marker (VolumeCopy discipline): a crash mid-rename leaves a
// .note that fails the load and sweeps the partial volume on restart.
noteFile := base + ".note"
if err := util.WriteFile(noteFile, []byte(fmt.Sprintf("adopting decoded volume %d", req.VolumeId)), 0644); err != nil {
return nil, fmt.Errorf("write .note for volume %d: %w", req.VolumeId, err)
}
// Rename staged files into place, then drop the .note before mounting — a
// volume that still carries a .note is swept by loadExistingVolume.
for _, ext := range []string{".vif", ".dat", ".idx"} {
if err := os.Rename(base+ext+".copying", base+ext); err != nil {
os.Remove(noteFile)
return nil, fmt.Errorf("rename staged %s for volume %d: %w", ext, req.VolumeId, err)
}
}
os.Remove(noteFile)
if err := vs.store.MountVolume(vid); err != nil {
return nil, fmt.Errorf("mount staged volume %d: %w", req.VolumeId, err)
}
glog.V(0).Infof("VolumeEcShardsToVolume: adopted decoded volume %d from staging (%s)", req.VolumeId, base)
return &volume_server_pb.VolumeEcShardsToVolumeResponse{}, nil
}
func (vs *VolumeServer) VolumeEcShardsInfo(ctx context.Context, req *volume_server_pb.VolumeEcShardsInfoRequest) (*volume_server_pb.VolumeEcShardsInfoResponse, error) {
glog.V(0).Infof("VolumeEcShardsInfo: volume %d", req.VolumeId)