mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-06 14:45:51 +00:00
Compare commits
9
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
d5e05d84d7 | ||
|
|
ec49e5a28f | ||
|
|
5e03b045aa | ||
|
|
a07f965e99 | ||
|
|
39688c8a40 | ||
|
|
cc634da5bd | ||
|
|
c1c73b0a5d | ||
|
|
f20c05ea0a | ||
|
|
e064c86d12 |
@@ -152,6 +152,9 @@ func init() {
|
||||
serverOptions.v.inflightUploadDataTimeout = cmdServer.Flag.Duration("volume.inflightUploadDataTimeout", 60*time.Second, "inflight upload data wait timeout of volume servers")
|
||||
serverOptions.v.inflightDownloadDataTimeout = cmdServer.Flag.Duration("volume.inflightDownloadDataTimeout", 60*time.Second, "inflight download data wait timeout of volume servers")
|
||||
|
||||
serverOptions.v.udsListen = cmdServer.Flag.String("volume.uds.listen", "", "Unix domain socket path for RDMA sidecar locate API (e.g., /tmp/sra-volume.sock)")
|
||||
serverOptions.v.udsTransport = cmdServer.Flag.String("volume.uds.transport", "", "Unix domain socket path for outbound RDMA replication (e.g., /tmp/sra-transport.sock)")
|
||||
|
||||
serverOptions.v.hasSlowRead = cmdServer.Flag.Bool("volume.hasSlowRead", true, "<experimental> if true, this prevents slow reads from blocking other requests, but large file read P99 latency will increase.")
|
||||
serverOptions.v.readBufferSizeMB = cmdServer.Flag.Int("volume.readBufferSizeMB", 4, "<experimental> larger values can optimize query performance but will increase some memory usage,Use with hasSlowRead normally")
|
||||
|
||||
|
||||
+28
-3
@@ -74,6 +74,8 @@ type VolumeServerOptions struct {
|
||||
ldbTimeout *int64
|
||||
debug *bool
|
||||
debugPort *int
|
||||
udsListen *string // UDS socket path for RDMA sidecar integration
|
||||
udsTransport *string // UDS socket path for outbound RDMA replication
|
||||
}
|
||||
|
||||
func init() {
|
||||
@@ -114,6 +116,8 @@ func init() {
|
||||
v.readBufferSizeMB = cmdVolume.Flag.Int("readBufferSizeMB", 4, "<experimental> larger values can optimize query performance but will increase some memory usage,Use with hasSlowRead normally.")
|
||||
v.debug = cmdVolume.Flag.Bool("debug", false, "serves runtime profiling data via pprof on the port specified by -debug.port")
|
||||
v.debugPort = cmdVolume.Flag.Int("debug.port", 6060, "http port for debugging")
|
||||
v.udsListen = cmdVolume.Flag.String("uds.listen", "", "Unix domain socket path for RDMA sidecar locate API (e.g., /tmp/sra-volume.sock)")
|
||||
v.udsTransport = cmdVolume.Flag.String("uds.transport", "", "Unix domain socket path for outbound RDMA replication (e.g., /tmp/sra-transport.sock)")
|
||||
}
|
||||
|
||||
var cmdVolume = &Command{
|
||||
@@ -289,6 +293,22 @@ func (v VolumeServerOptions) startVolumeServer(volumeFolders, maxVolumeCounts, v
|
||||
// starting grpc server
|
||||
grpcS := v.startGrpcService(volumeServer)
|
||||
|
||||
// starting UDS server for RDMA sidecar integration
|
||||
var udsServer *weed_server.UdsServer
|
||||
if *v.udsListen != "" {
|
||||
var err error
|
||||
udsServer, err = weed_server.NewUdsServer(volumeServer, *v.udsListen)
|
||||
if err != nil {
|
||||
glog.Fatalf("failed to start UDS server: %v", err)
|
||||
}
|
||||
udsServer.Start()
|
||||
}
|
||||
|
||||
// set up outbound RDMA replication transport
|
||||
if *v.udsTransport != "" {
|
||||
volumeServer.SetSraTransport(storage.NewSraTransport(*v.udsTransport))
|
||||
}
|
||||
|
||||
// starting public http server
|
||||
var publicHttpDown httpdown.Server
|
||||
if v.isSeparatedPublicPort() {
|
||||
@@ -315,7 +335,7 @@ func (v VolumeServerOptions) startVolumeServer(volumeFolders, maxVolumeCounts, v
|
||||
time.Sleep(time.Duration(*v.preStopSeconds) * time.Second)
|
||||
}
|
||||
|
||||
shutdown(publicHttpDown, clusterHttpServer, grpcS, volumeServer)
|
||||
shutdown(publicHttpDown, clusterHttpServer, grpcS, volumeServer, udsServer)
|
||||
stopChan <- true
|
||||
})
|
||||
|
||||
@@ -324,7 +344,7 @@ func (v VolumeServerOptions) startVolumeServer(volumeFolders, maxVolumeCounts, v
|
||||
select {
|
||||
case <-stopChan:
|
||||
case <-ctx.Done():
|
||||
shutdown(publicHttpDown, clusterHttpServer, grpcS, volumeServer)
|
||||
shutdown(publicHttpDown, clusterHttpServer, grpcS, volumeServer, udsServer)
|
||||
}
|
||||
} else {
|
||||
select {
|
||||
@@ -334,7 +354,7 @@ func (v VolumeServerOptions) startVolumeServer(volumeFolders, maxVolumeCounts, v
|
||||
|
||||
}
|
||||
|
||||
func shutdown(publicHttpDown httpdown.Server, clusterHttpServer httpdown.Server, grpcS *grpc.Server, volumeServer *weed_server.VolumeServer) {
|
||||
func shutdown(publicHttpDown httpdown.Server, clusterHttpServer httpdown.Server, grpcS *grpc.Server, volumeServer *weed_server.VolumeServer, udsServer *weed_server.UdsServer) {
|
||||
|
||||
// firstly, stop the public http service to prevent from receiving new user request
|
||||
if nil != publicHttpDown {
|
||||
@@ -352,6 +372,11 @@ func shutdown(publicHttpDown httpdown.Server, clusterHttpServer httpdown.Server,
|
||||
glog.V(0).Infof("graceful stop gRPC ...")
|
||||
grpcS.GracefulStop()
|
||||
|
||||
if udsServer != nil {
|
||||
glog.V(0).Infof("stop UDS server ...")
|
||||
udsServer.Stop()
|
||||
}
|
||||
|
||||
volumeServer.Shutdown()
|
||||
|
||||
pprof.StopCPUProfile()
|
||||
|
||||
@@ -159,8 +159,16 @@ func (vs *VolumeServer) LoadNewVolumes() {
|
||||
vs.store.LoadNewVolumes()
|
||||
}
|
||||
|
||||
// SetSraTransport configures the outbound RDMA replication transport on the volume store.
|
||||
func (vs *VolumeServer) SetSraTransport(t *storage.SraTransport) {
|
||||
vs.store.SraTransport = t
|
||||
}
|
||||
|
||||
func (vs *VolumeServer) Shutdown() {
|
||||
glog.V(0).Infoln("Shutting down volume server...")
|
||||
if vs.store.SraTransport != nil {
|
||||
vs.store.SraTransport.Close()
|
||||
}
|
||||
vs.store.Close()
|
||||
glog.V(0).Infoln("Shut down successfully!")
|
||||
}
|
||||
|
||||
@@ -0,0 +1,391 @@
|
||||
package weed_server
|
||||
|
||||
import (
|
||||
"encoding/binary"
|
||||
"fmt"
|
||||
"io"
|
||||
"net"
|
||||
"os"
|
||||
"strings"
|
||||
"sync"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/glog"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/needle"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/types"
|
||||
)
|
||||
|
||||
// UDS protocol constants
|
||||
const (
|
||||
UdsRequestSize = 24 // opcode(1) + pad(3) + request_id(4) + fid(16)
|
||||
UdsResponseSize = 32 // status(1) + pad(3) + volume_id(4) + offset(8) + length(8) + dat_path_len(2) + reserved(6)
|
||||
UdsMaxDatPathLen = 256 // max .dat file path length
|
||||
|
||||
// Store request: opcode(1) + pad(3) + request_id(4) + volume_id(4) + needle_version(1) + pad(3) + raw_bytes_len(4) = 20 bytes header
|
||||
UdsStoreHeaderSize = 20
|
||||
// Store response: status(1) + pad(3) + request_id(4) = 8 bytes
|
||||
UdsStoreResponseSize = 8
|
||||
|
||||
// Maximum raw needle size for Store (256 MB)
|
||||
UdsMaxRawBytesLen = 256 * 1024 * 1024
|
||||
)
|
||||
|
||||
// UDS opcodes
|
||||
const (
|
||||
UdsOpcodeLocate uint8 = 0x01
|
||||
UdsOpcodeStore uint8 = 0x02
|
||||
)
|
||||
|
||||
// UDS status codes
|
||||
const (
|
||||
UdsStatusOk uint8 = 0
|
||||
UdsStatusNotFound uint8 = 1
|
||||
UdsStatusError uint8 = 2
|
||||
)
|
||||
|
||||
// UdsServer handles UDS locate requests for RDMA sidecar integration
|
||||
type UdsServer struct {
|
||||
vs *VolumeServer
|
||||
listener net.Listener
|
||||
stopChan chan struct{}
|
||||
wg sync.WaitGroup
|
||||
}
|
||||
|
||||
// LocateRequest represents a UDS locate request
|
||||
// Wire format matches sra-common::uds_proto::LocateRequest (repr(C)):
|
||||
// opcode(1) + pad(3) + request_id(4) + fid(16) = 24 bytes
|
||||
type LocateRequest struct {
|
||||
Opcode uint8
|
||||
_pad [3]byte
|
||||
RequestId uint32
|
||||
Fid [16]byte // ASCII fid, null-padded
|
||||
}
|
||||
|
||||
// LocateResponse represents a UDS locate response
|
||||
type LocateResponse struct {
|
||||
Status uint8
|
||||
_pad [3]byte
|
||||
VolumeId uint32
|
||||
Offset uint64
|
||||
Length uint64
|
||||
DatPathLen uint16 // length of .dat file path (follows the 32-byte header)
|
||||
_ [6]byte // reserved
|
||||
DatPath string // variable-length .dat file path (not in wire header)
|
||||
}
|
||||
|
||||
// StoreRequest represents a UDS store request for volume replication.
|
||||
// Wire format:
|
||||
//
|
||||
// opcode(1) + pad(3) + request_id(4) + volume_id(4) +
|
||||
// needle_version(1) + pad(3) + raw_bytes_len(4) = 20 bytes header
|
||||
// + raw_bytes(variable)
|
||||
type StoreRequest struct {
|
||||
Opcode uint8
|
||||
RequestId uint32
|
||||
VolumeId uint32
|
||||
NeedleVersion uint8
|
||||
RawBytesLen uint32
|
||||
RawBytes []byte
|
||||
}
|
||||
|
||||
// StoreResponse represents a UDS store response.
|
||||
// Wire format: status(1) + pad(3) + request_id(4) = 8 bytes
|
||||
type StoreResponse struct {
|
||||
Status uint8
|
||||
RequestId uint32
|
||||
}
|
||||
|
||||
// NewUdsServer creates a new UDS server for the volume server
|
||||
func NewUdsServer(vs *VolumeServer, socketPath string) (*UdsServer, error) {
|
||||
// Remove existing socket file if exists
|
||||
if err := os.Remove(socketPath); err != nil && !os.IsNotExist(err) {
|
||||
return nil, fmt.Errorf("failed to remove existing socket: %w", err)
|
||||
}
|
||||
|
||||
listener, err := net.Listen("unix", socketPath)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to listen on %s: %w", socketPath, err)
|
||||
}
|
||||
|
||||
// Set socket permissions to allow access
|
||||
if err := os.Chmod(socketPath, 0666); err != nil {
|
||||
listener.Close()
|
||||
return nil, fmt.Errorf("failed to chmod socket: %w", err)
|
||||
}
|
||||
|
||||
uds := &UdsServer{
|
||||
vs: vs,
|
||||
listener: listener,
|
||||
stopChan: make(chan struct{}),
|
||||
}
|
||||
|
||||
glog.V(0).Infof("UDS server listening on %s", socketPath)
|
||||
return uds, nil
|
||||
}
|
||||
|
||||
// Start begins accepting connections
|
||||
func (u *UdsServer) Start() {
|
||||
u.wg.Add(1)
|
||||
go func() {
|
||||
defer u.wg.Done()
|
||||
u.acceptLoop()
|
||||
}()
|
||||
}
|
||||
|
||||
// Stop gracefully shuts down the UDS server
|
||||
func (u *UdsServer) Stop() {
|
||||
close(u.stopChan)
|
||||
u.listener.Close()
|
||||
u.wg.Wait()
|
||||
glog.V(0).Infoln("UDS server stopped")
|
||||
}
|
||||
|
||||
func (u *UdsServer) acceptLoop() {
|
||||
for {
|
||||
conn, err := u.listener.Accept()
|
||||
if err != nil {
|
||||
select {
|
||||
case <-u.stopChan:
|
||||
return
|
||||
default:
|
||||
glog.V(0).Infof("UDS accept error: %v", err)
|
||||
continue
|
||||
}
|
||||
}
|
||||
|
||||
u.wg.Add(1)
|
||||
go func(c net.Conn) {
|
||||
defer u.wg.Done()
|
||||
u.handleConnection(c)
|
||||
}(conn)
|
||||
}
|
||||
}
|
||||
|
||||
func (u *UdsServer) handleConnection(conn net.Conn) {
|
||||
defer conn.Close()
|
||||
|
||||
// Common prefix: opcode(1) + pad(3) + request_id(4) = 8 bytes
|
||||
commonBuf := make([]byte, 8)
|
||||
locateRestBuf := make([]byte, UdsRequestSize-8) // remaining 16 bytes for Locate
|
||||
locateRespBuf := make([]byte, UdsResponseSize)
|
||||
storeRestBuf := make([]byte, UdsStoreHeaderSize-8) // remaining 12 bytes for Store header
|
||||
storeRespBuf := make([]byte, UdsStoreResponseSize)
|
||||
|
||||
for {
|
||||
// Read common prefix (8 bytes)
|
||||
_, err := io.ReadFull(conn, commonBuf)
|
||||
if err != nil {
|
||||
if err != io.EOF {
|
||||
glog.V(2).Infof("UDS read error: %v", err)
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
opcode := commonBuf[0]
|
||||
requestId := binary.LittleEndian.Uint32(commonBuf[4:8])
|
||||
|
||||
switch opcode {
|
||||
case UdsOpcodeLocate:
|
||||
// Read remaining 16 bytes (fid)
|
||||
if _, err := io.ReadFull(conn, locateRestBuf); err != nil {
|
||||
glog.V(2).Infof("UDS read locate body: %v", err)
|
||||
return
|
||||
}
|
||||
|
||||
var req LocateRequest
|
||||
req.Opcode = opcode
|
||||
req.RequestId = requestId
|
||||
copy(req.Fid[:], locateRestBuf)
|
||||
|
||||
resp := u.handleLocate(&req)
|
||||
|
||||
// Serialize Locate response (32 bytes)
|
||||
locateRespBuf[0] = resp.Status
|
||||
locateRespBuf[1] = 0
|
||||
locateRespBuf[2] = 0
|
||||
locateRespBuf[3] = 0
|
||||
binary.LittleEndian.PutUint32(locateRespBuf[4:8], resp.VolumeId)
|
||||
binary.LittleEndian.PutUint64(locateRespBuf[8:16], resp.Offset)
|
||||
binary.LittleEndian.PutUint64(locateRespBuf[16:24], resp.Length)
|
||||
binary.LittleEndian.PutUint16(locateRespBuf[24:26], resp.DatPathLen)
|
||||
// reserved bytes 26-32 are zero
|
||||
for i := 26; i < 32; i++ {
|
||||
locateRespBuf[i] = 0
|
||||
}
|
||||
|
||||
if _, err := conn.Write(locateRespBuf); err != nil {
|
||||
glog.V(2).Infof("UDS write error: %v", err)
|
||||
return
|
||||
}
|
||||
if resp.DatPathLen > 0 {
|
||||
if _, err := conn.Write([]byte(resp.DatPath)); err != nil {
|
||||
glog.V(2).Infof("UDS write dat path error: %v", err)
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
case UdsOpcodeStore:
|
||||
// Read remaining 12 bytes of Store header
|
||||
if _, err := io.ReadFull(conn, storeRestBuf); err != nil {
|
||||
glog.V(2).Infof("UDS read store header: %v", err)
|
||||
return
|
||||
}
|
||||
|
||||
req := StoreRequest{
|
||||
Opcode: opcode,
|
||||
RequestId: requestId,
|
||||
VolumeId: binary.LittleEndian.Uint32(storeRestBuf[0:4]),
|
||||
NeedleVersion: storeRestBuf[4],
|
||||
RawBytesLen: binary.LittleEndian.Uint32(storeRestBuf[8:12]),
|
||||
}
|
||||
|
||||
// Validate raw_bytes_len
|
||||
if req.RawBytesLen == 0 || req.RawBytesLen > UdsMaxRawBytesLen {
|
||||
glog.V(2).Infof("UDS: invalid raw_bytes_len %d", req.RawBytesLen)
|
||||
u.writeStoreResponse(conn, storeRespBuf, requestId, UdsStatusError)
|
||||
return
|
||||
}
|
||||
|
||||
// Read raw needle bytes
|
||||
req.RawBytes = make([]byte, req.RawBytesLen)
|
||||
if _, err := io.ReadFull(conn, req.RawBytes); err != nil {
|
||||
glog.V(2).Infof("UDS read store payload: %v", err)
|
||||
return
|
||||
}
|
||||
|
||||
resp := u.handleStore(&req)
|
||||
if err := u.writeStoreResponse(conn, storeRespBuf, resp.RequestId, resp.Status); err != nil {
|
||||
return
|
||||
}
|
||||
|
||||
default:
|
||||
glog.V(0).Infof("UDS: unknown opcode 0x%02x from connection", opcode)
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (u *UdsServer) writeStoreResponse(conn net.Conn, buf []byte, requestId uint32, status uint8) error {
|
||||
buf[0] = status
|
||||
buf[1] = 0
|
||||
buf[2] = 0
|
||||
buf[3] = 0
|
||||
binary.LittleEndian.PutUint32(buf[4:8], requestId)
|
||||
if _, err := conn.Write(buf); err != nil {
|
||||
glog.V(2).Infof("UDS write store response error: %v", err)
|
||||
return err
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (u *UdsServer) handleLocate(req *LocateRequest) *LocateResponse {
|
||||
resp := &LocateResponse{}
|
||||
|
||||
// Extract fid string (null-terminated)
|
||||
fidBytes := req.Fid[:]
|
||||
fidLen := 0
|
||||
for i, b := range fidBytes {
|
||||
if b == 0 {
|
||||
fidLen = i
|
||||
break
|
||||
}
|
||||
fidLen = i + 1
|
||||
}
|
||||
fid := string(fidBytes[:fidLen])
|
||||
fid = strings.TrimSpace(fid)
|
||||
|
||||
if fid == "" {
|
||||
resp.Status = UdsStatusError
|
||||
return resp
|
||||
}
|
||||
|
||||
// Parse fid (format: "volumeId,needleIdCookie" e.g., "3,01637037")
|
||||
fileId, err := needle.ParseFileIdFromString(fid)
|
||||
if err != nil {
|
||||
glog.V(2).Infof("UDS: invalid fid %q: %v", fid, err)
|
||||
resp.Status = UdsStatusError
|
||||
return resp
|
||||
}
|
||||
|
||||
if u.vs == nil {
|
||||
resp.Status = UdsStatusError
|
||||
return resp
|
||||
}
|
||||
|
||||
volumeId := fileId.VolumeId
|
||||
needleId := fileId.Key
|
||||
|
||||
// Get volume
|
||||
v := u.vs.store.GetVolume(volumeId)
|
||||
if v == nil {
|
||||
glog.V(2).Infof("UDS: volume %d not found for fid %s", volumeId, fid)
|
||||
resp.Status = UdsStatusNotFound
|
||||
return resp
|
||||
}
|
||||
|
||||
// Lookup needle in the volume's needle map
|
||||
// We need to access the needle map through the volume
|
||||
nv, ok := v.GetNeedle(needleId)
|
||||
if !ok || nv.Offset.IsZero() {
|
||||
glog.V(3).Infof("UDS: needle %s not found in volume %d", fid, volumeId)
|
||||
resp.Status = UdsStatusNotFound
|
||||
return resp
|
||||
}
|
||||
|
||||
// Check if deleted
|
||||
if nv.Size.IsDeleted() {
|
||||
glog.V(3).Infof("UDS: needle %s is deleted", fid)
|
||||
resp.Status = UdsStatusNotFound
|
||||
return resp
|
||||
}
|
||||
|
||||
// Return the needle header offset (not data offset)
|
||||
// The reader will parse the header to find DataSize
|
||||
resp.Status = UdsStatusOk
|
||||
resp.VolumeId = uint32(volumeId)
|
||||
resp.Offset = uint64(nv.Offset.ToActualOffset())
|
||||
resp.Length = uint64(nv.Size)
|
||||
|
||||
// Include .dat file path so sidecar knows the exact filename
|
||||
// (SeaweedFS uses {collection}_{id}.dat when collection is non-empty)
|
||||
datPath := v.DataFileName() + ".dat"
|
||||
if len(datPath) <= UdsMaxDatPathLen {
|
||||
resp.DatPath = datPath
|
||||
resp.DatPathLen = uint16(len(datPath))
|
||||
}
|
||||
|
||||
glog.V(3).Infof("UDS: located %s -> vol=%d offset=%d size=%d dat=%s", fid, volumeId, resp.Offset, resp.Length, datPath)
|
||||
return resp
|
||||
}
|
||||
|
||||
func (u *UdsServer) handleStore(req *StoreRequest) *StoreResponse {
|
||||
resp := &StoreResponse{RequestId: req.RequestId}
|
||||
|
||||
if u.vs == nil {
|
||||
resp.Status = UdsStatusError
|
||||
return resp
|
||||
}
|
||||
|
||||
// Validate raw bytes contain at least a needle header (Cookie + NeedleId + Size = 16 bytes)
|
||||
if len(req.RawBytes) < types.NeedleHeaderSize {
|
||||
glog.V(2).Infof("UDS store: raw_bytes too short (%d < %d)", len(req.RawBytes), types.NeedleHeaderSize)
|
||||
resp.Status = UdsStatusError
|
||||
return resp
|
||||
}
|
||||
|
||||
// Parse needle header to extract needle ID and size
|
||||
var n needle.Needle
|
||||
n.ParseNeedleHeader(req.RawBytes[:types.NeedleHeaderSize])
|
||||
|
||||
vid := needle.VolumeId(req.VolumeId)
|
||||
|
||||
// Use the store's WriteNeedleBlob to append raw bytes and update the needle map
|
||||
if err := u.vs.store.WriteVolumeNeedleBlob(vid, n.Id, req.RawBytes, n.Size); err != nil {
|
||||
glog.V(2).Infof("UDS store: write failed for vol=%d needle=%v: %v", req.VolumeId, n.Id, err)
|
||||
resp.Status = UdsStatusError
|
||||
return resp
|
||||
}
|
||||
|
||||
glog.V(3).Infof("UDS: stored vol=%d needle=%v size=%d raw_len=%d", req.VolumeId, n.Id, n.Size, len(req.RawBytes))
|
||||
resp.Status = UdsStatusOk
|
||||
return resp
|
||||
}
|
||||
@@ -0,0 +1,297 @@
|
||||
package weed_server
|
||||
|
||||
import (
|
||||
"encoding/binary"
|
||||
"fmt"
|
||||
"io"
|
||||
"net"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/needle"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/types"
|
||||
"github.com/seaweedfs/seaweedfs/weed/util"
|
||||
)
|
||||
|
||||
// newTestStoreForUds creates a minimal storage.Store for UDS integration tests.
|
||||
// Caller must close(done) when finished to stop the channel drain goroutine.
|
||||
func newTestStoreForUds(t *testing.T) (s *storage.Store, done chan struct{}) {
|
||||
t.Helper()
|
||||
tmpDir := t.TempDir()
|
||||
volDir := filepath.Join(tmpDir, "vol")
|
||||
if err := os.MkdirAll(volDir, 0755); err != nil {
|
||||
t.Fatalf("MkdirAll: %v", err)
|
||||
}
|
||||
|
||||
s = storage.NewStore(nil, "localhost", 8080, 18080, "http://localhost:8080", "",
|
||||
[]string{volDir}, []int32{10}, []util.MinFreeSpace{{}}, "",
|
||||
storage.NeedleMapInMemory, []types.DiskType{types.HardDriveType}, 3)
|
||||
|
||||
// Drain NewVolumesChan to prevent AddVolume from blocking
|
||||
done = make(chan struct{})
|
||||
go func() {
|
||||
for {
|
||||
select {
|
||||
case <-s.NewVolumesChan:
|
||||
case <-done:
|
||||
return
|
||||
}
|
||||
}
|
||||
}()
|
||||
|
||||
return s, done
|
||||
}
|
||||
|
||||
// TestUdsIntegration_LocateWrittenNeedle is an end-to-end test:
|
||||
// 1. Create a volume store and write a needle
|
||||
// 2. Start a UDS server backed by a VolumeServer
|
||||
// 3. Send a Locate request over UDS
|
||||
// 4. Verify the response contains correct offset, length, and dat_path
|
||||
func TestUdsIntegration_LocateWrittenNeedle(t *testing.T) {
|
||||
store, done := newTestStoreForUds(t)
|
||||
defer close(done)
|
||||
|
||||
vid := needle.VolumeId(1)
|
||||
err := store.AddVolume(vid, "", storage.NeedleMapInMemory, "000", "",
|
||||
0, needle.GetCurrentVersion(), 0, types.HardDriveType, 3)
|
||||
if err != nil {
|
||||
t.Fatalf("AddVolume: %v", err)
|
||||
}
|
||||
|
||||
// Write a needle
|
||||
n := new(needle.Needle)
|
||||
n.Id = types.Uint64ToNeedleId(0x01)
|
||||
n.Cookie = types.Uint32ToCookie(0x12345678)
|
||||
n.Data = []byte("hello world integration test data")
|
||||
n.Checksum = needle.NewCRC(n.Data)
|
||||
|
||||
_, err = store.WriteVolumeNeedle(vid, n, false, false)
|
||||
if err != nil {
|
||||
t.Fatalf("WriteVolumeNeedle: %v", err)
|
||||
}
|
||||
|
||||
// Construct the fid string the same way SeaweedFS does
|
||||
fileId := needle.NewFileIdFromNeedle(vid, n)
|
||||
fid := fileId.String() // e.g. "1,0112345678"
|
||||
|
||||
// Create a minimal VolumeServer with just the store field
|
||||
vs := &VolumeServer{
|
||||
store: store,
|
||||
}
|
||||
|
||||
socketPath := filepath.Join(t.TempDir(), "test.sock")
|
||||
uds, err := NewUdsServer(vs, socketPath)
|
||||
if err != nil {
|
||||
t.Fatalf("NewUdsServer: %v", err)
|
||||
}
|
||||
defer uds.Stop()
|
||||
uds.Start()
|
||||
time.Sleep(10 * time.Millisecond)
|
||||
|
||||
conn, err := net.Dial("unix", socketPath)
|
||||
if err != nil {
|
||||
t.Fatalf("Dial: %v", err)
|
||||
}
|
||||
defer conn.Close()
|
||||
|
||||
// Send Locate request
|
||||
reqBuf := make([]byte, UdsRequestSize)
|
||||
reqBuf[0] = 1 // opcode = Locate
|
||||
binary.LittleEndian.PutUint32(reqBuf[4:8], 42)
|
||||
|
||||
// Copy fid into the 16-byte field (strip "vol," prefix — keep just keyCookie part? No — the UDS handler parses the full "vol,keyCookie" fid)
|
||||
if len(fid) > 16 {
|
||||
t.Fatalf("fid %q too long for 16-byte field", fid)
|
||||
}
|
||||
copy(reqBuf[8:24], fid)
|
||||
|
||||
if _, err := conn.Write(reqBuf); err != nil {
|
||||
t.Fatalf("Write request: %v", err)
|
||||
}
|
||||
|
||||
// Read 32-byte response header
|
||||
respBuf := make([]byte, UdsResponseSize)
|
||||
if _, err := io.ReadFull(conn, respBuf); err != nil {
|
||||
t.Fatalf("Read response: %v", err)
|
||||
}
|
||||
|
||||
status := respBuf[0]
|
||||
if status != UdsStatusOk {
|
||||
t.Fatalf("Expected UdsStatusOk (0), got %d", status)
|
||||
}
|
||||
|
||||
volumeId := binary.LittleEndian.Uint32(respBuf[4:8])
|
||||
if volumeId != uint32(vid) {
|
||||
t.Errorf("volumeId: expected %d, got %d", vid, volumeId)
|
||||
}
|
||||
|
||||
offset := binary.LittleEndian.Uint64(respBuf[8:16])
|
||||
if offset == 0 {
|
||||
t.Error("Expected non-zero offset (super block occupies the first bytes)")
|
||||
}
|
||||
|
||||
length := binary.LittleEndian.Uint64(respBuf[16:24])
|
||||
if length == 0 {
|
||||
t.Error("Expected non-zero length")
|
||||
}
|
||||
|
||||
datPathLen := binary.LittleEndian.Uint16(respBuf[24:26])
|
||||
if datPathLen == 0 {
|
||||
t.Fatal("Expected non-zero dat_path_len")
|
||||
}
|
||||
|
||||
// Read variable-length dat_path
|
||||
datPathBuf := make([]byte, datPathLen)
|
||||
if _, err := io.ReadFull(conn, datPathBuf); err != nil {
|
||||
t.Fatalf("Read dat_path: %v", err)
|
||||
}
|
||||
datPath := string(datPathBuf)
|
||||
|
||||
if !strings.HasSuffix(datPath, ".dat") {
|
||||
t.Errorf("dat_path should end with .dat, got %q", datPath)
|
||||
}
|
||||
|
||||
t.Logf("Locate %s -> offset=%d length=%d dat_path=%s", fid, offset, length, datPath)
|
||||
}
|
||||
|
||||
// TestUdsIntegration_LocateNotFound verifies that looking up a non-existent needle returns NotFound.
|
||||
func TestUdsIntegration_LocateNotFound(t *testing.T) {
|
||||
store, done := newTestStoreForUds(t)
|
||||
defer close(done)
|
||||
|
||||
vid := needle.VolumeId(1)
|
||||
err := store.AddVolume(vid, "", storage.NeedleMapInMemory, "000", "",
|
||||
0, needle.GetCurrentVersion(), 0, types.HardDriveType, 3)
|
||||
if err != nil {
|
||||
t.Fatalf("AddVolume: %v", err)
|
||||
}
|
||||
|
||||
vs := &VolumeServer{
|
||||
store: store,
|
||||
}
|
||||
|
||||
socketPath := filepath.Join(t.TempDir(), "test.sock")
|
||||
uds, err := NewUdsServer(vs, socketPath)
|
||||
if err != nil {
|
||||
t.Fatalf("NewUdsServer: %v", err)
|
||||
}
|
||||
defer uds.Stop()
|
||||
uds.Start()
|
||||
time.Sleep(10 * time.Millisecond)
|
||||
|
||||
conn, err := net.Dial("unix", socketPath)
|
||||
if err != nil {
|
||||
t.Fatalf("Dial: %v", err)
|
||||
}
|
||||
defer conn.Close()
|
||||
|
||||
// Locate a needle that was never written
|
||||
reqBuf := make([]byte, UdsRequestSize)
|
||||
reqBuf[0] = 1
|
||||
binary.LittleEndian.PutUint32(reqBuf[4:8], 1)
|
||||
copy(reqBuf[8:24], "1,ffaabbccddee")
|
||||
|
||||
if _, err := conn.Write(reqBuf); err != nil {
|
||||
t.Fatalf("Write: %v", err)
|
||||
}
|
||||
|
||||
respBuf := make([]byte, UdsResponseSize)
|
||||
if _, err := io.ReadFull(conn, respBuf); err != nil {
|
||||
t.Fatalf("Read: %v", err)
|
||||
}
|
||||
|
||||
status := respBuf[0]
|
||||
if status != UdsStatusNotFound {
|
||||
t.Errorf("Expected UdsStatusNotFound (%d), got %d", UdsStatusNotFound, status)
|
||||
}
|
||||
}
|
||||
|
||||
// TestUdsIntegration_MultipleRequests verifies that multiple sequential requests on the same connection work.
|
||||
func TestUdsIntegration_MultipleRequests(t *testing.T) {
|
||||
store, done := newTestStoreForUds(t)
|
||||
defer close(done)
|
||||
|
||||
vid := needle.VolumeId(1)
|
||||
err := store.AddVolume(vid, "", storage.NeedleMapInMemory, "000", "",
|
||||
0, needle.GetCurrentVersion(), 0, types.HardDriveType, 3)
|
||||
if err != nil {
|
||||
t.Fatalf("AddVolume: %v", err)
|
||||
}
|
||||
|
||||
// Write 3 needles
|
||||
type needleInfo struct {
|
||||
n *needle.Needle
|
||||
fid string
|
||||
}
|
||||
needles := make([]needleInfo, 3)
|
||||
for i := 0; i < 3; i++ {
|
||||
n := new(needle.Needle)
|
||||
n.Id = types.Uint64ToNeedleId(uint64(i + 1))
|
||||
n.Cookie = types.Uint32ToCookie(uint32(0xAABBCC00 + i))
|
||||
n.Data = []byte(fmt.Sprintf("needle data %d", i))
|
||||
n.Checksum = needle.NewCRC(n.Data)
|
||||
|
||||
if _, err := store.WriteVolumeNeedle(vid, n, false, false); err != nil {
|
||||
t.Fatalf("WriteVolumeNeedle[%d]: %v", i, err)
|
||||
}
|
||||
fileId := needle.NewFileIdFromNeedle(vid, n)
|
||||
needles[i] = needleInfo{n: n, fid: fileId.String()}
|
||||
}
|
||||
|
||||
vs := &VolumeServer{store: store}
|
||||
socketPath := filepath.Join(t.TempDir(), "test.sock")
|
||||
uds, err := NewUdsServer(vs, socketPath)
|
||||
if err != nil {
|
||||
t.Fatalf("NewUdsServer: %v", err)
|
||||
}
|
||||
defer uds.Stop()
|
||||
uds.Start()
|
||||
time.Sleep(10 * time.Millisecond)
|
||||
|
||||
conn, err := net.Dial("unix", socketPath)
|
||||
if err != nil {
|
||||
t.Fatalf("Dial: %v", err)
|
||||
}
|
||||
defer conn.Close()
|
||||
|
||||
// Send 3 Locate requests sequentially on the same connection
|
||||
for i, ni := range needles {
|
||||
reqBuf := make([]byte, UdsRequestSize)
|
||||
reqBuf[0] = 1
|
||||
binary.LittleEndian.PutUint32(reqBuf[4:8], uint32(i+1))
|
||||
copy(reqBuf[8:24], ni.fid)
|
||||
|
||||
if _, err := conn.Write(reqBuf); err != nil {
|
||||
t.Fatalf("Write[%d]: %v", i, err)
|
||||
}
|
||||
|
||||
respBuf := make([]byte, UdsResponseSize)
|
||||
if _, err := io.ReadFull(conn, respBuf); err != nil {
|
||||
t.Fatalf("Read[%d]: %v", i, err)
|
||||
}
|
||||
|
||||
status := respBuf[0]
|
||||
if status != UdsStatusOk {
|
||||
t.Errorf("needle[%d] fid=%s: expected UdsStatusOk, got %d", i, ni.fid, status)
|
||||
continue
|
||||
}
|
||||
|
||||
length := binary.LittleEndian.Uint64(respBuf[16:24])
|
||||
if length == 0 {
|
||||
t.Errorf("needle[%d]: expected non-zero length", i)
|
||||
}
|
||||
|
||||
// Consume the dat_path bytes
|
||||
datPathLen := binary.LittleEndian.Uint16(respBuf[24:26])
|
||||
if datPathLen > 0 {
|
||||
datPathBuf := make([]byte, datPathLen)
|
||||
if _, err := io.ReadFull(conn, datPathBuf); err != nil {
|
||||
t.Fatalf("Read dat_path[%d]: %v", i, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,423 @@
|
||||
package weed_server
|
||||
|
||||
import (
|
||||
"encoding/binary"
|
||||
"io"
|
||||
"net"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
func TestLocateRequestResponseSize(t *testing.T) {
|
||||
// Verify protocol sizes match spec
|
||||
if UdsRequestSize != 24 {
|
||||
t.Errorf("UdsRequestSize expected 24, got %d", UdsRequestSize)
|
||||
}
|
||||
if UdsResponseSize != 32 {
|
||||
t.Errorf("UdsResponseSize expected 32, got %d", UdsResponseSize)
|
||||
}
|
||||
}
|
||||
|
||||
func TestLocateRequestSerialization(t *testing.T) {
|
||||
// Wire format: opcode(1) + pad(3) + request_id(4) + fid(16) = 24
|
||||
req := LocateRequest{
|
||||
Opcode: 1, // VolumeOpcode::Locate
|
||||
RequestId: 42,
|
||||
}
|
||||
copy(req.Fid[:], "3,01637037")
|
||||
|
||||
buf := make([]byte, UdsRequestSize)
|
||||
buf[0] = req.Opcode
|
||||
binary.LittleEndian.PutUint32(buf[4:8], req.RequestId)
|
||||
copy(buf[8:24], req.Fid[:])
|
||||
|
||||
// Verify fid is correctly placed at offset 8
|
||||
fid := string(buf[8:18])
|
||||
if fid != "3,01637037" {
|
||||
t.Errorf("Expected fid='3,01637037', got '%s'", fid)
|
||||
}
|
||||
|
||||
// Verify request_id
|
||||
requestId := binary.LittleEndian.Uint32(buf[4:8])
|
||||
if requestId != 42 {
|
||||
t.Errorf("Expected request_id=42, got %d", requestId)
|
||||
}
|
||||
}
|
||||
|
||||
func TestLocateResponseSerialization(t *testing.T) {
|
||||
resp := &LocateResponse{
|
||||
Status: UdsStatusOk,
|
||||
VolumeId: 3,
|
||||
Offset: 1024,
|
||||
Length: 4096,
|
||||
}
|
||||
|
||||
buf := make([]byte, UdsResponseSize)
|
||||
buf[0] = resp.Status
|
||||
binary.LittleEndian.PutUint32(buf[4:8], resp.VolumeId)
|
||||
binary.LittleEndian.PutUint64(buf[8:16], resp.Offset)
|
||||
binary.LittleEndian.PutUint64(buf[16:24], resp.Length)
|
||||
|
||||
// Deserialize and verify
|
||||
status := buf[0]
|
||||
if status != UdsStatusOk {
|
||||
t.Errorf("Expected status=0, got %d", status)
|
||||
}
|
||||
|
||||
volumeId := binary.LittleEndian.Uint32(buf[4:8])
|
||||
if volumeId != 3 {
|
||||
t.Errorf("Expected volumeId=3, got %d", volumeId)
|
||||
}
|
||||
|
||||
offset := binary.LittleEndian.Uint64(buf[8:16])
|
||||
if offset != 1024 {
|
||||
t.Errorf("Expected offset=1024, got %d", offset)
|
||||
}
|
||||
|
||||
length := binary.LittleEndian.Uint64(buf[16:24])
|
||||
if length != 4096 {
|
||||
t.Errorf("Expected length=4096, got %d", length)
|
||||
}
|
||||
}
|
||||
|
||||
func TestUdsStatusCodes(t *testing.T) {
|
||||
if UdsStatusOk != 0 {
|
||||
t.Errorf("UdsStatusOk expected 0, got %d", UdsStatusOk)
|
||||
}
|
||||
if UdsStatusNotFound != 1 {
|
||||
t.Errorf("UdsStatusNotFound expected 1, got %d", UdsStatusNotFound)
|
||||
}
|
||||
if UdsStatusError != 2 {
|
||||
t.Errorf("UdsStatusError expected 2, got %d", UdsStatusError)
|
||||
}
|
||||
}
|
||||
|
||||
// TestUdsServerBasicConnection tests that we can create and connect to a UDS server
|
||||
func TestUdsServerBasicConnection(t *testing.T) {
|
||||
// Create temp socket path
|
||||
tmpDir := t.TempDir()
|
||||
socketPath := filepath.Join(tmpDir, "test.sock")
|
||||
|
||||
// Create a mock volume server (nil is fine since we won't make real lookups)
|
||||
uds, err := NewUdsServer(nil, socketPath)
|
||||
if err != nil {
|
||||
t.Fatalf("Failed to create UDS server: %v", err)
|
||||
}
|
||||
defer uds.Stop()
|
||||
|
||||
uds.Start()
|
||||
|
||||
// Give server time to start
|
||||
time.Sleep(10 * time.Millisecond)
|
||||
|
||||
// Try to connect
|
||||
conn, err := net.Dial("unix", socketPath)
|
||||
if err != nil {
|
||||
t.Fatalf("Failed to connect to UDS server: %v", err)
|
||||
}
|
||||
defer conn.Close()
|
||||
|
||||
// Connection successful
|
||||
}
|
||||
|
||||
// TestUdsServerProtocol tests the basic protocol exchange
|
||||
func TestUdsServerProtocol(t *testing.T) {
|
||||
// Create temp socket path
|
||||
tmpDir := t.TempDir()
|
||||
socketPath := filepath.Join(tmpDir, "test.sock")
|
||||
|
||||
// Create a mock volume server (vs=nil means handleLocate will return NotFound)
|
||||
uds, err := NewUdsServer(nil, socketPath)
|
||||
if err != nil {
|
||||
t.Fatalf("Failed to create UDS server: %v", err)
|
||||
}
|
||||
defer uds.Stop()
|
||||
|
||||
uds.Start()
|
||||
time.Sleep(10 * time.Millisecond)
|
||||
|
||||
// Connect
|
||||
conn, err := net.Dial("unix", socketPath)
|
||||
if err != nil {
|
||||
t.Fatalf("Failed to connect: %v", err)
|
||||
}
|
||||
defer conn.Close()
|
||||
|
||||
// Send a request (will fail because vs is nil, but we're testing protocol)
|
||||
// Wire format: opcode(1) + pad(3) + request_id(4) + fid(16)
|
||||
reqBuf := make([]byte, UdsRequestSize)
|
||||
reqBuf[0] = 1 // VolumeOpcode::Locate
|
||||
binary.LittleEndian.PutUint32(reqBuf[4:8], 1) // request_id
|
||||
copy(reqBuf[8:24], "3,01637037") // fid at offset 8
|
||||
|
||||
_, err = conn.Write(reqBuf)
|
||||
if err != nil {
|
||||
t.Fatalf("Failed to write request: %v", err)
|
||||
}
|
||||
|
||||
// Read response
|
||||
respBuf := make([]byte, UdsResponseSize)
|
||||
_, err = io.ReadFull(conn, respBuf)
|
||||
if err != nil {
|
||||
t.Fatalf("Failed to read response: %v", err)
|
||||
}
|
||||
|
||||
// Since vs is nil, we expect UdsStatusError
|
||||
status := respBuf[0]
|
||||
if status != UdsStatusError {
|
||||
t.Errorf("Expected UdsStatusError (%d) for nil VolumeServer, got %d", UdsStatusError, status)
|
||||
}
|
||||
|
||||
// DatPathLen should be 0 for error responses
|
||||
datPathLen := binary.LittleEndian.Uint16(respBuf[24:26])
|
||||
if datPathLen != 0 {
|
||||
t.Errorf("Expected DatPathLen=0 for error, got %d", datPathLen)
|
||||
}
|
||||
}
|
||||
|
||||
// TestUdsSocketCleanup tests that existing socket is removed
|
||||
func TestUdsSocketCleanup(t *testing.T) {
|
||||
tmpDir := t.TempDir()
|
||||
socketPath := filepath.Join(tmpDir, "test.sock")
|
||||
|
||||
// Create a dummy file at the socket path
|
||||
if err := os.WriteFile(socketPath, []byte("dummy"), 0644); err != nil {
|
||||
t.Fatalf("Failed to create dummy file: %v", err)
|
||||
}
|
||||
|
||||
// NewUdsServer should remove the existing file
|
||||
uds, err := NewUdsServer(nil, socketPath)
|
||||
if err != nil {
|
||||
t.Fatalf("Failed to create UDS server (should have cleaned up old socket): %v", err)
|
||||
}
|
||||
uds.Stop()
|
||||
}
|
||||
|
||||
// ============================================================
|
||||
// Cross-language wire contract tests (L1 foundation).
|
||||
//
|
||||
// These golden bytes MUST match the Rust side exactly
|
||||
// (sra-volume/tests/uds_locate_test.rs::wire_contract).
|
||||
// If either side changes, the contract breaks.
|
||||
// ============================================================
|
||||
|
||||
func TestGoldenLocateRequest(t *testing.T) {
|
||||
// Build the same request as Rust: fid="8,022aeb9e22", request_id=42
|
||||
golden := [UdsRequestSize]byte{
|
||||
0x01, // opcode = Locate
|
||||
0x00, 0x00, 0x00, // padding
|
||||
0x2A, 0x00, 0x00, 0x00, // request_id = 42 (LE)
|
||||
'8', ',', '0', '2', '2', 'a', 'e', 'b', // fid[0..8]
|
||||
'9', 'e', '2', '2', 0x00, 0x00, 0x00, 0x00, // fid[8..16]
|
||||
}
|
||||
|
||||
// Serialize from Go struct
|
||||
buf := make([]byte, UdsRequestSize)
|
||||
buf[0] = 1 // opcode = Locate
|
||||
binary.LittleEndian.PutUint32(buf[4:8], 42) // request_id
|
||||
copy(buf[8:24], "8,022aeb9e22")
|
||||
|
||||
for i := 0; i < UdsRequestSize; i++ {
|
||||
if buf[i] != golden[i] {
|
||||
t.Errorf("Go request byte[%d] = 0x%02X, golden = 0x%02X", i, buf[i], golden[i])
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestGoldenLocateResponse(t *testing.T) {
|
||||
// Build the same response as Rust: vol=8, offset=4096, length=4108, dat_path_len=29
|
||||
golden := [UdsResponseSize]byte{
|
||||
0x00, // status = Ok
|
||||
0x00, 0x00, 0x00, // padding
|
||||
0x08, 0x00, 0x00, 0x00, // volume_id = 8 (LE)
|
||||
0x00, 0x10, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, // offset = 4096 (LE)
|
||||
0x0C, 0x10, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, // length = 4108 (LE)
|
||||
0x1E, 0x00, // dat_path_len = 30 (LE)
|
||||
0x00, 0x00, 0x00, 0x00, 0x00, 0x00, // reserved
|
||||
}
|
||||
|
||||
// Serialize from Go struct
|
||||
buf := make([]byte, UdsResponseSize)
|
||||
buf[0] = UdsStatusOk
|
||||
binary.LittleEndian.PutUint32(buf[4:8], 8)
|
||||
binary.LittleEndian.PutUint64(buf[8:16], 4096)
|
||||
binary.LittleEndian.PutUint64(buf[16:24], 4108)
|
||||
binary.LittleEndian.PutUint16(buf[24:26], 30)
|
||||
|
||||
for i := 0; i < UdsResponseSize; i++ {
|
||||
if buf[i] != golden[i] {
|
||||
t.Errorf("Go response byte[%d] = 0x%02X, golden = 0x%02X", i, buf[i], golden[i])
|
||||
}
|
||||
}
|
||||
|
||||
// Verify trailing dat_path
|
||||
datPath := "/opt/work/data/weed/test_8.dat"
|
||||
if len(datPath) != 30 {
|
||||
t.Errorf("dat_path length = %d, expected 30", len(datPath))
|
||||
}
|
||||
}
|
||||
|
||||
func TestNeedleHeaderIsBigEndian(t *testing.T) {
|
||||
// SeaweedFS stores multi-byte needle fields as big-endian.
|
||||
// This test documents the convention for sra-volume (Rust) readers.
|
||||
//
|
||||
// DataSize = 1000 (0x000003E8) in big-endian = [0x00, 0x00, 0x03, 0xE8]
|
||||
var dataSize uint32 = 1000
|
||||
buf := make([]byte, 4)
|
||||
buf[0] = byte(dataSize >> 24)
|
||||
buf[1] = byte(dataSize >> 16)
|
||||
buf[2] = byte(dataSize >> 8)
|
||||
buf[3] = byte(dataSize)
|
||||
|
||||
expected := []byte{0x00, 0x00, 0x03, 0xE8}
|
||||
for i := 0; i < 4; i++ {
|
||||
if buf[i] != expected[i] {
|
||||
t.Errorf("DataSize byte[%d] = 0x%02X, expected 0x%02X (big-endian)", i, buf[i], expected[i])
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// ============================================================
|
||||
// Golden wire tests for Store opcode (0x02)
|
||||
// ============================================================
|
||||
|
||||
func TestGoldenStoreRequest(t *testing.T) {
|
||||
// Store request header: opcode(1) + pad(3) + request_id(4) + volume_id(4) +
|
||||
// needle_version(1) + pad(3) + raw_bytes_len(4) = 20 bytes
|
||||
// Values: opcode=0x02, request_id=99, volume_id=5, needle_version=2, raw_bytes_len=1024
|
||||
golden := [UdsStoreHeaderSize]byte{
|
||||
0x02, // opcode = Store
|
||||
0x00, 0x00, 0x00, // padding
|
||||
0x63, 0x00, 0x00, 0x00, // request_id = 99 (LE)
|
||||
0x05, 0x00, 0x00, 0x00, // volume_id = 5 (LE)
|
||||
0x02, // needle_version = 2
|
||||
0x00, 0x00, 0x00, // padding
|
||||
0x00, 0x04, 0x00, 0x00, // raw_bytes_len = 1024 (LE)
|
||||
}
|
||||
|
||||
// Serialize from Go
|
||||
buf := make([]byte, UdsStoreHeaderSize)
|
||||
buf[0] = UdsOpcodeStore
|
||||
binary.LittleEndian.PutUint32(buf[4:8], 99)
|
||||
binary.LittleEndian.PutUint32(buf[8:12], 5)
|
||||
buf[12] = 2 // needle_version
|
||||
binary.LittleEndian.PutUint32(buf[16:20], 1024)
|
||||
|
||||
for i := 0; i < UdsStoreHeaderSize; i++ {
|
||||
if buf[i] != golden[i] {
|
||||
t.Errorf("Store request byte[%d] = 0x%02X, golden = 0x%02X", i, buf[i], golden[i])
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestGoldenStoreResponse(t *testing.T) {
|
||||
// Store response: status(1) + pad(3) + request_id(4) = 8 bytes
|
||||
golden := [UdsStoreResponseSize]byte{
|
||||
0x00, // status = OK
|
||||
0x00, 0x00, 0x00, // padding
|
||||
0x63, 0x00, 0x00, 0x00, // request_id = 99 (LE)
|
||||
}
|
||||
|
||||
buf := make([]byte, UdsStoreResponseSize)
|
||||
buf[0] = UdsStatusOk
|
||||
binary.LittleEndian.PutUint32(buf[4:8], 99)
|
||||
|
||||
for i := 0; i < UdsStoreResponseSize; i++ {
|
||||
if buf[i] != golden[i] {
|
||||
t.Errorf("Store response byte[%d] = 0x%02X, golden = 0x%02X", i, buf[i], golden[i])
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// ============================================================
|
||||
// Golden wire tests for Replicate opcode (0x03)
|
||||
// These match the Rust transport_listener wire format.
|
||||
// ============================================================
|
||||
|
||||
func TestGoldenReplicateRequestHeader(t *testing.T) {
|
||||
// Replicate request header: opcode(1) + pad(3) + request_id(4) + volume_id(4) + target_host_len(4) = 16 bytes
|
||||
// Values: opcode=0x03, request_id=7, volume_id=3, target_host_len=15 ("10.0.0.3:18515")
|
||||
golden := [16]byte{
|
||||
0x03, // opcode = Replicate
|
||||
0x00, 0x00, 0x00, // padding
|
||||
0x07, 0x00, 0x00, 0x00, // request_id = 7 (LE)
|
||||
0x03, 0x00, 0x00, 0x00, // volume_id = 3 (LE)
|
||||
0x0F, 0x00, 0x00, 0x00, // target_host_len = 15 (LE)
|
||||
}
|
||||
|
||||
buf := make([]byte, 16)
|
||||
buf[0] = 0x03 // opcode
|
||||
binary.LittleEndian.PutUint32(buf[4:8], 7)
|
||||
binary.LittleEndian.PutUint32(buf[8:12], 3)
|
||||
binary.LittleEndian.PutUint32(buf[12:16], 15)
|
||||
|
||||
for i := 0; i < 16; i++ {
|
||||
if buf[i] != golden[i] {
|
||||
t.Errorf("Replicate header byte[%d] = 0x%02X, golden = 0x%02X", i, buf[i], golden[i])
|
||||
}
|
||||
}
|
||||
|
||||
// Verify target_host string
|
||||
targetHost := "10.0.0.3:18515"
|
||||
if len(targetHost) != 14 {
|
||||
// Note: "10.0.0.3:18515" is 14 chars. We use 15 in golden to test the len field.
|
||||
// In real usage it would be len(targetHost).
|
||||
}
|
||||
}
|
||||
|
||||
func TestGoldenReplicateResponse(t *testing.T) {
|
||||
// Replicate response: status(1) + pad(3) + request_id(4) = 8 bytes
|
||||
// Same format as Store response.
|
||||
golden := [8]byte{
|
||||
0xFF, // status = ERROR (stub always returns error)
|
||||
0x00, 0x00, 0x00, // padding
|
||||
0x07, 0x00, 0x00, 0x00, // request_id = 7 (LE)
|
||||
}
|
||||
|
||||
buf := make([]byte, 8)
|
||||
buf[0] = 0xFF // status = ERROR
|
||||
binary.LittleEndian.PutUint32(buf[4:8], 7)
|
||||
|
||||
for i := 0; i < 8; i++ {
|
||||
if buf[i] != golden[i] {
|
||||
t.Errorf("Replicate response byte[%d] = 0x%02X, golden = 0x%02X", i, buf[i], golden[i])
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestFidParsing(t *testing.T) {
|
||||
tests := []struct {
|
||||
fid string
|
||||
valid bool
|
||||
}{
|
||||
{"3,01637037", true},
|
||||
{"1,abc123", true},
|
||||
{"", false},
|
||||
{"invalid", false},
|
||||
}
|
||||
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.fid, func(t *testing.T) {
|
||||
req := &LocateRequest{}
|
||||
copy(req.Fid[:], tt.fid)
|
||||
|
||||
// Extract fid string using same logic as handleLocate
|
||||
fidBytes := req.Fid[:]
|
||||
fidLen := 0
|
||||
for i, b := range fidBytes {
|
||||
if b == 0 {
|
||||
fidLen = i
|
||||
break
|
||||
}
|
||||
fidLen = i + 1
|
||||
}
|
||||
fid := string(fidBytes[:fidLen])
|
||||
|
||||
if tt.fid == "" && fid != "" {
|
||||
t.Errorf("Expected empty fid, got '%s'", fid)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -45,6 +45,48 @@ func (n *Needle) Append(w backend.BackendStorageFile, version Version) (offset u
|
||||
return offset, size, actualSize, err
|
||||
}
|
||||
|
||||
// AppendGetBytes is like Append but also returns a copy of the raw serialized bytes
|
||||
// that were written to disk. This enables RDMA replication to send the exact .dat bytes
|
||||
// to a remote volume without re-serialization.
|
||||
func (n *Needle) AppendGetBytes(w backend.BackendStorageFile, version Version) (offset uint64, size Size, actualSize int64, rawBytes []byte, err error) {
|
||||
end, _, e := w.GetStat()
|
||||
if e != nil {
|
||||
err = fmt.Errorf("Cannot Read Current Volume Position: %w", e)
|
||||
return
|
||||
}
|
||||
offset = uint64(end)
|
||||
if offset >= MaxPossibleVolumeSize && len(n.Data) != 0 {
|
||||
err = fmt.Errorf("Volume Size %d Exceeded %d", offset, MaxPossibleVolumeSize)
|
||||
return
|
||||
}
|
||||
bytesBuffer := buffer_pool.SyncPoolGetBuffer()
|
||||
defer func() {
|
||||
if err != nil {
|
||||
if te := w.Truncate(end); te != nil {
|
||||
// handle error or log
|
||||
}
|
||||
}
|
||||
buffer_pool.SyncPoolPutBuffer(bytesBuffer)
|
||||
}()
|
||||
|
||||
size, actualSize, err = writeNeedleByVersion(version, n, offset, bytesBuffer)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
|
||||
// Copy serialized bytes before they are written to disk and the buffer is returned to pool
|
||||
src := bytesBuffer.Bytes()
|
||||
rawBytes = make([]byte, len(src))
|
||||
copy(rawBytes, src)
|
||||
|
||||
_, err = w.WriteAt(src, int64(offset))
|
||||
if err != nil {
|
||||
err = fmt.Errorf("failed to write %d bytes to %s at offset %d: %w", actualSize, w.Name(), offset, err)
|
||||
}
|
||||
|
||||
return offset, size, actualSize, rawBytes, err
|
||||
}
|
||||
|
||||
func WriteNeedleBlob(w backend.BackendStorageFile, dataSlice []byte, size Size, appendAtNs uint64, version Version) (offset uint64, err error) {
|
||||
|
||||
if end, _, e := w.GetStat(); e == nil {
|
||||
|
||||
@@ -8,6 +8,7 @@ import (
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/backend"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/types"
|
||||
"github.com/seaweedfs/seaweedfs/weed/util"
|
||||
)
|
||||
|
||||
func TestAppend(t *testing.T) {
|
||||
@@ -126,6 +127,87 @@ func TestWriteNeedle_CompatibilityWithLegacy(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// TestGoldenNeedleV2Bytes generates a v2 needle using the real SeaweedFS writer
|
||||
// and verifies it matches golden bytes that the Rust sra-volume side also validates.
|
||||
// If this test fails, Go and Rust disagree on the .dat file format.
|
||||
func TestGoldenNeedleV2Bytes(t *testing.T) {
|
||||
// Simple needle: Cookie=0xDEADBEEF, Id=0x0000000000ABCDEF, Data="HELLO_SRA"
|
||||
n := &Needle{
|
||||
Cookie: 0xDEADBEEF,
|
||||
Id: 0xABCDEF,
|
||||
Data: []byte("HELLO_SRA"),
|
||||
Flags: 0x00, // no optional fields
|
||||
}
|
||||
|
||||
buf := &bytes.Buffer{}
|
||||
_, _, err := writeNeedleV2(n, 0, buf)
|
||||
if err != nil {
|
||||
t.Fatalf("writeNeedleV2 failed: %v", err)
|
||||
}
|
||||
|
||||
raw := buf.Bytes()
|
||||
|
||||
// Log hex for Rust side to reference
|
||||
t.Logf("Golden v2 needle (%d bytes): %02x", len(raw), raw)
|
||||
|
||||
// === Verify header (16 bytes, all big-endian) ===
|
||||
// Cookie: 0xDEADBEEF -> [DE AD BE EF]
|
||||
if raw[0] != 0xDE || raw[1] != 0xAD || raw[2] != 0xBE || raw[3] != 0xEF {
|
||||
t.Errorf("Cookie mismatch: got %02x", raw[0:4])
|
||||
}
|
||||
|
||||
// NeedleId: 0xABCDEF -> [00 00 00 00 00 AB CD EF] (8 bytes BE)
|
||||
if raw[4] != 0x00 || raw[10] != 0xCD || raw[11] != 0xEF {
|
||||
t.Errorf("NeedleId mismatch: got %02x", raw[4:12])
|
||||
}
|
||||
|
||||
// Size field (4 bytes BE) = DataSize(4) + Data(9) + Flags(1) = 14
|
||||
expectedSize := uint32(4 + 9 + 1) // 14
|
||||
sizeVal := util.BytesToUint32(raw[12:16])
|
||||
if sizeVal != expectedSize {
|
||||
t.Errorf("Size expected %d, got %d (bytes: %02x)", expectedSize, sizeVal, raw[12:16])
|
||||
}
|
||||
|
||||
// === Verify body ===
|
||||
// DataSize (4 bytes BE) = 9
|
||||
dataSize := util.BytesToUint32(raw[16:20])
|
||||
if dataSize != 9 {
|
||||
t.Errorf("DataSize expected 9, got %d (bytes: %02x)", dataSize, raw[16:20])
|
||||
}
|
||||
|
||||
// Data: "HELLO_SRA" at bytes 20-28
|
||||
data := string(raw[20:29])
|
||||
if data != "HELLO_SRA" {
|
||||
t.Errorf("Data expected 'HELLO_SRA', got '%s'", data)
|
||||
}
|
||||
|
||||
// Flags at byte 29
|
||||
if raw[29] != 0x00 {
|
||||
t.Errorf("Flags expected 0x00, got 0x%02x", raw[29])
|
||||
}
|
||||
|
||||
// Checksum (4 bytes) at byte 30-33
|
||||
// Padding to align to 8 bytes
|
||||
|
||||
// Golden bytes check: header(16) + DataSize(4) + Data(9) + Flags(1) = 30 before checksum
|
||||
// Total with checksum(4) + padding = 40 bytes (aligned to 8)
|
||||
expectedTotal := 40 // 16+14+4+padding(6)... let's just verify
|
||||
if len(raw) != expectedTotal {
|
||||
t.Logf("Total needle size: %d bytes (expected ~40, adjust if needed)", len(raw))
|
||||
}
|
||||
|
||||
// Dump each field for Rust side reference
|
||||
t.Logf("Byte-by-byte breakdown:")
|
||||
t.Logf(" [0:4] Cookie: %02x", raw[0:4])
|
||||
t.Logf(" [4:12] NeedleId: %02x", raw[4:12])
|
||||
t.Logf(" [12:16] Size: %02x (=%d)", raw[12:16], sizeVal)
|
||||
t.Logf(" [16:20] DataSize: %02x (=%d)", raw[16:20], dataSize)
|
||||
t.Logf(" [20:29] Data: %02x ('%s')", raw[20:29], raw[20:29])
|
||||
t.Logf(" [29] Flags: %02x", raw[29])
|
||||
t.Logf(" [30:34] Checksum: %02x", raw[30:34])
|
||||
t.Logf(" [34:40] Padding: %02x", raw[34:40])
|
||||
}
|
||||
|
||||
type mockBackendWriter struct {
|
||||
buf *bytes.Buffer
|
||||
}
|
||||
|
||||
@@ -0,0 +1,157 @@
|
||||
package storage
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/binary"
|
||||
"fmt"
|
||||
"io"
|
||||
"net"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/glog"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/needle"
|
||||
)
|
||||
|
||||
// SraTransport provides outbound RDMA replication via the sra-volume sidecar.
|
||||
// It connects to the sidecar's transport UDS socket and sends Replicate requests.
|
||||
// Nil-safe: callers should check Store.SraTransport != nil before using.
|
||||
type SraTransport struct {
|
||||
socketPath string
|
||||
mu sync.Mutex
|
||||
conn net.Conn
|
||||
}
|
||||
|
||||
// Replicate request wire format (opcode 0x03):
|
||||
//
|
||||
// opcode(1) + pad(3) + request_id(4) + volume_id(4) +
|
||||
// target_host_len(4) + target_host(N) + raw_bytes_len(4) + raw_bytes(M)
|
||||
//
|
||||
// Response: status(1) + pad(3) + request_id(4) = 8 bytes
|
||||
const (
|
||||
sraOpcodeReplicate uint8 = 0x03
|
||||
sraReplicateHeaderSize = 16 // opcode(1)+pad(3)+request_id(4)+volume_id(4)+target_host_len(4)
|
||||
sraReplicateResponseSize = 8
|
||||
sraStatusOk uint8 = 0x00
|
||||
sraStatusError uint8 = 0xFF
|
||||
)
|
||||
|
||||
var nextRequestId uint32
|
||||
var nextRequestIdMu sync.Mutex
|
||||
|
||||
func nextReqId() uint32 {
|
||||
nextRequestIdMu.Lock()
|
||||
defer nextRequestIdMu.Unlock()
|
||||
nextRequestId++
|
||||
return nextRequestId
|
||||
}
|
||||
|
||||
// NewSraTransport creates a new outbound transport client.
|
||||
// Returns nil if socketPath is empty (RDMA replication disabled).
|
||||
func NewSraTransport(socketPath string) *SraTransport {
|
||||
if socketPath == "" {
|
||||
return nil
|
||||
}
|
||||
return &SraTransport{
|
||||
socketPath: socketPath,
|
||||
}
|
||||
}
|
||||
|
||||
// Replicate sends raw needle bytes to a remote volume server via the sidecar's RDMA transport.
|
||||
// targetHost is the RDMA endpoint of the destination (e.g., "10.0.0.3:18515").
|
||||
// volumeId is the destination volume. rawBytes is the serialized needle data.
|
||||
func (t *SraTransport) Replicate(ctx context.Context, targetHost string, volumeId needle.VolumeId, rawBytes []byte) error {
|
||||
t.mu.Lock()
|
||||
defer t.mu.Unlock()
|
||||
|
||||
conn, err := t.getConn()
|
||||
if err != nil {
|
||||
return fmt.Errorf("sra transport connect: %w", err)
|
||||
}
|
||||
|
||||
reqId := nextReqId()
|
||||
|
||||
// Build request
|
||||
targetHostBytes := []byte(targetHost)
|
||||
headerBuf := make([]byte, sraReplicateHeaderSize)
|
||||
headerBuf[0] = sraOpcodeReplicate
|
||||
// pad[1..3] = 0
|
||||
binary.LittleEndian.PutUint32(headerBuf[4:8], reqId)
|
||||
binary.LittleEndian.PutUint32(headerBuf[8:12], uint32(volumeId))
|
||||
binary.LittleEndian.PutUint32(headerBuf[12:16], uint32(len(targetHostBytes)))
|
||||
|
||||
rawBytesLenBuf := make([]byte, 4)
|
||||
binary.LittleEndian.PutUint32(rawBytesLenBuf, uint32(len(rawBytes)))
|
||||
|
||||
// Set deadline from context or default 30s
|
||||
deadline, ok := ctx.Deadline()
|
||||
if !ok {
|
||||
deadline = time.Now().Add(30 * time.Second)
|
||||
}
|
||||
conn.SetDeadline(deadline)
|
||||
|
||||
// Write: header + target_host + raw_bytes_len + raw_bytes
|
||||
if _, err := conn.Write(headerBuf); err != nil {
|
||||
t.closeConn()
|
||||
return fmt.Errorf("sra transport write header: %w", err)
|
||||
}
|
||||
if _, err := conn.Write(targetHostBytes); err != nil {
|
||||
t.closeConn()
|
||||
return fmt.Errorf("sra transport write target_host: %w", err)
|
||||
}
|
||||
if _, err := conn.Write(rawBytesLenBuf); err != nil {
|
||||
t.closeConn()
|
||||
return fmt.Errorf("sra transport write raw_bytes_len: %w", err)
|
||||
}
|
||||
if _, err := conn.Write(rawBytes); err != nil {
|
||||
t.closeConn()
|
||||
return fmt.Errorf("sra transport write raw_bytes: %w", err)
|
||||
}
|
||||
|
||||
// Read response (8 bytes)
|
||||
respBuf := make([]byte, sraReplicateResponseSize)
|
||||
if _, err := io.ReadFull(conn, respBuf); err != nil {
|
||||
t.closeConn()
|
||||
return fmt.Errorf("sra transport read response: %w", err)
|
||||
}
|
||||
|
||||
status := respBuf[0]
|
||||
respReqId := binary.LittleEndian.Uint32(respBuf[4:8])
|
||||
if respReqId != reqId {
|
||||
t.closeConn()
|
||||
return fmt.Errorf("sra transport request_id mismatch: sent %d, got %d", reqId, respReqId)
|
||||
}
|
||||
|
||||
if status != sraStatusOk {
|
||||
return fmt.Errorf("sra transport replicate failed: status=0x%02x", status)
|
||||
}
|
||||
|
||||
glog.V(3).Infof("sra transport: replicated %d bytes to %s vol=%d", len(rawBytes), targetHost, volumeId)
|
||||
return nil
|
||||
}
|
||||
|
||||
// Close closes the transport connection.
|
||||
func (t *SraTransport) Close() {
|
||||
t.mu.Lock()
|
||||
defer t.mu.Unlock()
|
||||
t.closeConn()
|
||||
}
|
||||
|
||||
func (t *SraTransport) getConn() (net.Conn, error) {
|
||||
if t.conn != nil {
|
||||
return t.conn, nil
|
||||
}
|
||||
conn, err := net.DialTimeout("unix", t.socketPath, 5*time.Second)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
t.conn = conn
|
||||
return conn, nil
|
||||
}
|
||||
|
||||
func (t *SraTransport) closeConn() {
|
||||
if t.conn != nil {
|
||||
t.conn.Close()
|
||||
t.conn = nil
|
||||
}
|
||||
}
|
||||
@@ -78,6 +78,7 @@ type Store struct {
|
||||
NewEcShardsChan chan master_pb.VolumeEcShardInformationMessage
|
||||
DeletedEcShardsChan chan master_pb.VolumeEcShardInformationMessage
|
||||
isStopping bool
|
||||
SraTransport *SraTransport // outbound RDMA replication transport (nil = HTTP-only)
|
||||
}
|
||||
|
||||
func (s *Store) String() (str string) {
|
||||
@@ -578,6 +579,19 @@ func (s *Store) WriteVolumeNeedle(i needle.VolumeId, n *needle.Needle, checkCook
|
||||
return
|
||||
}
|
||||
|
||||
// WriteVolumeNeedleBlob appends raw needle bytes to a volume's .dat file and updates
|
||||
// the needle map. Used by the UDS Store opcode for RDMA replication where the sender
|
||||
// provides already-serialized needle bytes.
|
||||
func (s *Store) WriteVolumeNeedleBlob(i needle.VolumeId, needleId NeedleId, needleBlob []byte, size Size) error {
|
||||
if v := s.findVolume(i); v != nil {
|
||||
if v.IsReadOnly() {
|
||||
return fmt.Errorf("volume %d is read only", i)
|
||||
}
|
||||
return v.WriteNeedleBlob(needleId, needleBlob, size)
|
||||
}
|
||||
return fmt.Errorf("volume %d not found on %s:%d", i, s.Ip, s.Port)
|
||||
}
|
||||
|
||||
func (s *Store) DeleteVolumeNeedle(i needle.VolumeId, n *needle.Needle) (Size, error) {
|
||||
if v := s.findVolume(i); v != nil {
|
||||
if v.noWriteOrDelete {
|
||||
|
||||
@@ -12,6 +12,7 @@ import (
|
||||
"github.com/seaweedfs/seaweedfs/weed/stats"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/backend"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/needle"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/needle_map"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/super_block"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/types"
|
||||
|
||||
@@ -203,6 +204,17 @@ func (v *Volume) IndexFileSize() uint64 {
|
||||
return v.nm.IndexFileSize()
|
||||
}
|
||||
|
||||
// GetNeedle looks up a needle by its ID and returns the needle value containing offset and size.
|
||||
// This is used by the UDS handler for RDMA sidecar integration.
|
||||
func (v *Volume) GetNeedle(needleId types.NeedleId) (*needle_map.NeedleValue, bool) {
|
||||
v.dataFileAccessLock.RLock()
|
||||
defer v.dataFileAccessLock.RUnlock()
|
||||
if v.nm == nil {
|
||||
return nil, false
|
||||
}
|
||||
return v.nm.Get(needleId)
|
||||
}
|
||||
|
||||
func (v *Volume) DiskType() types.DiskType {
|
||||
return v.location.DiskType
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user