mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-01 20:26:27 +00:00
Compare commits
4
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
3cf04eabef | ||
|
|
69df7a9182 | ||
|
|
b25f95df0a | ||
|
|
b7427c7d76 |
@@ -96,21 +96,12 @@ func (n *Needle) ReadBytes(bytes []byte, offset int64, size Size, version Versio
|
||||
}
|
||||
|
||||
// ReadData hydrates the needle from the file, with only n.Id is set.
|
||||
func (n *Needle) ReadData(r backend.BackendStorageFile, offset int64, size Size, version Version) (err error) {
|
||||
bytes, err := ReadNeedleBlob(r, offset, size, version)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
err = n.ReadBytes(bytes, offset, size, version)
|
||||
if err == ErrorSizeMismatch && OffsetSize == 4 {
|
||||
offset = offset + int64(MaxPossibleVolumeSize)
|
||||
bytes, err = ReadNeedleBlob(r, offset, size, version)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
err = n.ReadBytes(bytes, offset, size, version)
|
||||
}
|
||||
return err
|
||||
func (n *Needle) ReadData(r backend.BackendStorageFile, offset int64, size Size, version Version) error {
|
||||
return n.ReadFromFile(r, offset, size, version, NeedleReadOptions{
|
||||
ReadHeader: true,
|
||||
ReadData: true,
|
||||
ReadMeta: true,
|
||||
})
|
||||
}
|
||||
|
||||
func (n *Needle) ParseNeedleHeader(bytes []byte) {
|
||||
|
||||
@@ -0,0 +1,108 @@
|
||||
package needle
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"io"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/backend"
|
||||
. "github.com/seaweedfs/seaweedfs/weed/storage/types"
|
||||
"github.com/seaweedfs/seaweedfs/weed/util"
|
||||
)
|
||||
|
||||
// NeedleReadOptions specifies which parts of the Needle to read.
|
||||
type NeedleReadOptions struct {
|
||||
ReadHeader bool // always true for any read
|
||||
ReadData bool // read the Data field
|
||||
ReadMeta bool // read metadata fields (Name, Mime, LastModified, Ttl, Pairs, etc.)
|
||||
}
|
||||
|
||||
// ReadFromFile reads the Needle from the backend file according to the specified options.
|
||||
// - If only ReadHeader is true, only the header is read and parsed.
|
||||
// - If ReadData or ReadMeta is true, reads GetActualSize(size, version) bytes from disk (size is the logical body size).
|
||||
func (n *Needle) ReadFromFile(r backend.BackendStorageFile, offset int64, size Size, version Version, opts NeedleReadOptions) (err error) {
|
||||
defer func() {
|
||||
if r := recover(); r != nil {
|
||||
err = fmt.Errorf("panic occurred: %+v", r)
|
||||
}
|
||||
}()
|
||||
|
||||
if opts.ReadHeader && !opts.ReadData && !opts.ReadMeta {
|
||||
// Only read the header
|
||||
header := make([]byte, NeedleHeaderSize)
|
||||
count, err := r.ReadAt(header, offset)
|
||||
if err == io.EOF && count == NeedleHeaderSize {
|
||||
err = nil
|
||||
}
|
||||
if count != NeedleHeaderSize || err != nil {
|
||||
return err
|
||||
}
|
||||
n.ParseNeedleHeader(header)
|
||||
return nil
|
||||
}
|
||||
if opts.ReadHeader && opts.ReadMeta && !opts.ReadData {
|
||||
// Optimized: Read header and DataSize in one call
|
||||
buf := make([]byte, NeedleHeaderSize+DataSizeSize)
|
||||
count, err := r.ReadAt(buf, offset)
|
||||
if err == io.EOF && count == NeedleHeaderSize+DataSizeSize {
|
||||
err = nil
|
||||
}
|
||||
if count != NeedleHeaderSize+DataSizeSize || err != nil {
|
||||
return err
|
||||
}
|
||||
n.ParseNeedleHeader(buf[:NeedleHeaderSize])
|
||||
if n.Size != size {
|
||||
if OffsetSize == 4 && offset < int64(MaxPossibleVolumeSize) {
|
||||
return ErrorSizeMismatch
|
||||
}
|
||||
}
|
||||
|
||||
// Now read meta fields after DataSize+Data
|
||||
if version == Version2 || version == Version3 {
|
||||
n.DataSize = util.BytesToUint32(buf[NeedleHeaderSize : NeedleHeaderSize+DataSizeSize])
|
||||
|
||||
startOffset := offset + NeedleHeaderSize
|
||||
if size.IsValid() {
|
||||
startOffset = offset + NeedleHeaderSize + DataSizeSize + int64(n.DataSize)
|
||||
}
|
||||
dataSize := GetActualSize(size, version)
|
||||
stopOffset := offset + dataSize
|
||||
metaFieldsLen := stopOffset - startOffset
|
||||
|
||||
metaFieldsBuf := make([]byte, metaFieldsLen)
|
||||
count, err = r.ReadAt(metaFieldsBuf, startOffset)
|
||||
if err == io.EOF && int64(count) == metaFieldsLen {
|
||||
err = nil
|
||||
}
|
||||
if count <= 0 || err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
var index int
|
||||
if size.IsValid() {
|
||||
index, err = n.readNeedleDataVersion2NonData(metaFieldsBuf)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
n.Checksum = CRC(util.BytesToUint32(metaFieldsBuf[index : index+NeedleChecksumSize]))
|
||||
if version == Version3 {
|
||||
n.AppendAtNs = util.BytesToUint64(metaFieldsBuf[index+NeedleChecksumSize : index+NeedleChecksumSize+TimestampSize])
|
||||
}
|
||||
return nil
|
||||
}
|
||||
// For v1, just skip Data
|
||||
return nil
|
||||
}
|
||||
// Otherwise, read the full on-disk entry size
|
||||
readLen := int(GetActualSize(size, version))
|
||||
bytes := make([]byte, readLen)
|
||||
count, err := r.ReadAt(bytes, offset)
|
||||
if err == io.EOF && count == readLen {
|
||||
err = nil
|
||||
}
|
||||
if count != readLen || err != nil {
|
||||
return err
|
||||
}
|
||||
return n.ReadBytes(bytes, offset, size, version)
|
||||
}
|
||||
@@ -0,0 +1,174 @@
|
||||
package needle
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"io"
|
||||
"reflect"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
type mockBackend struct {
|
||||
data []byte
|
||||
}
|
||||
|
||||
func (m *mockBackend) ReadAt(p []byte, off int64) (n int, err error) {
|
||||
if int(off) >= len(m.data) {
|
||||
return 0, io.EOF
|
||||
}
|
||||
n = copy(p, m.data[off:])
|
||||
if n < len(p) {
|
||||
return n, io.EOF
|
||||
}
|
||||
return n, nil
|
||||
}
|
||||
|
||||
func (m *mockBackend) GetStat() (int64, time.Time, error) {
|
||||
return int64(len(m.data)), time.Time{}, nil
|
||||
}
|
||||
|
||||
func (m *mockBackend) Name() string {
|
||||
return "mock"
|
||||
}
|
||||
|
||||
func (m *mockBackend) Close() error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (m *mockBackend) Sync() error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (m *mockBackend) Truncate(size int64) error {
|
||||
m.data = m.data[:size]
|
||||
return nil
|
||||
}
|
||||
|
||||
func (m *mockBackend) WriteAt(p []byte, off int64) (n int, err error) {
|
||||
return 0, nil
|
||||
}
|
||||
|
||||
func TestReadFromFile_EquivalenceWithReadData(t *testing.T) {
|
||||
n := &Needle{
|
||||
Cookie: 0x12345678,
|
||||
Id: 0x1122334455667788,
|
||||
Data: []byte("hello world"),
|
||||
Flags: 0xFF,
|
||||
Name: []byte("filename.txt"),
|
||||
Mime: []byte("text/plain"),
|
||||
LastModified: 0x1234567890,
|
||||
Ttl: nil,
|
||||
Pairs: []byte("key=value"),
|
||||
PairsSize: 9,
|
||||
Checksum: 0,
|
||||
AppendAtNs: 0xDEADBEEF,
|
||||
}
|
||||
// remove the TTL bit in the flags
|
||||
n.Flags = n.Flags &^ FlagHasTtl
|
||||
// n.Checksum = NewCRC(n.Data)
|
||||
|
||||
buf := &bytes.Buffer{}
|
||||
_, _, err := writeNeedleV2(n, 0, buf)
|
||||
if err != nil {
|
||||
t.Fatalf("writeNeedleV2 failed: %v", err)
|
||||
}
|
||||
backend := &mockBackend{data: buf.Bytes()}
|
||||
size := n.Size
|
||||
|
||||
// Old method
|
||||
nOld := &Needle{}
|
||||
errOld := nOld.ReadData(backend, 0, size, Version2)
|
||||
|
||||
// New method
|
||||
nNew := &Needle{}
|
||||
opts := NeedleReadOptions{ReadHeader: true, ReadData: true, ReadMeta: true}
|
||||
errNew := nNew.ReadFromFile(backend, 0, size, Version2, opts)
|
||||
|
||||
if (errOld != nil) != (errNew != nil) || (errOld != nil && errOld.Error() != errNew.Error()) {
|
||||
t.Errorf("error mismatch: old=%v new=%v", errOld, errNew)
|
||||
}
|
||||
|
||||
if !reflect.DeepEqual(nOld, nNew) {
|
||||
t.Errorf("needle mismatch: old=%+v new=%+v", nOld, nNew)
|
||||
}
|
||||
}
|
||||
|
||||
func TestReadFromFile_OptionsMatrix(t *testing.T) {
|
||||
n := &Needle{
|
||||
Cookie: 0x12345678,
|
||||
Id: 0x1122334455667788,
|
||||
Data: []byte("hello world"),
|
||||
Flags: 0xFF,
|
||||
Name: []byte("filename.txt"),
|
||||
Mime: []byte("text/plain"),
|
||||
LastModified: 0x1234567890,
|
||||
Ttl: nil,
|
||||
Pairs: []byte("key=value"),
|
||||
PairsSize: 9,
|
||||
Checksum: 0,
|
||||
AppendAtNs: 0xDEADBEEF,
|
||||
}
|
||||
n.Flags = n.Flags &^ FlagHasTtl
|
||||
n.Checksum = NewCRC(n.Data)
|
||||
|
||||
buf := &bytes.Buffer{}
|
||||
_, _, err := writeNeedleV2(n, 0, buf)
|
||||
if err != nil {
|
||||
t.Fatalf("writeNeedleV2 failed: %v", err)
|
||||
}
|
||||
backend := &mockBackend{data: buf.Bytes()}
|
||||
size := n.Size
|
||||
|
||||
t.Run("ReadHeader only", func(t *testing.T) {
|
||||
nHeader := &Needle{}
|
||||
opts := NeedleReadOptions{ReadHeader: true}
|
||||
err := nHeader.ReadFromFile(backend, 0, size, Version2, opts)
|
||||
if err != nil {
|
||||
t.Fatalf("ReadFromFile header only failed: %v", err)
|
||||
}
|
||||
if nHeader.Cookie != n.Cookie || nHeader.Id != n.Id || nHeader.Size != n.Size {
|
||||
t.Errorf("header fields mismatch: got %+v want %+v", nHeader, n)
|
||||
}
|
||||
if nHeader.Data != nil || nHeader.Name != nil || nHeader.Mime != nil || nHeader.Pairs != nil || nHeader.Checksum != 0 {
|
||||
t.Errorf("non-header fields should be zero, got %+v", nHeader)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("ReadHeader+ReadMeta", func(t *testing.T) {
|
||||
nMeta := &Needle{}
|
||||
opts := NeedleReadOptions{ReadHeader: true, ReadMeta: true}
|
||||
err := nMeta.ReadFromFile(backend, 0, size, Version2, opts)
|
||||
if err != nil {
|
||||
t.Fatalf("ReadFromFile header+meta failed: %v", err)
|
||||
}
|
||||
if nMeta.Data != nil {
|
||||
t.Errorf("Data should not be set when only meta is read")
|
||||
}
|
||||
if nMeta.Name == nil || nMeta.Mime == nil || nMeta.Pairs == nil {
|
||||
t.Errorf("meta fields should be set, got %+v", nMeta)
|
||||
}
|
||||
if nMeta.Cookie != n.Cookie || nMeta.Id != n.Id || nMeta.Size != n.Size {
|
||||
t.Errorf("header fields mismatch: got %+v want %+v", nMeta, n)
|
||||
}
|
||||
if nMeta.Checksum != n.Checksum {
|
||||
t.Errorf("checksum mismatch: got %d want %d", nMeta.Checksum, n.Checksum)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("ReadHeader+ReadData", func(t *testing.T) {
|
||||
// this is the same as ReadHeader+ReadData+ReadMeta
|
||||
})
|
||||
|
||||
t.Run("ReadHeader+ReadData+ReadMeta", func(t *testing.T) {
|
||||
nFull := &Needle{}
|
||||
opts := NeedleReadOptions{ReadHeader: true, ReadData: true, ReadMeta: true}
|
||||
err := nFull.ReadFromFile(backend, 0, size, Version2, opts)
|
||||
if err != nil {
|
||||
t.Fatalf("ReadFromFile header+data+meta failed: %v", err)
|
||||
}
|
||||
nFull.AppendAtNs = n.AppendAtNs
|
||||
if !reflect.DeepEqual(nFull, n) {
|
||||
t.Errorf("needle mismatch: got %+v want %+v", nFull, n)
|
||||
}
|
||||
})
|
||||
}
|
||||
@@ -1,12 +1,11 @@
|
||||
package needle
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"io"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/glog"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/backend"
|
||||
. "github.com/seaweedfs/seaweedfs/weed/storage/types"
|
||||
"github.com/seaweedfs/seaweedfs/weed/util"
|
||||
"io"
|
||||
)
|
||||
|
||||
// ReadNeedleData uses a needle without n.Data to read the content
|
||||
@@ -33,57 +32,12 @@ func (n *Needle) ReadNeedleData(r backend.BackendStorageFile, volumeOffset int64
|
||||
}
|
||||
|
||||
// ReadNeedleMeta fills all metadata except the n.Data
|
||||
func (n *Needle) ReadNeedleMeta(r backend.BackendStorageFile, offset int64, size Size, version Version) (err error) {
|
||||
defer func() {
|
||||
if r := recover(); r != nil {
|
||||
err = fmt.Errorf("panic occurred: %+v", r)
|
||||
}
|
||||
}()
|
||||
|
||||
bytes := make([]byte, NeedleHeaderSize+DataSizeSize)
|
||||
|
||||
count, err := r.ReadAt(bytes, offset)
|
||||
if err == io.EOF && count == NeedleHeaderSize+DataSizeSize {
|
||||
err = nil
|
||||
}
|
||||
if count != NeedleHeaderSize+DataSizeSize || err != nil {
|
||||
return err
|
||||
}
|
||||
n.ParseNeedleHeader(bytes)
|
||||
if n.Size != size {
|
||||
if OffsetSize == 4 && offset < int64(MaxPossibleVolumeSize) {
|
||||
return ErrorSizeMismatch
|
||||
}
|
||||
}
|
||||
n.DataSize = util.BytesToUint32(bytes[NeedleHeaderSize : NeedleHeaderSize+DataSizeSize])
|
||||
startOffset := offset + NeedleHeaderSize
|
||||
if size.IsValid() {
|
||||
startOffset = offset + NeedleHeaderSize + DataSizeSize + int64(n.DataSize)
|
||||
}
|
||||
dataSize := GetActualSize(size, version)
|
||||
stopOffset := offset + dataSize
|
||||
metaSize := stopOffset - startOffset
|
||||
metaSlice := make([]byte, int(metaSize))
|
||||
|
||||
count, err = r.ReadAt(metaSlice, startOffset)
|
||||
if err != nil && int64(count) == metaSize {
|
||||
err = nil
|
||||
}
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
var index int
|
||||
if size.IsValid() {
|
||||
index, err = n.readNeedleDataVersion2NonData(metaSlice)
|
||||
}
|
||||
|
||||
n.Checksum = CRC(util.BytesToUint32(metaSlice[index : index+NeedleChecksumSize]))
|
||||
if version == Version3 {
|
||||
n.AppendAtNs = util.BytesToUint64(metaSlice[index+NeedleChecksumSize : index+NeedleChecksumSize+TimestampSize])
|
||||
}
|
||||
return err
|
||||
|
||||
func (n *Needle) ReadNeedleMeta(r backend.BackendStorageFile, offset int64, size Size, version Version) error {
|
||||
return n.ReadFromFile(r, offset, size, version, NeedleReadOptions{
|
||||
ReadHeader: true,
|
||||
ReadMeta: true,
|
||||
ReadData: false,
|
||||
})
|
||||
}
|
||||
|
||||
func min(x, y int64) int64 {
|
||||
|
||||
Reference in New Issue
Block a user