mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-07-21 20:26:41 +00:00
DPA initial test passes
- allocate byte slice for root key in dpa.Store - close chunk.C in retrieveloops even if chunk not found (fixes hanging) in both common_test and dpa - make dpa.Store synchronous - block until chunker finishes with root key
This commit is contained in:
parent
22249ac2b9
commit
d82157ec86
3 changed files with 42 additions and 50 deletions
|
|
@ -72,15 +72,13 @@ SPLIT:
|
||||||
go func() {
|
go func() {
|
||||||
storedChunk, err := m.Get(chunk.Key)
|
storedChunk, err := m.Get(chunk.Key)
|
||||||
if err == notFound {
|
if err == notFound {
|
||||||
t.Errorf("Chunk not found: %v", err)
|
dpaLogger.DebugDetailf("chunk '%x' not found", chunk.Key)
|
||||||
return
|
} else if err != nil {
|
||||||
|
dpaLogger.DebugDetailf("error retrieving chunk %x: %v", chunk.Key, err)
|
||||||
|
} else {
|
||||||
|
chunk.Reader = NewChunkReaderFromBytes(storedChunk.Data)
|
||||||
|
chunk.Size = storedChunk.Size
|
||||||
}
|
}
|
||||||
if err != nil {
|
|
||||||
t.Errorf("GET error: %v", err)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
chunk.Reader = NewChunkReaderFromBytes(storedChunk.Data)
|
|
||||||
chunk.Size = storedChunk.Size
|
|
||||||
close(chunk.C)
|
close(chunk.C)
|
||||||
}()
|
}()
|
||||||
case err, ok := <-errC:
|
case err, ok := <-errC:
|
||||||
|
|
|
||||||
70
bzz/dpa.go
70
bzz/dpa.go
|
|
@ -4,10 +4,9 @@ import (
|
||||||
"errors"
|
"errors"
|
||||||
"sync"
|
"sync"
|
||||||
// "time"
|
// "time"
|
||||||
"fmt"
|
// "fmt"
|
||||||
|
|
||||||
ethlogger "github.com/ethereum/go-ethereum/logger"
|
ethlogger "github.com/ethereum/go-ethereum/logger"
|
||||||
// "github.com/ethereum/go-ethereum/rlp"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
/*
|
/*
|
||||||
|
|
@ -70,50 +69,49 @@ type ChunkStore interface {
|
||||||
Get(Key) (*Chunk, error)
|
Get(Key) (*Chunk, error)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (self *DPA) Retrieve(key Key) (data LazySectionReader) {
|
func (self *DPA) Retrieve(key Key) LazySectionReader {
|
||||||
|
|
||||||
reader, errC := self.Chunker.Join(key, self.retrieveC)
|
reader, errC := self.Chunker.Join(key, self.retrieveC)
|
||||||
data = reader
|
|
||||||
// we can add subscriptions etc. or timeout here
|
// we can add subscriptions etc. or timeout here
|
||||||
go func() {
|
go func() {
|
||||||
LOOP:
|
JOIN:
|
||||||
for {
|
for {
|
||||||
select {
|
select {
|
||||||
case err, ok := <-errC:
|
case err, ok := <-errC:
|
||||||
if err != nil {
|
if err != nil {
|
||||||
dpaLogger.Warnf("%v", err)
|
dpaLogger.Warnf("chunker join error: %v", err)
|
||||||
}
|
}
|
||||||
if !ok {
|
if !ok {
|
||||||
break LOOP
|
break JOIN
|
||||||
}
|
}
|
||||||
case <-self.quitC:
|
case <-self.quitC:
|
||||||
return
|
break JOIN
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}()
|
}()
|
||||||
|
|
||||||
return
|
return reader
|
||||||
}
|
}
|
||||||
|
|
||||||
func (self *DPA) Store(data SectionReader) (key Key, err error) {
|
func (self *DPA) Store(data SectionReader) (key Key, err error) {
|
||||||
|
key = make([]byte, self.Chunker.KeySize())
|
||||||
errC := self.Chunker.Split(key, data, self.storeC)
|
errC := self.Chunker.Split(key, data, self.storeC)
|
||||||
|
|
||||||
go func() {
|
SPLIT:
|
||||||
LOOP:
|
for {
|
||||||
for {
|
select {
|
||||||
select {
|
case err, ok := <-errC:
|
||||||
case err, ok := <-errC:
|
if err != nil {
|
||||||
dpaLogger.Warnf("%v", err)
|
dpaLogger.Warnf("chunkner split error: %v", err)
|
||||||
if !ok {
|
|
||||||
break LOOP
|
|
||||||
}
|
|
||||||
|
|
||||||
case <-self.quitC:
|
|
||||||
break LOOP
|
|
||||||
}
|
}
|
||||||
|
if !ok {
|
||||||
|
break SPLIT
|
||||||
|
}
|
||||||
|
|
||||||
|
case <-self.quitC:
|
||||||
|
break SPLIT
|
||||||
}
|
}
|
||||||
}()
|
}
|
||||||
return
|
return
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
@ -146,18 +144,17 @@ func (self *DPA) retrieveLoop() {
|
||||||
go func() {
|
go func() {
|
||||||
RETRIEVE:
|
RETRIEVE:
|
||||||
for chunk := range self.retrieveC {
|
for chunk := range self.retrieveC {
|
||||||
|
|
||||||
go func() {
|
go func() {
|
||||||
storedChunk, err := self.ChunkStore.Get(chunk.Key)
|
storedChunk, err := self.ChunkStore.Get(chunk.Key)
|
||||||
if err == notFound {
|
if err == notFound {
|
||||||
dpaLogger.DebugDetailf("chunk %x not found", chunk.Key)
|
dpaLogger.DebugDetailf("chunk '%x' not found", chunk.Key)
|
||||||
return
|
} else if err != nil {
|
||||||
}
|
|
||||||
if err != nil {
|
|
||||||
dpaLogger.DebugDetailf("error retrieving chunk %x: %v", chunk.Key, err)
|
dpaLogger.DebugDetailf("error retrieving chunk %x: %v", chunk.Key, err)
|
||||||
return
|
} else {
|
||||||
|
chunk.Reader = NewChunkReaderFromBytes(storedChunk.Data)
|
||||||
|
chunk.Size = storedChunk.Size
|
||||||
}
|
}
|
||||||
chunk.Reader = NewChunkReaderFromBytes(storedChunk.Data)
|
|
||||||
chunk.Size = storedChunk.Size
|
|
||||||
close(chunk.C)
|
close(chunk.C)
|
||||||
}()
|
}()
|
||||||
select {
|
select {
|
||||||
|
|
@ -172,18 +169,15 @@ func (self *DPA) retrieveLoop() {
|
||||||
func (self *DPA) storeLoop() {
|
func (self *DPA) storeLoop() {
|
||||||
self.storeC = make(chan *Chunk)
|
self.storeC = make(chan *Chunk)
|
||||||
go func() {
|
go func() {
|
||||||
fmt.Printf("StoreLoop started.\n")
|
|
||||||
STORE:
|
STORE:
|
||||||
for {
|
for chunk := range self.storeC {
|
||||||
|
chunk.Data = make([]byte, chunk.Reader.Size())
|
||||||
|
chunk.Reader.ReadAt(chunk.Data, 0)
|
||||||
|
self.ChunkStore.Put(chunk)
|
||||||
select {
|
select {
|
||||||
case chunk := <-self.storeC:
|
|
||||||
fmt.Printf("StoreLoop reader size %d\n", chunk.Reader.Size())
|
|
||||||
chunk.Data = make([]byte, chunk.Reader.Size())
|
|
||||||
chunk.Reader.ReadAt(chunk.Data, 0)
|
|
||||||
self.ChunkStore.Put(chunk)
|
|
||||||
case <-self.quitC:
|
case <-self.quitC:
|
||||||
break STORE
|
break STORE
|
||||||
// default:
|
default:
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}()
|
}()
|
||||||
|
|
|
||||||
|
|
@ -6,7 +6,7 @@ import (
|
||||||
"github.com/ethereum/go-ethereum/bzz/test"
|
"github.com/ethereum/go-ethereum/bzz/test"
|
||||||
"os"
|
"os"
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
// "time"
|
||||||
)
|
)
|
||||||
|
|
||||||
func TestDPA(t *testing.T) {
|
func TestDPA(t *testing.T) {
|
||||||
|
|
@ -31,13 +31,13 @@ func TestDPA(t *testing.T) {
|
||||||
reader, slice := testDataReader(0x100)
|
reader, slice := testDataReader(0x100)
|
||||||
fmt.Printf("Chunk size: %d.", len(slice))
|
fmt.Printf("Chunk size: %d.", len(slice))
|
||||||
key, err := dpa.Store(reader)
|
key, err := dpa.Store(reader)
|
||||||
// _, err = dpa.Store(reader)
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Errorf("Store error: %v", err)
|
t.Errorf("Store error: %v", err)
|
||||||
}
|
}
|
||||||
|
// time.Sleep(2 * time.Second)
|
||||||
resultReader := dpa.Retrieve(key)
|
resultReader := dpa.Retrieve(key)
|
||||||
resultSlice := make([]byte, len(slice))
|
resultSlice := make([]byte, len(slice))
|
||||||
n, err := resultReader.Read(resultSlice)
|
n, err := resultReader.ReadAt(resultSlice, 0)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Errorf("Retrieve error: %v", err)
|
t.Errorf("Retrieve error: %v", err)
|
||||||
}
|
}
|
||||||
|
|
@ -47,5 +47,5 @@ func TestDPA(t *testing.T) {
|
||||||
if !bytes.Equal(slice, resultSlice) {
|
if !bytes.Equal(slice, resultSlice) {
|
||||||
t.Errorf("Comparison error.")
|
t.Errorf("Comparison error.")
|
||||||
}
|
}
|
||||||
time.Sleep(time.Second)
|
// time.Sleep(time.Second)
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue