mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-01 20:26:27 +00:00
Compare commits
25
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
d1031e7190 | ||
|
|
8b9a48c1b5 | ||
|
|
498922dbb2 | ||
|
|
bba8931d64 | ||
|
|
c238153a3b | ||
|
|
c4c4d227a1 | ||
|
|
60f9dea60a | ||
|
|
cb952ff107 | ||
|
|
fe972eab44 | ||
|
|
7275613cc3 | ||
|
|
3b1755a1ee | ||
|
|
f383230766 | ||
|
|
d00308399d | ||
|
|
e8e50ebc75 | ||
|
|
6cbfb24022 | ||
|
|
d1428df1df | ||
|
|
6b2f5147b9 | ||
|
|
422ecb8953 | ||
|
|
0667c4964e | ||
|
|
5737f77d5f | ||
|
|
55bf25e338 | ||
|
|
7df898e00b | ||
|
|
3be91b2125 | ||
|
|
d9e19cc49b | ||
|
|
aded68d69e |
@@ -116,7 +116,7 @@ func (ce *CommandEnv) AdjustedUrl(location *filer_pb.Location) string {
|
||||
}
|
||||
|
||||
func (ce *CommandEnv) GetDataCenter() string {
|
||||
return ce.MasterClient.DataCenter
|
||||
return ce.MasterClient.GetDataCenter()
|
||||
}
|
||||
|
||||
func parseFilerUrl(entryPath string) (filerServer string, filerPort int64, path string, err error) {
|
||||
|
||||
+112
-19
@@ -35,10 +35,10 @@ type MasterClient struct {
|
||||
masters pb.ServerDiscovery
|
||||
grpcDialOption grpc.DialOption
|
||||
|
||||
// TODO: CRITICAL - Data race: resetVidMap() writes to vidMap while other methods read concurrently
|
||||
// This embedded *vidMap should be changed to a private field protected by sync.RWMutex
|
||||
// See: https://github.com/seaweedfs/seaweedfs/issues/[ISSUE_NUMBER]
|
||||
*vidMap
|
||||
// vidMap stores volume location mappings
|
||||
// Protected by vidMapLock to prevent race conditions during pointer swaps in resetVidMap
|
||||
vidMap *vidMap
|
||||
vidMapLock sync.RWMutex
|
||||
vidMapCacheSize int
|
||||
OnPeerUpdate func(update *master_pb.ClusterNodeUpdate, startFrom time.Time)
|
||||
OnPeerUpdateLock sync.RWMutex
|
||||
@@ -71,8 +71,13 @@ func (mc *MasterClient) GetLookupFileIdFunction() LookupFileIdFunctionType {
|
||||
}
|
||||
|
||||
func (mc *MasterClient) LookupFileIdWithFallback(ctx context.Context, fileId string) (fullUrls []string, err error) {
|
||||
// Try cache first using the fast path
|
||||
fullUrls, err = mc.vidMap.LookupFileId(ctx, fileId)
|
||||
// Try cache first using the fast path - grab both vidMap and dataCenter in one lock
|
||||
mc.vidMapLock.RLock()
|
||||
vm := mc.vidMap
|
||||
dataCenter := vm.DataCenter
|
||||
mc.vidMapLock.RUnlock()
|
||||
|
||||
fullUrls, err = vm.LookupFileId(ctx, fileId)
|
||||
if err == nil && len(fullUrls) > 0 {
|
||||
return
|
||||
}
|
||||
@@ -99,7 +104,7 @@ func (mc *MasterClient) LookupFileIdWithFallback(ctx context.Context, fileId str
|
||||
var sameDcUrls, otherDcUrls []string
|
||||
for _, loc := range locations {
|
||||
httpUrl := "http://" + loc.Url + "/" + fileId
|
||||
if mc.DataCenter != "" && mc.DataCenter == loc.DataCenter {
|
||||
if dataCenter != "" && dataCenter == loc.DataCenter {
|
||||
sameDcUrls = append(sameDcUrls, httpUrl)
|
||||
} else {
|
||||
otherDcUrls = append(otherDcUrls, httpUrl)
|
||||
@@ -120,6 +125,10 @@ func (mc *MasterClient) LookupVolumeIdsWithFallback(ctx context.Context, volumeI
|
||||
|
||||
// Check cache first and parse volume IDs once
|
||||
vidStringToUint := make(map[string]uint32, len(volumeIds))
|
||||
|
||||
// Get stable pointer to vidMap with minimal lock hold time
|
||||
vm := mc.getStableVidMap()
|
||||
|
||||
for _, vidString := range volumeIds {
|
||||
vid, err := strconv.ParseUint(vidString, 10, 32)
|
||||
if err != nil {
|
||||
@@ -127,7 +136,7 @@ func (mc *MasterClient) LookupVolumeIdsWithFallback(ctx context.Context, volumeI
|
||||
}
|
||||
vidStringToUint[vidString] = uint32(vid)
|
||||
|
||||
locations, found := mc.GetLocations(uint32(vid))
|
||||
locations, found := vm.GetLocations(uint32(vid))
|
||||
if found && len(locations) > 0 {
|
||||
result[vidString] = locations
|
||||
} else {
|
||||
@@ -149,9 +158,12 @@ func (mc *MasterClient) LookupVolumeIdsWithFallback(ctx context.Context, volumeI
|
||||
stillNeedLookup := make([]string, 0, len(needsLookup))
|
||||
batchResult := make(map[string][]Location)
|
||||
|
||||
// Get stable pointer with minimal lock hold time
|
||||
vm := mc.getStableVidMap()
|
||||
|
||||
for _, vidString := range needsLookup {
|
||||
vid := vidStringToUint[vidString] // Use pre-parsed value
|
||||
if locations, found := mc.GetLocations(vid); found && len(locations) > 0 {
|
||||
if locations, found := vm.GetLocations(vid); found && len(locations) > 0 {
|
||||
batchResult[vidString] = locations
|
||||
} else {
|
||||
stillNeedLookup = append(stillNeedLookup, vidString)
|
||||
@@ -196,7 +208,7 @@ func (mc *MasterClient) LookupVolumeIdsWithFallback(ctx context.Context, volumeI
|
||||
GrpcPort: int(masterLoc.GrpcPort),
|
||||
DataCenter: masterLoc.DataCenter,
|
||||
}
|
||||
mc.vidMap.addLocation(uint32(vid), loc)
|
||||
mc.addLocation(uint32(vid), loc)
|
||||
locations = append(locations, loc)
|
||||
}
|
||||
|
||||
@@ -351,7 +363,7 @@ func (mc *MasterClient) tryConnectToMaster(ctx context.Context, master pb.Server
|
||||
|
||||
if err = stream.Send(&master_pb.KeepConnectedRequest{
|
||||
FilerGroup: mc.FilerGroup,
|
||||
DataCenter: mc.DataCenter,
|
||||
DataCenter: mc.GetDataCenter(),
|
||||
Rack: mc.rack,
|
||||
ClientType: mc.clientType,
|
||||
ClientAddress: string(mc.clientHost),
|
||||
@@ -482,15 +494,96 @@ func (mc *MasterClient) WithClientCustomGetMaster(getMasterF func() pb.ServerAdd
|
||||
})
|
||||
}
|
||||
|
||||
func (mc *MasterClient) resetVidMap() {
|
||||
tail := &vidMap{
|
||||
vid2Locations: mc.vid2Locations,
|
||||
ecVid2Locations: mc.ecVid2Locations,
|
||||
DataCenter: mc.DataCenter,
|
||||
cache: mc.cache,
|
||||
}
|
||||
// getStableVidMap gets a stable pointer to the vidMap, releasing the lock immediately.
|
||||
// This is safe for read operations as the returned pointer is a stable snapshot,
|
||||
// and the underlying vidMap methods have their own internal locking.
|
||||
func (mc *MasterClient) getStableVidMap() *vidMap {
|
||||
mc.vidMapLock.RLock()
|
||||
vm := mc.vidMap
|
||||
mc.vidMapLock.RUnlock()
|
||||
return vm
|
||||
}
|
||||
|
||||
nvm := newVidMap(mc.DataCenter)
|
||||
// withCurrentVidMap executes a function with the current vidMap under a read lock.
|
||||
// This is for methods that modify vidMap's internal state, ensuring the pointer
|
||||
// is not swapped by resetVidMap during the operation. The actual map mutations
|
||||
// are protected by vidMap's internal mutex.
|
||||
func (mc *MasterClient) withCurrentVidMap(f func(vm *vidMap)) {
|
||||
mc.vidMapLock.RLock()
|
||||
defer mc.vidMapLock.RUnlock()
|
||||
f(mc.vidMap)
|
||||
}
|
||||
|
||||
// Public methods for external packages to access vidMap safely
|
||||
|
||||
// GetLocations safely retrieves volume locations
|
||||
func (mc *MasterClient) GetLocations(vid uint32) (locations []Location, found bool) {
|
||||
return mc.getStableVidMap().GetLocations(vid)
|
||||
}
|
||||
|
||||
// GetLocationsClone safely retrieves a clone of volume locations
|
||||
func (mc *MasterClient) GetLocationsClone(vid uint32) (locations []Location, found bool) {
|
||||
return mc.getStableVidMap().GetLocationsClone(vid)
|
||||
}
|
||||
|
||||
// GetVidLocations safely retrieves volume locations by string ID
|
||||
func (mc *MasterClient) GetVidLocations(vid string) (locations []Location, err error) {
|
||||
return mc.getStableVidMap().GetVidLocations(vid)
|
||||
}
|
||||
|
||||
// LookupFileId safely looks up URLs for a file ID
|
||||
func (mc *MasterClient) LookupFileId(ctx context.Context, fileId string) (fullUrls []string, err error) {
|
||||
return mc.getStableVidMap().LookupFileId(ctx, fileId)
|
||||
}
|
||||
|
||||
// LookupVolumeServerUrl safely looks up volume server URLs
|
||||
func (mc *MasterClient) LookupVolumeServerUrl(vid string) (serverUrls []string, err error) {
|
||||
return mc.getStableVidMap().LookupVolumeServerUrl(vid)
|
||||
}
|
||||
|
||||
// GetDataCenter safely retrieves the data center
|
||||
func (mc *MasterClient) GetDataCenter() string {
|
||||
return mc.getStableVidMap().DataCenter
|
||||
}
|
||||
|
||||
// Thread-safe helpers for vidMap operations
|
||||
|
||||
// addLocation adds a volume location
|
||||
func (mc *MasterClient) addLocation(vid uint32, location Location) {
|
||||
mc.withCurrentVidMap(func(vm *vidMap) {
|
||||
vm.addLocation(vid, location)
|
||||
})
|
||||
}
|
||||
|
||||
// deleteLocation removes a volume location
|
||||
func (mc *MasterClient) deleteLocation(vid uint32, location Location) {
|
||||
mc.withCurrentVidMap(func(vm *vidMap) {
|
||||
vm.deleteLocation(vid, location)
|
||||
})
|
||||
}
|
||||
|
||||
// addEcLocation adds an EC volume location
|
||||
func (mc *MasterClient) addEcLocation(vid uint32, location Location) {
|
||||
mc.withCurrentVidMap(func(vm *vidMap) {
|
||||
vm.addEcLocation(vid, location)
|
||||
})
|
||||
}
|
||||
|
||||
// deleteEcLocation removes an EC volume location
|
||||
func (mc *MasterClient) deleteEcLocation(vid uint32, location Location) {
|
||||
mc.withCurrentVidMap(func(vm *vidMap) {
|
||||
vm.deleteEcLocation(vid, location)
|
||||
})
|
||||
}
|
||||
|
||||
func (mc *MasterClient) resetVidMap() {
|
||||
mc.vidMapLock.Lock()
|
||||
defer mc.vidMapLock.Unlock()
|
||||
|
||||
// Create a shallow clone to preserve in the cache chain
|
||||
tail := mc.vidMap.shallowClone()
|
||||
|
||||
nvm := newVidMap(tail.DataCenter)
|
||||
nvm.cache = tail
|
||||
mc.vidMap = nvm
|
||||
|
||||
|
||||
@@ -53,6 +53,17 @@ func newVidMap(dataCenter string) *vidMap {
|
||||
}
|
||||
}
|
||||
|
||||
// shallowClone creates a shallow copy of the vidMap for use in cache chaining.
|
||||
// The caller is responsible for ensuring thread safety.
|
||||
func (vc *vidMap) shallowClone() *vidMap {
|
||||
return &vidMap{
|
||||
vid2Locations: vc.vid2Locations,
|
||||
ecVid2Locations: vc.ecVid2Locations,
|
||||
DataCenter: vc.DataCenter,
|
||||
cache: vc.cache,
|
||||
}
|
||||
}
|
||||
|
||||
func (vc *vidMap) getLocationIndex(length int) (int, error) {
|
||||
if length <= 0 {
|
||||
return 0, fmt.Errorf("invalid length: %d", length)
|
||||
|
||||
Reference in New Issue
Block a user