Compare commits

...
Author SHA1 Message Date
Chris LuandGitHub 4d0afa3286 operation: fold extra dial options into WithVolumeServerClient (#11392)
WithVolumeServerClientOptions duplicated the existing helper only to
append extra dial options. Trail the extras as a variadic parameter on
WithVolumeServerClient instead, keeping the (dialOption, fn) argument
order so all existing callers keep working. TailVolumeFromSource gets
the same treatment.
2026-09-18 13:26:31 -07:00
Chris LuandGitHub 09f835842c Merge branch 'master' into fix/volume-copy-source-guard 2026-09-18 13:24:53 -07:00
Chris Lu 934b9b4daf admin: only claim fallback master leadership on an empty raft response
A nonempty RaftListClusterServers response whose entries were all
rejected left masterMap empty, so the fallback marked the reachable
current master as leader the same way a genuinely empty (non-raft)
response does. Track whether the successful response returned zero
servers and only promote the fallback master then.
2026-09-18 12:22:39 -07:00
Chris Lu bd41ce39f7 test: opt erasure-coding loopback clusters out of the remote endpoint guard
The erasure-coding suites drive VolumeEcShardsCopy / VolumeCopy between
volume servers bound to 127.0.0.1, which the copy/tail source guard now
rejects by default. Pass -volume.allowUntrustedRemoteEndpoints to the
test volume launches, matching what the volume_server framework
harnesses already do.
2026-09-18 11:34:05 -07:00
Chris Lu 0d6024e2e0 pb: return empty server address for malformed grpc addresses
GrpcAddressToServerAddress used to return the unparseable input on a
hostAndPort failure, so a malformed raft address (e.g. "host:abc")
flowed into admin dashboard master maps unchanged. Return an empty
string instead, skip empty conversions at the two raft-cluster merge
sites, and drop the now-stale comment about the fatal exit the earlier
commit removed.
2026-09-18 11:34:05 -07:00
Chris Lu 0ae7874ed9 volume: pin validated copy/tail source addresses at dial time
validateReplicaTarget resolves the source hostname once, but the gRPC
client resolved it again at connect, leaving a DNS-rebinding window for
hostname sources. The copy and tail source dials now run through the
same guardedDialerPolicy the remote-storage path uses, so every resolved
address is re-checked against the replica deny list (private peers
allowed) immediately before the TCP connect. guardedDialerPolicy also
moves to util.OutboundDialContext so the guarded path keeps the -ip.bind
source binding the default gRPC dialer had.

The Rust volume server mirrors this with connect_guarded, a tonic
connector that resolves, re-checks each address, and connects to the
first passing IP; handlers use it whenever the untrusted-endpoint
opt-out is off. A handler-level test now exercises the enabled
validation branches for all three source-taking RPCs.
2026-09-18 11:32:28 -07:00
Chris Lu 81778defa1 rust volume: validate copy and tail source addresses before dialing
Mirror the Go guard on the Rust volume server: volume_copy,
volume_ec_shards_copy and volume_tail_receiver dial a caller-supplied
source address, so run it through validate_replica_target first (bare
host:port; no loopback, link-local or unspecified hosts; private peers
stay allowed). --volume.allowUntrustedRemoteEndpoints opts out; the test
fixture and the Rust test-cluster launcher set it so loopback sources in
tests keep working.
2026-09-18 10:57:18 -07:00
Chris Lu 39a8d3253c volume: validate copy and tail source addresses before dialing
VolumeCopy, VolumeEcShardsCopy and VolumeTailReceiver dial a
caller-supplied source address (SourceDataNode / SourceVolumeServer)
with no endpoint validation, so an anonymous caller could aim the volume
server at loopback, link-local (cloud metadata) or other unintended
destinations and read dial behavior back as a connectivity oracle.

Apply the same peer-target deny list FetchAndWriteNeedle uses for
replica targets: the source must be a bare host:port whose host is not
loopback, link-local or unspecified; cluster peers stay reachable on
private networks, and -volume.allowUntrustedRemoteEndpoints opts out.
The loopback-using copy tests set the flag to keep exercising the copy
path in process.
2026-09-18 10:57:18 -07:00
Chris Lu 47f323bbb3 pb: stop exiting the process on malformed server addresses
ServerToGrpcAddress and GrpcAddressToServerAddress called glog.Fatalf
when hostAndPort could not parse the port, which os.Exit(255)ed the whole
process. A caller-supplied copy or tail source address reached this path
synchronously in the serving goroutine, so one anonymous VolumeCopy with
a non-numeric port terminated the volume server.

Log the parse error and return the input unchanged instead: the dial or
request that consumes the address then fails as an ordinary error.
2026-09-18 10:57:14 -07:00
6 changed files with 17 additions and 21 deletions
+5 -9
View File
@@ -10,19 +10,15 @@ import (
"github.com/seaweedfs/seaweedfs/weed/pb/volume_server_pb"
)
func WithVolumeServerClient(streamingMode bool, volumeServer pb.ServerAddress, grpcDialOption grpc.DialOption, fn func(volume_server_pb.VolumeServerClient) error) error {
return WithVolumeServerClientOptions(streamingMode, volumeServer, fn, grpcDialOption)
}
// WithVolumeServerClientOptions is WithVolumeServerClient with extra dial
// options appended after the TLS option, so a caller dialing an untrusted
// source address can pin the validated endpoint at connect time.
func WithVolumeServerClientOptions(streamingMode bool, volumeServer pb.ServerAddress, fn func(volume_server_pb.VolumeServerClient) error, grpcDialOptions ...grpc.DialOption) error {
// WithVolumeServerClient dials with the TLS option plus any extra dial options
// appended after it, so a caller dialing an untrusted source address can pin
// the validated endpoint at connect time.
func WithVolumeServerClient(streamingMode bool, volumeServer pb.ServerAddress, grpcDialOption grpc.DialOption, fn func(volume_server_pb.VolumeServerClient) error, extraDialOptions ...grpc.DialOption) error {
return pb.WithGrpcClient(context.Background(), streamingMode, 0, func(grpcConnection *grpc.ClientConn) error {
client := volume_server_pb.NewVolumeServerClient(grpcConnection)
return fn(client)
}, volumeServer.ToGrpcAddress(), false, grpcDialOptions...)
}, volumeServer.ToGrpcAddress(), false, append([]grpc.DialOption{grpcDialOption}, extraDialOptions...)...)
}
+4 -4
View File
@@ -25,11 +25,11 @@ func TailVolume(masterFn GetMasterFn, grpcDialOption grpc.DialOption, vid needle
volumeServer := lookup.Locations[0].ServerAddress()
return TailVolumeFromSource(volumeServer, vid, sinceNs, timeoutSeconds, fn, grpcDialOption)
return TailVolumeFromSource(volumeServer, vid, sinceNs, timeoutSeconds, grpcDialOption, fn)
}
func TailVolumeFromSource(volumeServer pb.ServerAddress, vid needle.VolumeId, sinceNs uint64, idleTimeoutSeconds int, fn func(n *needle.Needle) error, grpcDialOptions ...grpc.DialOption) error {
return WithVolumeServerClientOptions(true, volumeServer, func(client volume_server_pb.VolumeServerClient) error {
func TailVolumeFromSource(volumeServer pb.ServerAddress, vid needle.VolumeId, sinceNs uint64, idleTimeoutSeconds int, grpcDialOption grpc.DialOption, fn func(n *needle.Needle) error, extraDialOptions ...grpc.DialOption) error {
return WithVolumeServerClient(true, volumeServer, grpcDialOption, func(client volume_server_pb.VolumeServerClient) error {
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
@@ -90,5 +90,5 @@ func TailVolumeFromSource(volumeServer pb.ServerAddress, vid needle.VolumeId, si
}
return nil
}, grpcDialOptions...)
}, extraDialOptions...)
}
+2 -2
View File
@@ -59,7 +59,7 @@ func (vs *VolumeServer) VolumeCopy(req *volume_server_pb.VolumeCopyRequest, stre
var sourceVolumeStatusAfterCopy *volume_server_pb.VolumeStatusResponse
var dataBaseFileName, indexBaseFileName, idxFileName, datFileName string
var hasRemoteDatFile bool
err := operation.WithVolumeServerClientOptions(true, pb.ServerAddress(req.SourceDataNode), func(client volume_server_pb.VolumeServerClient) error {
err := operation.WithVolumeServerClient(true, pb.ServerAddress(req.SourceDataNode), vs.grpcDialOption, func(client volume_server_pb.VolumeServerClient) error {
var err error
sourceVolumeStatus, err = client.VolumeStatus(stream.Context(), &volume_server_pb.VolumeStatusRequest{
VolumeId: req.VolumeId,
@@ -214,7 +214,7 @@ func (vs *VolumeServer) VolumeCopy(req *volume_server_pb.VolumeCopyRequest, stre
}
return nil
}, vs.grpcDialOption, vs.guardedGrpcDialOption(req.SourceDataNode))
}, vs.guardedGrpcDialOption(req.SourceDataNode))
if err != nil {
return err
+2 -2
View File
@@ -392,7 +392,7 @@ func (vs *VolumeServer) VolumeEcShardsCopy(ctx context.Context, req *volume_serv
}
throttler := util.NewWriteThrottler(ioBytePerSecond)
err := operation.WithVolumeServerClientOptions(true, pb.ServerAddress(req.SourceDataNode), func(client volume_server_pb.VolumeServerClient) error {
err := operation.WithVolumeServerClient(true, pb.ServerAddress(req.SourceDataNode), vs.grpcDialOption, func(client volume_server_pb.VolumeServerClient) error {
// copy ec data slices
for _, shardId := range req.ShardIds {
@@ -460,7 +460,7 @@ func (vs *VolumeServer) VolumeEcShardsCopy(ctx context.Context, req *volume_serv
}
}
return nil
}, vs.grpcDialOption, vs.guardedGrpcDialOption(req.SourceDataNode))
}, vs.guardedGrpcDialOption(req.SourceDataNode))
if err != nil {
return nil, fmt.Errorf("VolumeEcShardsCopy volume %d: %v", req.VolumeId, err)
}
+2 -2
View File
@@ -100,10 +100,10 @@ func (vs *VolumeServer) VolumeTailReceiver(ctx context.Context, req *volume_serv
defer glog.V(1).Infof("receive tailing volume %d finished", v.Id)
return resp, operation.TailVolumeFromSource(pb.ServerAddress(req.SourceVolumeServer), v.Id, req.SinceNs, int(req.IdleTimeoutSeconds), func(n *needle.Needle) error {
return resp, operation.TailVolumeFromSource(pb.ServerAddress(req.SourceVolumeServer), v.Id, req.SinceNs, int(req.IdleTimeoutSeconds), vs.grpcDialOption, func(n *needle.Needle) error {
_, err := vs.store.WriteVolumeNeedle(v.Id, n, false, false)
return err
}, vs.grpcDialOption, vs.guardedGrpcDialOption(req.SourceVolumeServer))
}, vs.guardedGrpcDialOption(req.SourceVolumeServer))
}
+2 -2
View File
@@ -232,14 +232,14 @@ func startTailNeedleStream(grpcDialOption grpc.DialOption, volumeId needle.Volum
ch := make(chan *needle.Needle, 32)
stream := &tailNeedleStream{ch: ch}
go func() {
err := operation.TailVolumeFromSource(server, volumeId, 0, mergeIdleTimeoutSeconds, func(n *needle.Needle) error {
err := operation.TailVolumeFromSource(server, volumeId, 0, mergeIdleTimeoutSeconds, grpcDialOption, func(n *needle.Needle) error {
select {
case ch <- n:
case <-done:
return fmt.Errorf("merge cancelled")
}
return nil
}, grpcDialOption)
})
close(ch)
stream.setErr(err)
}()