fix join and chunker interface change

This commit is contained in:
zelig 2015-02-11 12:12:24 +01:00
parent ce53889384
commit c4c8cc9bf6
5 changed files with 189 additions and 327 deletions

View file

@ -4,8 +4,6 @@ import (
"bytes"
"errors"
"io"
"sync"
"time"
)
type Bounded interface {
@ -24,12 +22,6 @@ type SectionReader interface {
io.ReaderAt
}
type LazySectionReader interface {
SectionReader
GetTimeout() time.Time
SetTimeout(time.Time) error
}
// ChunkReader implements SectionReader on a section
// of an underlying ReaderAt.
type ChunkReader struct {
@ -172,23 +164,6 @@ func (r *ChunkReader) WriteTo(w io.Writer) (n int64, err error) {
return
}
// LazyChunkReader implements LazySectionReader
type LazyChunkReader struct {
sections []LazySectionReader
readerFs [](func() LazySectionReader)
size int64
treeSize int64
off int64
}
func (self *LazyChunkReader) GetTimeout() (t time.Time) {
return
}
func (self *LazyChunkReader) SetTimeout(t time.Time) error {
return nil
}
func (self *LazyChunkReader) Size() (n int64) {
return self.size
}
@ -216,89 +191,3 @@ func (s *LazyChunkReader) Seek(offset int64, whence int) (int64, error) {
s.off = offset
return offset, nil
}
func (self *LazyChunkReader) ReadAt(b []byte, off int64) (read int, err error) {
if off < 0 || off >= self.size {
err = io.EOF
return
}
want := len(b)
got := int(self.size - off)
if want > got {
err = io.EOF
}
var index int
index = int(off / self.treeSize)
off -= int64(index) * self.treeSize
var limit int
var reader LazySectionReader
for i := index; i < len(self.sections) && want > 0; i++ {
if reader = self.sections[i]; reader == nil {
reader = self.readerFs[i]() // this hooks into the go routines and does a wg.Done once the needed chunks are fetched and processed
self.sections[i] = reader // memoize this reader
}
if want > int(reader.Size()-off) {
limit = int(reader.Size() - off)
} else {
limit = want
}
if got, err = reader.ReadAt(b[read:read+limit], off); err != nil {
return
}
read += got
want -= got
off = 0
}
return
}
type EmbeddedReader struct {
LazySectionReader
lock sync.Mutex
}
type NotReadyReader struct {
initF func()
r LazySectionReader
}
func (self *NotReadyReader) Size() (n int64) {
self.initF()
return self.r.Size()
}
func (self *NotReadyReader) Read(b []byte) (read int, err error) {
self.initF()
return self.r.Read(b)
}
func (self *NotReadyReader) ReadAt(b []byte, off int64) (read int, err error) {
self.initF()
return self.r.ReadAt(b, off)
}
func (self *NotReadyReader) Seek(offset int64, whence int) (int64, error) {
self.initF()
return self.r.Seek(offset, whence)
}
// just a wrapper to make a SectionReader conform to LazySectionReader
// with dummy implementation
type lazyReader struct {
SectionReader
}
func LazyReader(r SectionReader) (l LazySectionReader) {
return &lazyReader{
SectionReader: r,
}
}
func (self *lazyReader) GetTimeout() (t time.Time) {
return
}
func (self *lazyReader) SetTimeout(t time.Time) error {
return nil
}

View file

@ -26,6 +26,7 @@ import (
"crypto"
"encoding/binary"
"fmt"
"io"
"sync"
"time"
)
@ -63,14 +64,12 @@ type Chunker interface {
Split(key Key, data SectionReader, chunkC chan *Chunk) chan error
/*
Join reconstructs original content based on a root key.
When joining, the caller gets returned a LazySectionReader and an error channel.
When joining, the caller gets returned a Lazy SectionReader
New chunks to retrieve are coming to caller via the Chunk channel, which the caller provides.
If an error is encountered during joining, it is fed to errC error channel.
A closed error signals process completion at which point the data can be considered final and fully reconstructed if there were no errors.
The LazySectionReader provides on-demand fetching of content including chunks.
Lifecycle of the reader can be modified with SetTimeout()
If an error is encountered during joining, it appears as a reader error.
The SectionReader provides on-demand fetching of chunks.
*/
Join(key Key, chunkC chan *Chunk) (LazySectionReader, chan error)
Join(key Key, chunkC chan *Chunk) SectionReader
// returns the key length
KeySize() int64
@ -109,7 +108,7 @@ func (self *TreeChunker) Init() {
}
self.hashSize = int64(self.HashFunc.New().Size())
self.chunkSize = self.hashSize * self.Branches
// dpaLogger.Debugf("Chunker initialised: branches: %v, hashsize: %v, chunksize: %v, join timeout: %v , split timeout: %v", self.Branches, self.hashSize, self.chunkSize, self.JoinTimeout, self.SplitTimeout)
dpaLogger.Debugf("Chunker initialised: branches: %v, hashsize: %v, chunksize: %v, join timeout: %v , split timeout: %v", self.Branches, self.hashSize, self.chunkSize, self.JoinTimeout, self.SplitTimeout)
}
@ -120,9 +119,7 @@ func (self *TreeChunker) KeySize() int64 {
// String() for pretty printing
func (self *Chunk) String() string {
var size int64
var slice []byte
var n int
var err error
return fmt.Sprintf("Key: [%x..] TreeSize: %v Chunksize: %v Data: %x\n", self.Key[:4], self.Size, size, self.Data[:n])
}
@ -164,7 +161,7 @@ func (self *TreeChunker) Split(key Key, data SectionReader, chunkC chan *Chunk)
depth++
}
// dpaLogger.Debugf("split request received for data (%v bytes, depth: %v)", size, depth)
dpaLogger.Debugf("split request received for data (%v bytes, depth: %v)", size, depth)
//launch actual recursive function passing the workgroup
self.split(depth, treeSize/self.Branches, key, data, chunkC, rerrC, wg)
@ -200,7 +197,7 @@ func (self *TreeChunker) split(depth int, treeSize int64, key Key, data SectionR
size := data.Size()
var newChunk *Chunk
var hash Key
// dpaLogger.Debugf("depth: %v, max subtree size: %v, data size: %v", depth, treeSize, size)
dpaLogger.Debugf("depth: %v, max subtree size: %v, data size: %v", depth, treeSize, size)
for depth > 0 && size < treeSize {
treeSize /= self.Branches
@ -212,7 +209,7 @@ func (self *TreeChunker) split(depth int, treeSize int64, key Key, data SectionR
chunkData := make([]byte, data.Size())
data.ReadAt(chunkData, 0)
hash = self.Hash(size, chunkData)
// dpaLogger.Debugf("content chunk: max subtree size: %v, data size: %v", treeSize, size)
dpaLogger.Debugf("content chunk: max subtree size: %v, data size: %v", treeSize, size)
newChunk = &Chunk{
Key: hash,
Data: chunkData,
@ -220,15 +217,15 @@ func (self *TreeChunker) split(depth int, treeSize int64, key Key, data SectionR
}
} else {
// intermediate chunk containing child nodes hashes
branches := int64((size-1)/treeSize) + 1
// dpaLogger.Debugf("intermediate node: setting branches: %v, depth: %v, max subtree size: %v, data size: %v", branches, depth, treeSize, size)
branchCnt := int64((size-1)/treeSize) + 1
dpaLogger.Debugf("intermediate node: setting branches: %v, depth: %v, max subtree size: %v, data size: %v", branches, depth, treeSize, size)
var chunk []byte = make([]byte, branches*self.hashSize)
var pos, i int64
childrenWg := &sync.WaitGroup{}
var secSize int64
for i < branches {
for i < branchCnt {
// the last item can have shorter data
if size-pos < treeSize {
secSize = size - pos
@ -269,132 +266,142 @@ func (self *TreeChunker) split(depth int, treeSize int64, key Key, data SectionR
}
func (self *TreeChunker) Join(key Key, chunkC chan *Chunk) (data LazySectionReader, errC chan error) {
// initialise return parameters
errC = make(chan error)
// timer to time out the operation (needed within so as to avoid process leakage)
timeout := time.After(self.JoinTimeout)
wg := &sync.WaitGroup{}
// initialise internal error channel
rerrC := make(chan error)
quitC := make(chan bool)
func (self *TreeChunker) Join(key Key, chunkC chan *Chunk) SectionReader {
// launch lazy recursive call on root chunk
notReadyReader := &NotReadyReader{}
reader := &EmbeddedReader{
LazySectionReader: &lazyReader{notReadyReader},
return &LazyChunkReader{
key: key,
chunkC: chunkC,
quitC: make(chan bool),
errC: make(chan error),
chunker: self,
}
fun := self.join(0, 0, key, chunkC, rerrC, wg, quitC)
data = reader
}
// processing is triggered by reads on the LazySectionReader
// LazyChunkReader implements LazySectionReader
type LazyChunkReader struct {
key Key
chunkC chan *Chunk
size int64
off int64
quitC chan bool
errC chan error
chunker *TreeChunker
}
func (self *LazyChunkReader) ReadAt(b []byte, off int64) (read int, err error) {
chunk := &Chunk{
Key: self.key,
C: make(chan bool), // close channel to signal data delivery
}
self.chunkC <- chunk // submit retrieval request, someone should be listening on the other side (or we will time out globally)
dpaLogger.Debugf("readAt %x", chunk.Key[:4])
// waiting for the chunk retrieval
select {
case <-self.quitC:
// this is how we control process leakage (quitC is closed once join is finished (after timeout))
dpaLogger.Debugf("quit")
return
case <-chunk.C: // bells are ringing, data have been delivered
dpaLogger.Debugf("chunk data received for %x", chunk.Key[:4])
fmt.Printf("chunk data received for %x\n", chunk.Key[:4])
}
if chunk.Data == nil {
return 0, notFound
}
want := int64(len(b))
if off < 0 || want+off > chunk.Size {
return 0, io.EOF
}
var treeSize int64
var depth int
// calculate depth and max treeSize
treeSize = self.chunker.chunkSize
for ; treeSize < chunk.Size; treeSize *= self.chunker.Branches {
depth++
}
wg := sync.WaitGroup{}
wg.Add(1)
notReadyReader.r = reader
notReadyReader.initF = func() {
reader.lock.Lock()
if fun != nil {
reader.LazySectionReader = fun() // replace the reader while caller is waiting
wg.Done() // just to have one in waitgroup until reader funcs are called
fun = nil // to be only called once
}
reader.lock.Unlock()
}
// waits for all the processes to finish and signals by closing internal rerrc
go self.join(b, off, off+want, depth, treeSize/self.chunker.Branches, chunk, &wg)
go func() {
wg.Wait()
close(rerrC)
close(self.errC)
}()
go func() {
select {
case err := <-rerrC:
if err != nil {
errC <- err
} // otherwise channel is closed, data joining complete
case <-timeout:
errC <- fmt.Errorf("join time out")
close(quitC)
}
// this will indicate to the caller that processing is finished (with or without error)
close(errC)
}()
return
}
func (self *TreeChunker) join(depth int, treeSize int64, key Key, chunkC chan *Chunk, errC chan error, wg *sync.WaitGroup, quitC chan bool) (readerF func() LazySectionReader) {
wg.Add(1)
readerF = func() (r LazySectionReader) {
defer wg.Done()
chunk := &Chunk{
Key: key,
C: make(chan bool, 1), // close channel to signal data delivery
}
chunkC <- chunk // submit retrieval request, someone should be listening on the other side (or we will time out globally)
// waiting for the chunk retrieval
select {
case <-quitC:
// this is how we control process leakage (quitC is closed once join is finished (after timeout))
return
case <-chunk.C: // bells are ringing, data have been delivered
}
// if data == nil {
// return
// }
// calculate depth and max treeSize
var depth int
var treeSize int64 = self.hashSize
for ; treeSize*self.Branches < chunk.Size; treeSize *= self.Branches {
depth++
}
if depth == 0 {
r = LazyReader(NewByteSliceReader(chunk.Data))
return // simply give back the chunks reader for content chunks
}
// find appropriate block level
for chunk.Size < treeSize {
treeSize /= self.Branches
depth--
}
// boooo
// intermediate chunk, chunk containing hashes of child nodes
var i int64
var childKey Key
var readerF func() (r LazySectionReader)
var readerFs [](func() (r LazySectionReader))
branches := int64((chunk.Size-1)/treeSize) + 1
// dpaLogger.DebugDetailf("tree node - size %v, chunk size: %v, subtreeSize %v, branches %v", chunk.Size, chunk.Reader.Size(), treeSize, branches)
// iterate through the chunk containing the keys of children
// create lazy init functions that give back readers
for i < branches {
// create partial Chunk in order to send a retrieval request
childKey = make([]byte, self.hashSize) // preallocate hashSize long slice for key
// read the Hash of the subtree from the relevant section of the Chunk into the allocated byte slice in subtree.Key
if _, err := chunk.Reader.ReadAt(childKey, i*self.hashSize); err != nil {
// dpaLogger.DebugDetailf("Read error: %v", err)
errC <- err
break
}
// call lazy reader function recursively on the subtree
readerF = self.join(depth-1, treeSize/self.Branches, childKey, chunkC, errC, wg, quitC)
readerFs = append(readerFs, readerF)
i++
}
// new reader created on demand:
r = &LazyChunkReader{
readerFs: readerFs,
sections: make([]LazySectionReader, branches),
size: chunk.Size,
treeSize: treeSize,
}
return
select {
case err = <-self.errC:
case <-self.quitC:
read = len(b)
}
return
}
func (self *LazyChunkReader) join(b []byte, off int64, eoff int64, depth int, treeSize int64, chunk *Chunk, parentWg *sync.WaitGroup) {
defer parentWg.Done()
dpaLogger.Debugf("depth: %v, loff: %v, eoff: %v, chunk.Size: %v, treeSize: %v", depth, off, eoff, chunk.Size, treeSize)
fmt.Printf("depth: %v, loff: %v, eoff: %v, chunk.Size: %v, treeSize: %v\n", depth, off, eoff, chunk.Size, treeSize)
// find appropriate block level
for chunk.Size < treeSize && depth > 0 {
treeSize /= self.chunker.Branches
depth--
}
dpaLogger.Debugf("-> depth: %v, loff: %v, eoff: %v, chunk.Size: %v, treeSize: %v", depth, off, eoff, chunk.Size, treeSize)
fmt.Printf("-> depth: %v, loff: %v, eoff: %v, chunk.Size: %v, treeSize: %v\n", depth, off, eoff, chunk.Size, treeSize)
if depth == 0 {
copy(b, chunk.Data[off:eoff])
return // simply give back the chunks reader for content chunks
}
// subtree index
start := off / treeSize
end := (eoff + treeSize - 1) / treeSize
wg := sync.WaitGroup{}
for i := start; i < end; i++ {
soff := i * treeSize
roff := soff
seoff := soff + treeSize
if soff < off {
soff = off
}
if seoff > eoff {
seoff = eoff
}
wg.Add(1)
go func(j int64) {
dpaLogger.Debugf("subtree index: %v", j)
childKey := chunk.Data[j*self.chunker.hashSize : (j+1)*self.chunker.hashSize]
chunk := &Chunk{
Key: childKey,
C: make(chan bool, 1), // close channel to signal data delivery
}
dpaLogger.Debugf("chunk data sent for %x (key interval in chunk %v-%v)", chunk.Key[:4], j*self.chunker.hashSize, (j+1)*self.chunker.hashSize)
self.chunkC <- chunk // submit retrieval request, someone should be listening on the other side (or we will time out globally)
// waiting for the chunk retrieval
select {
case <-self.quitC:
// this is how we control process leakage (quitC is closed once join is finished (after timeout))
return
case <-chunk.C: // bells are ringing, data have been delivered
dpaLogger.Debugf("chunk data received")
}
if soff < off {
soff = off
}
if chunk.Data == nil {
self.errC <- fmt.Errorf("chunk %v-%v not found", off, off+treeSize)
return
}
self.join(b[soff-off:seoff-off], soff-roff, seoff-roff, depth-1, treeSize/self.chunker.Branches, chunk, &wg)
}(i)
} //for
wg.Wait()
}

View file

@ -59,6 +59,7 @@ func (self *chunkerTester) Split(chunker *TreeChunker, l int) (key Key, input []
if err != nil {
self.errors = append(self.errors, err)
}
fmt.Printf("err %v", err)
if !ok {
close(chunkC)
errC = nil
@ -71,13 +72,13 @@ func (self *chunkerTester) Split(chunker *TreeChunker, l int) (key Key, input []
return
}
func (self *chunkerTester) Join(chunker *TreeChunker, key Key, c int) (LazySectionReader, chan bool) {
func (self *chunkerTester) Join(chunker *TreeChunker, key Key, c int) SectionReader {
// reset but not the chunks
self.errors = nil
self.timeout = false
chunkC := make(chan *Chunk, 1000)
reader, errC := chunker.Join(key, chunkC)
reader := chunker.Join(key, chunkC)
quitC := make(chan bool)
timeout := time.After(600 * time.Second)
@ -95,55 +96,50 @@ func (self *chunkerTester) Join(chunker *TreeChunker, key Key, c int) (LazySecti
case chunk := <-chunkC:
i++
dpaLogger.DebugDetailf("TESTER: chunk request %x", chunk.Key[:4])
// this just mocks the behaviour of a chunk store retrieval
var found bool
for _, ch := range self.chunks {
if bytes.Compare(chunk.Key, ch.Key) == 0 {
if bytes.Equal(chunk.Key, ch.Key) {
found = true
chunk.Reader = ch.Reader
chunk.Data = ch.Data
chunk.Size = ch.Size
close(chunk.C)
break
}
}
if !found {
fmt.Printf("chunk request unknown for %x", chunk.Key[:4])
}
case err, ok := <-errC:
if err != nil {
fmt.Printf("error %v", err)
self.errors = append(self.errors, err)
}
if !ok {
close(chunkC)
errC = nil
break LOOP
fmt.Printf("TESTER: chunk unknown for %x", chunk.Key[:4])
}
close(chunk.C)
dpaLogger.DebugDetailf("TESTER: chunk request served %x", chunk.Key[:4])
}
}
}()
return reader, quitC
return reader
}
func testRandomData(chunker *TreeChunker, tester *chunkerTester, n int, chunks int, t *testing.T) {
key, input := tester.Split(chunker, n)
fmt.Printf("split returned\n")
fmt.Printf("input\n%x\n", input)
tester.checkChunks(t, chunks)
t.Logf("chunks: %v", tester.chunks)
reader, quitC := tester.Join(chunker, key, 0)
output := make([]byte, reader.Size())
fmt.Printf("chunks: %v", tester.chunks)
reader := tester.Join(chunker, key, 0)
output := make([]byte, n)
_, err := reader.Read(output)
if err != nil {
t.Errorf("read error %v\n", err)
}
t.Logf(" IN: %x\nOUT: %x\n", input, output)
if bytes.Compare(output, input) != 0 {
if !bytes.Equal(output, input) {
t.Errorf("input and output mismatch\n IN: %x\nOUT: %x\n", input, output)
}
close(quitC)
}
func TestRandomData(t *testing.T) {
defer test.Testlog(t).Detach()
test.LogInit()
chunker := &TreeChunker{
Branches: 2,
SplitTimeout: 10 * time.Second,
@ -168,11 +164,17 @@ func chunkerAndTester() (chunker *TreeChunker, tester *chunkerTester) {
return
}
func readAll(reader SectionReader) {
size := reader.Size()
output := make([]byte, 1000)
func readAll(reader SectionReader, result []byte) {
size := int64(len(result))
var end int64
for pos := int64(0); pos < size; pos += 1000 {
reader.ReadAt(output, pos)
if pos+1000 > size {
end = size
} else {
end = pos + 1000
}
reader.ReadAt(result[pos:end], pos)
}
}
@ -182,9 +184,9 @@ func benchmarkJoinRandomData(n int, chunks int, t *testing.B) {
chunker, tester := chunkerAndTester()
key, _ := tester.Split(chunker, n)
t.StartTimer()
reader, quitC := tester.Join(chunker, key, i)
readAll(reader)
close(quitC)
reader := tester.Join(chunker, key, i)
result := make([]byte, n)
readAll(reader, result)
}
}

View file

@ -38,10 +38,7 @@ SPLIT:
for {
select {
case chunk := <-chunkC:
chunk.Data = make([]byte, chunk.Reader.Size())
chunk.Reader.ReadAt(chunk.Data, 0)
m.Put(chunk)
case err, ok := <-errC:
if err != nil {
t.Errorf("Chunker error: %v", err)
@ -53,45 +50,30 @@ SPLIT:
}
}
}
chunker := &TreeChunker{
Branches: branches,
}
chunker.Init()
chunkC = make(chan *Chunk)
var r LazySectionReader
r, errC = chunker.Join(key, chunkC)
var r SectionReader
r = chunker.Join(key, chunkC)
quit := make(chan bool)
go func() {
JOIN:
for {
select {
case chunk := <-chunkC:
go func() {
storedChunk, err := m.Get(chunk.Key)
if err == notFound {
dpaLogger.DebugDetailf("chunk '%x' not found", chunk.Key)
} else if err != nil {
dpaLogger.DebugDetailf("error retrieving chunk %x: %v", chunk.Key, err)
} else {
chunk.Reader = NewChunkReaderFromBytes(storedChunk.Data)
chunk.Size = storedChunk.Size
}
close(chunk.C)
}()
case err, ok := <-errC:
if err != nil {
t.Errorf("Chunker error: %v", err)
return
for chunk := range chunkC {
go func() {
storedChunk, err := m.Get(chunk.Key)
if err == notFound {
dpaLogger.DebugDetailf("chunk '%x' not found", chunk.Key)
} else if err != nil {
dpaLogger.DebugDetailf("error retrieving chunk %x: %v", chunk.Key, err)
} else {
chunk.Data = storedChunk.Data
chunk.Size = storedChunk.Size
}
if !ok {
break JOIN
}
case <-quit:
break JOIN
}
close(chunk.C)
}()
}
}()

View file

@ -68,28 +68,10 @@ type ChunkStore interface {
Get(Key) (*Chunk, error)
}
func (self *DPA) Retrieve(key Key) LazySectionReader {
func (self *DPA) Retrieve(key Key) SectionReader {
reader, errC := self.Chunker.Join(key, self.retrieveC)
return self.Chunker.Join(key, self.retrieveC)
// we can add subscriptions etc. or timeout here
go func() {
JOIN:
for {
select {
case err, ok := <-errC:
if err != nil {
dpaLogger.Warnf("chunker join error: %v", err)
}
if !ok {
break JOIN
}
case <-self.quitC:
break JOIN
}
}
}()
return reader
}
func (self *DPA) Store(data SectionReader) (key Key, err error) {
@ -151,7 +133,7 @@ func (self *DPA) retrieveLoop() {
} else if err != nil {
dpaLogger.DebugDetailf("error retrieving chunk %x: %v", chunk.Key, err)
} else {
chunk.Reader = NewChunkReaderFromBytes(storedChunk.Data)
chunk.Data = storedChunk.Data
chunk.Size = storedChunk.Size
}
close(chunk.C)