|
@ -1,6 +1,7 @@ |
|
|
package needle |
|
|
package needle |
|
|
|
|
|
|
|
|
import ( |
|
|
import ( |
|
|
|
|
|
"bytes" |
|
|
"errors" |
|
|
"errors" |
|
|
"fmt" |
|
|
"fmt" |
|
|
"github.com/chrislusf/seaweedfs/weed/glog" |
|
|
"github.com/chrislusf/seaweedfs/weed/glog" |
|
@ -9,6 +10,7 @@ import ( |
|
|
"github.com/chrislusf/seaweedfs/weed/util" |
|
|
"github.com/chrislusf/seaweedfs/weed/util" |
|
|
"io" |
|
|
"io" |
|
|
"math" |
|
|
"math" |
|
|
|
|
|
"sync" |
|
|
) |
|
|
) |
|
|
|
|
|
|
|
|
const ( |
|
|
const ( |
|
@ -29,10 +31,14 @@ func (n *Needle) DiskSize(version Version) int64 { |
|
|
return GetActualSize(n.Size, version) |
|
|
return GetActualSize(n.Size, version) |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
func (n *Needle) prepareWriteBuffer(version Version) ([]byte, Size, int64, error) { |
|
|
|
|
|
|
|
|
|
|
|
writeBytes := make([]byte, 0) |
|
|
|
|
|
|
|
|
var bufPool = sync.Pool{ |
|
|
|
|
|
New: func() interface{} { |
|
|
|
|
|
return new(bytes.Buffer) |
|
|
|
|
|
}, |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
func (n *Needle) prepareWriteBuffer(version Version, writeBytes *bytes.Buffer) (Size, int64, error) { |
|
|
|
|
|
writeBytes.Reset() |
|
|
switch version { |
|
|
switch version { |
|
|
case Version1: |
|
|
case Version1: |
|
|
header := make([]byte, NeedleHeaderSize) |
|
|
header := make([]byte, NeedleHeaderSize) |
|
@ -42,12 +48,12 @@ func (n *Needle) prepareWriteBuffer(version Version) ([]byte, Size, int64, error |
|
|
SizeToBytes(header[CookieSize+NeedleIdSize:CookieSize+NeedleIdSize+SizeSize], n.Size) |
|
|
SizeToBytes(header[CookieSize+NeedleIdSize:CookieSize+NeedleIdSize+SizeSize], n.Size) |
|
|
size := n.Size |
|
|
size := n.Size |
|
|
actualSize := NeedleHeaderSize + int64(n.Size) |
|
|
actualSize := NeedleHeaderSize + int64(n.Size) |
|
|
writeBytes = append(writeBytes, header...) |
|
|
|
|
|
writeBytes = append(writeBytes, n.Data...) |
|
|
|
|
|
|
|
|
writeBytes.Write(header) |
|
|
|
|
|
writeBytes.Write(n.Data) |
|
|
padding := PaddingLength(n.Size, version) |
|
|
padding := PaddingLength(n.Size, version) |
|
|
util.Uint32toBytes(header[0:NeedleChecksumSize], n.Checksum.Value()) |
|
|
util.Uint32toBytes(header[0:NeedleChecksumSize], n.Checksum.Value()) |
|
|
writeBytes = append(writeBytes, header[0:NeedleChecksumSize+padding]...) |
|
|
|
|
|
return writeBytes, size, actualSize, nil |
|
|
|
|
|
|
|
|
writeBytes.Write(header[0:NeedleChecksumSize+padding]) |
|
|
|
|
|
return size, actualSize, nil |
|
|
case Version2, Version3: |
|
|
case Version2, Version3: |
|
|
header := make([]byte, NeedleHeaderSize+TimestampSize) // adding timestamp to reuse it and avoid extra allocation
|
|
|
header := make([]byte, NeedleHeaderSize+TimestampSize) // adding timestamp to reuse it and avoid extra allocation
|
|
|
CookieToBytes(header[0:CookieSize], n.Cookie) |
|
|
CookieToBytes(header[0:CookieSize], n.Cookie) |
|
@ -79,51 +85,51 @@ func (n *Needle) prepareWriteBuffer(version Version) ([]byte, Size, int64, error |
|
|
n.Size = 0 |
|
|
n.Size = 0 |
|
|
} |
|
|
} |
|
|
SizeToBytes(header[CookieSize+NeedleIdSize:CookieSize+NeedleIdSize+SizeSize], n.Size) |
|
|
SizeToBytes(header[CookieSize+NeedleIdSize:CookieSize+NeedleIdSize+SizeSize], n.Size) |
|
|
writeBytes = append(writeBytes, header[0:NeedleHeaderSize]...) |
|
|
|
|
|
|
|
|
writeBytes.Write(header[0:NeedleHeaderSize]) |
|
|
if n.DataSize > 0 { |
|
|
if n.DataSize > 0 { |
|
|
util.Uint32toBytes(header[0:4], n.DataSize) |
|
|
util.Uint32toBytes(header[0:4], n.DataSize) |
|
|
writeBytes = append(writeBytes, header[0:4]...) |
|
|
|
|
|
writeBytes = append(writeBytes, n.Data...) |
|
|
|
|
|
|
|
|
writeBytes.Write(header[0:4]) |
|
|
|
|
|
writeBytes.Write(n.Data) |
|
|
util.Uint8toBytes(header[0:1], n.Flags) |
|
|
util.Uint8toBytes(header[0:1], n.Flags) |
|
|
writeBytes = append(writeBytes, header[0:1]...) |
|
|
|
|
|
|
|
|
writeBytes.Write(header[0:1]) |
|
|
if n.HasName() { |
|
|
if n.HasName() { |
|
|
util.Uint8toBytes(header[0:1], n.NameSize) |
|
|
util.Uint8toBytes(header[0:1], n.NameSize) |
|
|
writeBytes = append(writeBytes, header[0:1]...) |
|
|
|
|
|
writeBytes = append(writeBytes, n.Name[:n.NameSize]...) |
|
|
|
|
|
|
|
|
writeBytes.Write(header[0:1]) |
|
|
|
|
|
writeBytes.Write(n.Name[:n.NameSize]) |
|
|
} |
|
|
} |
|
|
if n.HasMime() { |
|
|
if n.HasMime() { |
|
|
util.Uint8toBytes(header[0:1], n.MimeSize) |
|
|
util.Uint8toBytes(header[0:1], n.MimeSize) |
|
|
writeBytes = append(writeBytes, header[0:1]...) |
|
|
|
|
|
writeBytes = append(writeBytes, n.Mime...) |
|
|
|
|
|
|
|
|
writeBytes.Write(header[0:1]) |
|
|
|
|
|
writeBytes.Write(n.Mime) |
|
|
} |
|
|
} |
|
|
if n.HasLastModifiedDate() { |
|
|
if n.HasLastModifiedDate() { |
|
|
util.Uint64toBytes(header[0:8], n.LastModified) |
|
|
util.Uint64toBytes(header[0:8], n.LastModified) |
|
|
writeBytes = append(writeBytes, header[8-LastModifiedBytesLength:8]...) |
|
|
|
|
|
|
|
|
writeBytes.Write(header[8-LastModifiedBytesLength:8]) |
|
|
} |
|
|
} |
|
|
if n.HasTtl() && n.Ttl != nil { |
|
|
if n.HasTtl() && n.Ttl != nil { |
|
|
n.Ttl.ToBytes(header[0:TtlBytesLength]) |
|
|
n.Ttl.ToBytes(header[0:TtlBytesLength]) |
|
|
writeBytes = append(writeBytes, header[0:TtlBytesLength]...) |
|
|
|
|
|
|
|
|
writeBytes.Write(header[0:TtlBytesLength]) |
|
|
} |
|
|
} |
|
|
if n.HasPairs() { |
|
|
if n.HasPairs() { |
|
|
util.Uint16toBytes(header[0:2], n.PairsSize) |
|
|
util.Uint16toBytes(header[0:2], n.PairsSize) |
|
|
writeBytes = append(writeBytes, header[0:2]...) |
|
|
|
|
|
writeBytes = append(writeBytes, n.Pairs...) |
|
|
|
|
|
|
|
|
writeBytes.Write(header[0:2]) |
|
|
|
|
|
writeBytes.Write(n.Pairs) |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
padding := PaddingLength(n.Size, version) |
|
|
padding := PaddingLength(n.Size, version) |
|
|
util.Uint32toBytes(header[0:NeedleChecksumSize], n.Checksum.Value()) |
|
|
util.Uint32toBytes(header[0:NeedleChecksumSize], n.Checksum.Value()) |
|
|
if version == Version2 { |
|
|
if version == Version2 { |
|
|
writeBytes = append(writeBytes, header[0:NeedleChecksumSize+padding]...) |
|
|
|
|
|
|
|
|
writeBytes.Write(header[0:NeedleChecksumSize+padding]) |
|
|
} else { |
|
|
} else { |
|
|
// version3
|
|
|
// version3
|
|
|
util.Uint64toBytes(header[NeedleChecksumSize:NeedleChecksumSize+TimestampSize], n.AppendAtNs) |
|
|
util.Uint64toBytes(header[NeedleChecksumSize:NeedleChecksumSize+TimestampSize], n.AppendAtNs) |
|
|
writeBytes = append(writeBytes, header[0:NeedleChecksumSize+TimestampSize+padding]...) |
|
|
|
|
|
|
|
|
writeBytes.Write(header[0:NeedleChecksumSize+TimestampSize+padding]) |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
return writeBytes, Size(n.DataSize), GetActualSize(n.Size, version), nil |
|
|
|
|
|
|
|
|
return Size(n.DataSize), GetActualSize(n.Size, version), nil |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
return writeBytes, 0, 0, fmt.Errorf("Unsupported Version! (%d)", version) |
|
|
|
|
|
|
|
|
return 0, 0, fmt.Errorf("Unsupported Version! (%d)", version) |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
func (n *Needle) Append(w backend.BackendStorageFile, version Version) (offset uint64, size Size, actualSize int64, err error) { |
|
|
func (n *Needle) Append(w backend.BackendStorageFile, version Version) (offset uint64, size Size, actualSize int64, err error) { |
|
@ -146,10 +152,13 @@ func (n *Needle) Append(w backend.BackendStorageFile, version Version) (offset u |
|
|
return |
|
|
return |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
bytesToWrite, size, actualSize, err := n.prepareWriteBuffer(version) |
|
|
|
|
|
|
|
|
bytesBuffer := bufPool.Get().(*bytes.Buffer) |
|
|
|
|
|
defer bufPool.Put(bytesBuffer) |
|
|
|
|
|
|
|
|
|
|
|
size, actualSize, err = n.prepareWriteBuffer(version, bytesBuffer) |
|
|
|
|
|
|
|
|
if err == nil { |
|
|
if err == nil { |
|
|
_, err = w.WriteAt(bytesToWrite, int64(offset)) |
|
|
|
|
|
|
|
|
_, err = w.WriteAt(bytesBuffer.Bytes(), int64(offset)) |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
return offset, size, actualSize, err |
|
|
return offset, size, actualSize, err |
|
|