diff --git a/bzz/common_test.go b/bzz/common_test.go index b1931d5589..f39c2442c1 100644 --- a/bzz/common_test.go +++ b/bzz/common_test.go @@ -72,15 +72,13 @@ SPLIT: go func() { storedChunk, err := m.Get(chunk.Key) if err == notFound { - t.Errorf("Chunk not found: %v", err) - return + 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 } - if err != nil { - t.Errorf("GET error: %v", err) - return - } - chunk.Reader = NewChunkReaderFromBytes(storedChunk.Data) - chunk.Size = storedChunk.Size close(chunk.C) }() case err, ok := <-errC: diff --git a/bzz/dpa.go b/bzz/dpa.go index 496919a4f6..8a0a416740 100644 --- a/bzz/dpa.go +++ b/bzz/dpa.go @@ -4,10 +4,9 @@ import ( "errors" "sync" // "time" - "fmt" + // "fmt" ethlogger "github.com/ethereum/go-ethereum/logger" - // "github.com/ethereum/go-ethereum/rlp" ) /* @@ -70,50 +69,49 @@ type ChunkStore interface { 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) - data = reader // we can add subscriptions etc. or timeout here go func() { - LOOP: + JOIN: for { select { case err, ok := <-errC: if err != nil { - dpaLogger.Warnf("%v", err) + dpaLogger.Warnf("chunker join error: %v", err) } if !ok { - break LOOP + break JOIN } case <-self.quitC: - return + break JOIN } } }() - return + return reader } func (self *DPA) Store(data SectionReader) (key Key, err error) { - + key = make([]byte, self.Chunker.KeySize()) errC := self.Chunker.Split(key, data, self.storeC) - go func() { - LOOP: - for { - select { - case err, ok := <-errC: - dpaLogger.Warnf("%v", err) - if !ok { - break LOOP - } - - case <-self.quitC: - break LOOP +SPLIT: + for { + select { + case err, ok := <-errC: + if err != nil { + dpaLogger.Warnf("chunkner split error: %v", err) } + if !ok { + break SPLIT + } + + case <-self.quitC: + break SPLIT } - }() + } return } @@ -146,18 +144,17 @@ func (self *DPA) retrieveLoop() { go func() { RETRIEVE: for chunk := range self.retrieveC { + go func() { storedChunk, err := self.ChunkStore.Get(chunk.Key) if err == notFound { - dpaLogger.DebugDetailf("chunk %x not found", chunk.Key) - return - } - if err != nil { + dpaLogger.DebugDetailf("chunk '%x' not found", chunk.Key) + } else if err != nil { 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) }() select { @@ -172,18 +169,15 @@ func (self *DPA) retrieveLoop() { func (self *DPA) storeLoop() { self.storeC = make(chan *Chunk) go func() { - fmt.Printf("StoreLoop started.\n") 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 { - 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: break STORE - // default: + default: } } }() diff --git a/bzz/dpa_test.go b/bzz/dpa_test.go index 10cfd0773a..a514654124 100644 --- a/bzz/dpa_test.go +++ b/bzz/dpa_test.go @@ -6,7 +6,7 @@ import ( "github.com/ethereum/go-ethereum/bzz/test" "os" "testing" - "time" + // "time" ) func TestDPA(t *testing.T) { @@ -31,13 +31,13 @@ func TestDPA(t *testing.T) { reader, slice := testDataReader(0x100) fmt.Printf("Chunk size: %d.", len(slice)) key, err := dpa.Store(reader) - // _, err = dpa.Store(reader) if err != nil { t.Errorf("Store error: %v", err) } + // time.Sleep(2 * time.Second) resultReader := dpa.Retrieve(key) resultSlice := make([]byte, len(slice)) - n, err := resultReader.Read(resultSlice) + n, err := resultReader.ReadAt(resultSlice, 0) if err != nil { t.Errorf("Retrieve error: %v", err) } @@ -47,5 +47,5 @@ func TestDPA(t *testing.T) { if !bytes.Equal(slice, resultSlice) { t.Errorf("Comparison error.") } - time.Sleep(time.Second) + // time.Sleep(time.Second) }