mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-07-23 05:06:43 +00:00
fix chunker test io.EOF
This commit is contained in:
parent
2b1b2bfdb7
commit
68b6ee18f8
3 changed files with 24 additions and 24 deletions
|
|
@ -108,7 +108,7 @@ func (self *TreeChunker) Init() {
|
||||||
}
|
}
|
||||||
self.hashSize = int64(self.HashFunc.New().Size())
|
self.hashSize = int64(self.HashFunc.New().Size())
|
||||||
self.chunkSize = self.hashSize * self.Branches
|
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)
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -161,7 +161,7 @@ func (self *TreeChunker) Split(key Key, data SectionReader, chunkC chan *Chunk)
|
||||||
depth++
|
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
|
//launch actual recursive function passing the workgroup
|
||||||
self.split(depth, treeSize/self.Branches, key, data, chunkC, rerrC, wg)
|
self.split(depth, treeSize/self.Branches, key, data, chunkC, rerrC, wg)
|
||||||
|
|
@ -197,7 +197,7 @@ func (self *TreeChunker) split(depth int, treeSize int64, key Key, data SectionR
|
||||||
size := data.Size()
|
size := data.Size()
|
||||||
var newChunk *Chunk
|
var newChunk *Chunk
|
||||||
var hash Key
|
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 {
|
for depth > 0 && size < treeSize {
|
||||||
treeSize /= self.Branches
|
treeSize /= self.Branches
|
||||||
|
|
@ -209,7 +209,7 @@ func (self *TreeChunker) split(depth int, treeSize int64, key Key, data SectionR
|
||||||
chunkData := make([]byte, data.Size())
|
chunkData := make([]byte, data.Size())
|
||||||
data.ReadAt(chunkData, 0)
|
data.ReadAt(chunkData, 0)
|
||||||
hash = self.Hash(size, chunkData)
|
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{
|
newChunk = &Chunk{
|
||||||
Key: hash,
|
Key: hash,
|
||||||
Data: chunkData,
|
Data: chunkData,
|
||||||
|
|
@ -218,7 +218,7 @@ func (self *TreeChunker) split(depth int, treeSize int64, key Key, data SectionR
|
||||||
} else {
|
} else {
|
||||||
// intermediate chunk containing child nodes hashes
|
// intermediate chunk containing child nodes hashes
|
||||||
branchCnt := int64((size-1)/treeSize) + 1
|
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)
|
// 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 chunk []byte = make([]byte, branches*self.hashSize)
|
||||||
var pos, i int64
|
var pos, i int64
|
||||||
|
|
@ -295,24 +295,24 @@ func (self *LazyChunkReader) ReadAt(b []byte, off int64) (read int, err error) {
|
||||||
C: make(chan bool), // close channel to signal data delivery
|
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)
|
self.chunkC <- chunk // submit retrieval request, someone should be listening on the other side (or we will time out globally)
|
||||||
dpaLogger.Debugf("readAt %x into %d bytes at %d.", chunk.Key[:4], len(b), off)
|
// dpaLogger.Debugf("readAt %x into %d bytes at %d.", chunk.Key[:4], len(b), off)
|
||||||
|
|
||||||
// waiting for the chunk retrieval
|
// waiting for the chunk retrieval
|
||||||
select {
|
select {
|
||||||
case <-self.quitC:
|
case <-self.quitC:
|
||||||
// this is how we control process leakage (quitC is closed once join is finished (after timeout))
|
// this is how we control process leakage (quitC is closed once join is finished (after timeout))
|
||||||
dpaLogger.Debugf("quit")
|
// dpaLogger.Debugf("quit")
|
||||||
return
|
return
|
||||||
case <-chunk.C: // bells are ringing, data have been delivered
|
case <-chunk.C: // bells are ringing, data have been delivered
|
||||||
dpaLogger.Debugf("chunk data received for %x", chunk.Key[:4])
|
// dpaLogger.Debugf("chunk data received for %x", chunk.Key[:4])
|
||||||
}
|
}
|
||||||
if len(chunk.Data) == 0 {
|
if len(chunk.Data) == 0 {
|
||||||
dpaLogger.Debugf("No payload.")
|
// dpaLogger.Debugf("No payload.")
|
||||||
return 0, notFound
|
return 0, notFound
|
||||||
}
|
}
|
||||||
self.size = chunk.Size
|
self.size = chunk.Size
|
||||||
if b == nil {
|
if b == nil {
|
||||||
dpaLogger.Debugf("Size query for %x.", chunk.Key[:4])
|
// dpaLogger.Debugf("Size query for %x.", chunk.Key[:4])
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
want := int64(len(b))
|
want := int64(len(b))
|
||||||
|
|
@ -335,14 +335,14 @@ func (self *LazyChunkReader) ReadAt(b []byte, off int64) (read int, err error) {
|
||||||
}()
|
}()
|
||||||
select {
|
select {
|
||||||
case err = <-self.errC:
|
case err = <-self.errC:
|
||||||
dpaLogger.Debugf("ReadAt received %v.", err)
|
// dpaLogger.Debugf("ReadAt received %v.", err)
|
||||||
read = len(b)
|
read = len(b)
|
||||||
if off+int64(read) == self.size {
|
if off+int64(read) == self.size {
|
||||||
err = io.EOF
|
err = io.EOF
|
||||||
}
|
}
|
||||||
dpaLogger.Debugf("ReadAt returning with %d, %v.", read, err)
|
// dpaLogger.Debugf("ReadAt returning with %d, %v.", read, err)
|
||||||
case <-self.quitC:
|
case <-self.quitC:
|
||||||
dpaLogger.Debugf("ReadAt aborted with %d, %v.", read, err)
|
// dpaLogger.Debugf("ReadAt aborted with %d, %v.", read, err)
|
||||||
}
|
}
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
@ -350,7 +350,7 @@ func (self *LazyChunkReader) ReadAt(b []byte, off int64) (read int, err error) {
|
||||||
func (self *LazyChunkReader) join(b []byte, off int64, eoff int64, depth int, treeSize int64, chunk *Chunk, parentWg *sync.WaitGroup) {
|
func (self *LazyChunkReader) join(b []byte, off int64, eoff int64, depth int, treeSize int64, chunk *Chunk, parentWg *sync.WaitGroup) {
|
||||||
defer parentWg.Done()
|
defer parentWg.Done()
|
||||||
|
|
||||||
dpaLogger.Debugf("depth: %v, loff: %v, eoff: %v, chunk.Size: %v, treeSize: %v", depth, off, eoff, chunk.Size, treeSize)
|
// dpaLogger.Debugf("depth: %v, loff: %v, eoff: %v, chunk.Size: %v, treeSize: %v", depth, off, eoff, chunk.Size, treeSize)
|
||||||
|
|
||||||
// find appropriate block level
|
// find appropriate block level
|
||||||
for chunk.Size < treeSize && depth > 0 {
|
for chunk.Size < treeSize && depth > 0 {
|
||||||
|
|
@ -359,7 +359,7 @@ func (self *LazyChunkReader) join(b []byte, off int64, eoff int64, depth int, tr
|
||||||
}
|
}
|
||||||
|
|
||||||
if depth == 0 {
|
if depth == 0 {
|
||||||
dpaLogger.Debugf("len(b): %v, off: %v, eoff: %v, chunk.Size: %v, treeSize: %v", depth, len(b), off, eoff, chunk.Size, treeSize)
|
// dpaLogger.Debugf("len(b): %v, off: %v, eoff: %v, chunk.Size: %v, treeSize: %v", depth, len(b), off, eoff, chunk.Size, treeSize)
|
||||||
|
|
||||||
copy(b, chunk.Data[off:eoff])
|
copy(b, chunk.Data[off:eoff])
|
||||||
return // simply give back the chunks reader for content chunks
|
return // simply give back the chunks reader for content chunks
|
||||||
|
|
@ -386,13 +386,13 @@ func (self *LazyChunkReader) join(b []byte, off int64, eoff int64, depth int, tr
|
||||||
wg.Add(1)
|
wg.Add(1)
|
||||||
go func(j int64) {
|
go func(j int64) {
|
||||||
childKey := chunk.Data[j*self.chunker.hashSize : (j+1)*self.chunker.hashSize]
|
childKey := chunk.Data[j*self.chunker.hashSize : (j+1)*self.chunker.hashSize]
|
||||||
dpaLogger.Debugf("subtree index: %v -> %x", j, childKey[:4])
|
// dpaLogger.Debugf("subtree index: %v -> %x", j, childKey[:4])
|
||||||
|
|
||||||
ch := &Chunk{
|
ch := &Chunk{
|
||||||
Key: childKey,
|
Key: childKey,
|
||||||
C: make(chan bool), // close channel to signal data delivery
|
C: make(chan bool), // close channel to signal data delivery
|
||||||
}
|
}
|
||||||
dpaLogger.Debugf("chunk data sent for %x (key interval in chunk %v-%v)", ch.Key[:4], j*self.chunker.hashSize, (j+1)*self.chunker.hashSize)
|
// dpaLogger.Debugf("chunk data sent for %x (key interval in chunk %v-%v)", ch.Key[:4], j*self.chunker.hashSize, (j+1)*self.chunker.hashSize)
|
||||||
self.chunkC <- ch // submit retrieval request, someone should be listening on the other side (or we will time out globally)
|
self.chunkC <- ch // submit retrieval request, someone should be listening on the other side (or we will time out globally)
|
||||||
|
|
||||||
// waiting for the chunk retrieval
|
// waiting for the chunk retrieval
|
||||||
|
|
@ -401,7 +401,7 @@ func (self *LazyChunkReader) join(b []byte, off int64, eoff int64, depth int, tr
|
||||||
// this is how we control process leakage (quitC is closed once join is finished (after timeout))
|
// this is how we control process leakage (quitC is closed once join is finished (after timeout))
|
||||||
return
|
return
|
||||||
case <-ch.C: // bells are ringing, data have been delivered
|
case <-ch.C: // bells are ringing, data have been delivered
|
||||||
dpaLogger.Debugf("chunk data received")
|
// dpaLogger.Debugf("chunk data received")
|
||||||
}
|
}
|
||||||
if soff < off {
|
if soff < off {
|
||||||
soff = off
|
soff = off
|
||||||
|
|
|
||||||
|
|
@ -3,6 +3,7 @@ package bzz
|
||||||
import (
|
import (
|
||||||
"bytes"
|
"bytes"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"io"
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
|
|
@ -59,7 +60,7 @@ func (self *chunkerTester) Split(chunker *TreeChunker, l int) (key Key, input []
|
||||||
if err != nil {
|
if err != nil {
|
||||||
self.errors = append(self.errors, err)
|
self.errors = append(self.errors, err)
|
||||||
}
|
}
|
||||||
fmt.Printf("err %v", err)
|
// fmt.Printf("err %v", err)
|
||||||
if !ok {
|
if !ok {
|
||||||
close(chunkC)
|
close(chunkC)
|
||||||
errC = nil
|
errC = nil
|
||||||
|
|
@ -96,7 +97,7 @@ func (self *chunkerTester) Join(chunker *TreeChunker, key Key, c int) SectionRea
|
||||||
|
|
||||||
case chunk := <-chunkC:
|
case chunk := <-chunkC:
|
||||||
i++
|
i++
|
||||||
dpaLogger.DebugDetailf("TESTER: chunk request %x", chunk.Key[:4])
|
// dpaLogger.DebugDetailf("TESTER: chunk request %x", chunk.Key[:4])
|
||||||
// this just mocks the behaviour of a chunk store retrieval
|
// this just mocks the behaviour of a chunk store retrieval
|
||||||
var found bool
|
var found bool
|
||||||
for _, ch := range self.chunks {
|
for _, ch := range self.chunks {
|
||||||
|
|
@ -108,10 +109,10 @@ func (self *chunkerTester) Join(chunker *TreeChunker, key Key, c int) SectionRea
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if !found {
|
if !found {
|
||||||
fmt.Printf("TESTER: chunk unknown for %x", chunk.Key[:4])
|
// fmt.Printf("TESTER: chunk unknown for %x", chunk.Key[:4])
|
||||||
}
|
}
|
||||||
close(chunk.C)
|
close(chunk.C)
|
||||||
dpaLogger.DebugDetailf("TESTER: chunk request served %x", chunk.Key[:4])
|
// dpaLogger.DebugDetailf("TESTER: chunk request served %x", chunk.Key[:4])
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}()
|
}()
|
||||||
|
|
@ -129,7 +130,7 @@ func testRandomData(chunker *TreeChunker, tester *chunkerTester, n int, chunks i
|
||||||
reader := tester.Join(chunker, key, 0)
|
reader := tester.Join(chunker, key, 0)
|
||||||
output := make([]byte, n)
|
output := make([]byte, n)
|
||||||
_, err := reader.Read(output)
|
_, err := reader.Read(output)
|
||||||
if err != nil {
|
if err != io.EOF {
|
||||||
t.Errorf("read error %v\n", err)
|
t.Errorf("read error %v\n", err)
|
||||||
}
|
}
|
||||||
t.Logf(" IN: %x\nOUT: %x\n", input, output)
|
t.Logf(" IN: %x\nOUT: %x\n", input, output)
|
||||||
|
|
|
||||||
|
|
@ -46,7 +46,6 @@ SPLIT:
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
if !ok {
|
if !ok {
|
||||||
t.Logf("quitting SPLIT loop\n")
|
|
||||||
break SPLIT
|
break SPLIT
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue