From 75f0248db01458a285fbef29f25a88717395268c Mon Sep 17 00:00:00 2001 From: zelig Date: Fri, 13 Feb 2015 12:29:32 +0100 Subject: [PATCH 01/12] add pull cli flag to start swarm download --- bzz/chunker_test.go | 12 ++++++++++-- cmd/ethereum/flags.go | 2 ++ cmd/ethereum/main.go | 2 +- cmd/utils/cmd.go | 4 ++-- eth/backend.go | 22 ++++++++++++++++++---- 5 files changed, 33 insertions(+), 9 deletions(-) diff --git a/bzz/chunker_test.go b/bzz/chunker_test.go index d533317f0e..5886dfd6f5 100644 --- a/bzz/chunker_test.go +++ b/bzz/chunker_test.go @@ -178,14 +178,22 @@ func readAll(reader SectionReader, result []byte) { } func benchmarkJoinRandomData(n int, chunks int, t *testing.B) { + t.StopTimer() for i := 0; i < t.N; i++ { - t.StopTimer() + // fmt.Printf("round %v\n", i) chunker, tester := chunkerAndTester() - key, _ := tester.Split(chunker, n) + key, slice := tester.Split(chunker, n) + // fmt.Printf("split done %v, joining...\n", i) t.StartTimer() reader := tester.Join(chunker, key, i) + // fmt.Printf("join done %v, reading...\n", i) result := make([]byte, n) readAll(reader, result) + // fmt.Printf("read done %v\n", i) + t.StopTimer() + if !bytes.Equal(slice, result) { + t.Errorf("input output mismatch") + } } } diff --git a/cmd/ethereum/flags.go b/cmd/ethereum/flags.go index d69ea5cfb6..169756e272 100644 --- a/cmd/ethereum/flags.go +++ b/cmd/ethereum/flags.go @@ -66,6 +66,7 @@ var ( Dial bool PrintVersion bool Peers string + Pull string ) // flags specific to cli client @@ -105,6 +106,7 @@ func Init() { flag.BoolVar(&Dial, "dial", true, "dial out connections (on)") flag.BoolVar(&GenAddr, "genaddr", false, "create a new priv/pub key") flag.StringVar(&Peers, "peers", "", "imports the file given (hex or mnemonic formats)") + flag.StringVar(&Peers, "pull", "", "swarm pull key") flag.StringVar(&SecretFile, "import", "", "imports the file given (hex or mnemonic formats)") flag.StringVar(&ExportDir, "export", "", "exports the session keyring to files in the directory given") flag.StringVar(&LogFile, "logfile", "", "log file (defaults to standard output)") diff --git a/cmd/ethereum/main.go b/cmd/ethereum/main.go index ea53b9a304..8a549dbea0 100644 --- a/cmd/ethereum/main.go +++ b/cmd/ethereum/main.go @@ -134,7 +134,7 @@ func main() { utils.StartWebSockets(ethereum) } - utils.StartEthereum(ethereum, UseSeed, Peers) + utils.StartEthereum(ethereum, UseSeed, Peers, Pull) if StartJsConsole { InitJsConsole(ethereum) diff --git a/cmd/utils/cmd.go b/cmd/utils/cmd.go index 06e3fb93f1..f6e7999e75 100644 --- a/cmd/utils/cmd.go +++ b/cmd/utils/cmd.go @@ -120,9 +120,9 @@ func exit(err error) { os.Exit(status) } -func StartEthereum(ethereum *eth.Ethereum, UseSeed bool, Peers string) { +func StartEthereum(ethereum *eth.Ethereum, UseSeed bool, Peers string, Pull string) { clilogger.Infof("Starting %s", ethereum.ClientIdentity()) - err := ethereum.Start(UseSeed, Peers) + err := ethereum.Start(UseSeed, Peers, Pull) if err != nil { exit(err) } diff --git a/eth/backend.go b/eth/backend.go index f5163bdfd5..535444b02a 100644 --- a/eth/backend.go +++ b/eth/backend.go @@ -2,7 +2,9 @@ package eth import ( "fmt" + "io" "net" + "os" "strings" "sync" @@ -70,6 +72,7 @@ type Ethereum struct { RpcServer *rpc.JsonRpcServer keyManager *crypto.KeyManager + dpa *bzz.DPA clientIdentity p2p.ClientIdentity logger ethlogger.LogSystem @@ -144,12 +147,12 @@ func New(config *Config) (*Ethereum, error) { } chunker := &bzz.TreeChunker{} chunker.Init() - dpa := &bzz.DPA{ + eth.dpa = &bzz.DPA{ Chunker: chunker, ChunkStore: netStore, } - dpa.Start() - go bzz.StartHttpServer(dpa) + eth.dpa.Start() + go bzz.StartHttpServer(eth.dpa) nat, err := p2p.ParseNAT(config.NATType, config.PMPGateway) if err != nil { @@ -234,7 +237,7 @@ func (s *Ethereum) MaxPeers() int { } // Start the ethereum -func (s *Ethereum) Start(seed bool, p string) error { +func (s *Ethereum) Start(seed bool, p string, pull string) error { err := s.net.Start() if err != nil { return err @@ -264,6 +267,17 @@ func (s *Ethereum) Start(seed bool, p string) error { } } + if len(pull) > 0 { + key := make([]byte, s.dpa.Chunker.KeySize()) + reader := s.dpa.Retrieve(key) + fo, err := os.Open("/tmp/swarm.tmp") + if err != nil { + logger.Warnf("file open error %v", err) + } else { + io.Copy(fo, reader) + } + } + // TODO: read peers here if seed { logger.Infof("Connect to seed node %v", seedNodeAddress) From 706f3eddad1d33129047d5ee95c93327c9aaf1eb Mon Sep 17 00:00:00 2001 From: zelig Date: Fri, 13 Feb 2015 12:29:32 +0100 Subject: [PATCH 02/12] add pull cli flag to start swarm download --- bzz/chunker_test.go | 12 ++++++++++-- cmd/ethereum/flags.go | 2 ++ cmd/ethereum/main.go | 2 +- cmd/utils/cmd.go | 4 ++-- eth/backend.go | 21 +++++++++++++++++---- 5 files changed, 32 insertions(+), 9 deletions(-) diff --git a/bzz/chunker_test.go b/bzz/chunker_test.go index d533317f0e..5886dfd6f5 100644 --- a/bzz/chunker_test.go +++ b/bzz/chunker_test.go @@ -178,14 +178,22 @@ func readAll(reader SectionReader, result []byte) { } func benchmarkJoinRandomData(n int, chunks int, t *testing.B) { + t.StopTimer() for i := 0; i < t.N; i++ { - t.StopTimer() + // fmt.Printf("round %v\n", i) chunker, tester := chunkerAndTester() - key, _ := tester.Split(chunker, n) + key, slice := tester.Split(chunker, n) + // fmt.Printf("split done %v, joining...\n", i) t.StartTimer() reader := tester.Join(chunker, key, i) + // fmt.Printf("join done %v, reading...\n", i) result := make([]byte, n) readAll(reader, result) + // fmt.Printf("read done %v\n", i) + t.StopTimer() + if !bytes.Equal(slice, result) { + t.Errorf("input output mismatch") + } } } diff --git a/cmd/ethereum/flags.go b/cmd/ethereum/flags.go index d69ea5cfb6..b3be4de885 100644 --- a/cmd/ethereum/flags.go +++ b/cmd/ethereum/flags.go @@ -66,6 +66,7 @@ var ( Dial bool PrintVersion bool Peers string + Pull string ) // flags specific to cli client @@ -105,6 +106,7 @@ func Init() { flag.BoolVar(&Dial, "dial", true, "dial out connections (on)") flag.BoolVar(&GenAddr, "genaddr", false, "create a new priv/pub key") flag.StringVar(&Peers, "peers", "", "imports the file given (hex or mnemonic formats)") + flag.StringVar(&Pull, "pull", "", "swarm pull key") flag.StringVar(&SecretFile, "import", "", "imports the file given (hex or mnemonic formats)") flag.StringVar(&ExportDir, "export", "", "exports the session keyring to files in the directory given") flag.StringVar(&LogFile, "logfile", "", "log file (defaults to standard output)") diff --git a/cmd/ethereum/main.go b/cmd/ethereum/main.go index ea53b9a304..8a549dbea0 100644 --- a/cmd/ethereum/main.go +++ b/cmd/ethereum/main.go @@ -134,7 +134,7 @@ func main() { utils.StartWebSockets(ethereum) } - utils.StartEthereum(ethereum, UseSeed, Peers) + utils.StartEthereum(ethereum, UseSeed, Peers, Pull) if StartJsConsole { InitJsConsole(ethereum) diff --git a/cmd/utils/cmd.go b/cmd/utils/cmd.go index 06e3fb93f1..f6e7999e75 100644 --- a/cmd/utils/cmd.go +++ b/cmd/utils/cmd.go @@ -120,9 +120,9 @@ func exit(err error) { os.Exit(status) } -func StartEthereum(ethereum *eth.Ethereum, UseSeed bool, Peers string) { +func StartEthereum(ethereum *eth.Ethereum, UseSeed bool, Peers string, Pull string) { clilogger.Infof("Starting %s", ethereum.ClientIdentity()) - err := ethereum.Start(UseSeed, Peers) + err := ethereum.Start(UseSeed, Peers, Pull) if err != nil { exit(err) } diff --git a/eth/backend.go b/eth/backend.go index f5163bdfd5..26a8f144d3 100644 --- a/eth/backend.go +++ b/eth/backend.go @@ -2,7 +2,9 @@ package eth import ( "fmt" + "io" "net" + "os" "strings" "sync" @@ -70,6 +72,7 @@ type Ethereum struct { RpcServer *rpc.JsonRpcServer keyManager *crypto.KeyManager + dpa *bzz.DPA clientIdentity p2p.ClientIdentity logger ethlogger.LogSystem @@ -144,12 +147,12 @@ func New(config *Config) (*Ethereum, error) { } chunker := &bzz.TreeChunker{} chunker.Init() - dpa := &bzz.DPA{ + eth.dpa = &bzz.DPA{ Chunker: chunker, ChunkStore: netStore, } - dpa.Start() - go bzz.StartHttpServer(dpa) + eth.dpa.Start() + go bzz.StartHttpServer(eth.dpa) nat, err := p2p.ParseNAT(config.NATType, config.PMPGateway) if err != nil { @@ -234,7 +237,7 @@ func (s *Ethereum) MaxPeers() int { } // Start the ethereum -func (s *Ethereum) Start(seed bool, p string) error { +func (s *Ethereum) Start(seed bool, p string, pull string) error { err := s.net.Start() if err != nil { return err @@ -264,6 +267,16 @@ func (s *Ethereum) Start(seed bool, p string) error { } } + if len(pull) > 0 { + reader := s.dpa.Retrieve(pull) + fo, err := os.Open("/tmp/swarm.tmp") + if err != nil { + logger.Warnf("file open error %v", err) + } else { + io.Copy(fo, reader) + } + } + // TODO: read peers here if seed { logger.Infof("Connect to seed node %v", seedNodeAddress) From e2fe4748567eee2418d91a737bc02e02b20f0719 Mon Sep 17 00:00:00 2001 From: zelig Date: Fri, 13 Feb 2015 12:56:18 +0100 Subject: [PATCH 03/12] fix open file for swarm pull --- eth/backend.go | 11 ++++++++--- 1 file changed, 8 insertions(+), 3 deletions(-) diff --git a/eth/backend.go b/eth/backend.go index b4f41ff8b6..4bc201300a 100644 --- a/eth/backend.go +++ b/eth/backend.go @@ -268,12 +268,17 @@ func (s *Ethereum) Start(seed bool, p string, pull string) error { } if len(pull) > 0 { - reader := s.dpa.Retrieve([]byte(pull)) - fo, err := os.Open("/tmp/swarm.tmp") + key := []byte(pull) + reader := s.dpa.Retrieve(key) + logger.Debugf("retrieved reader for %064x", key) + fo, err := os.OpenFile("/tmp/swarm.tmp", os.O_CREATE|os.O_RDWR, 0666) if err != nil { logger.Warnf("file open error %v", err) } else { - io.Copy(fo, reader) + n, err := io.Copy(fo, reader) + if err != nil && err != io.EOF { + logger.Debugf("read %v bytes. read error for %064x: %v", n, key, err) + } } } From 33126ae7a21863ede216080ad8e644c3fe7488bf Mon Sep 17 00:00:00 2001 From: zelig Date: Fri, 13 Feb 2015 12:59:50 +0100 Subject: [PATCH 04/12] fix hex2bytes in backend --- eth/backend.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/eth/backend.go b/eth/backend.go index 4bc201300a..cd6560a315 100644 --- a/eth/backend.go +++ b/eth/backend.go @@ -268,7 +268,7 @@ func (s *Ethereum) Start(seed bool, p string, pull string) error { } if len(pull) > 0 { - key := []byte(pull) + key := ethutil.Hex2Bytes(pull) reader := s.dpa.Retrieve(key) logger.Debugf("retrieved reader for %064x", key) fo, err := os.OpenFile("/tmp/swarm.tmp", os.O_CREATE|os.O_RDWR, 0666) From 31ed59d20b57a0875b6a1684cc6688c6bca8a1f7 Mon Sep 17 00:00:00 2001 From: zelig Date: Fri, 13 Feb 2015 13:10:52 +0100 Subject: [PATCH 05/12] sleep before launch swarm pull and no file written out --- eth/backend.go | 31 +++++++++++++++++++------------ 1 file changed, 19 insertions(+), 12 deletions(-) diff --git a/eth/backend.go b/eth/backend.go index cd6560a315..4c6cd22d06 100644 --- a/eth/backend.go +++ b/eth/backend.go @@ -4,9 +4,9 @@ import ( "fmt" "io" "net" - "os" "strings" "sync" + "time" "github.com/ethereum/go-ethereum/bzz" "github.com/ethereum/go-ethereum/core" @@ -268,18 +268,25 @@ func (s *Ethereum) Start(seed bool, p string, pull string) error { } if len(pull) > 0 { - key := ethutil.Hex2Bytes(pull) - reader := s.dpa.Retrieve(key) - logger.Debugf("retrieved reader for %064x", key) - fo, err := os.OpenFile("/tmp/swarm.tmp", os.O_CREATE|os.O_RDWR, 0666) - if err != nil { - logger.Warnf("file open error %v", err) - } else { - n, err := io.Copy(fo, reader) - if err != nil && err != io.EOF { - logger.Debugf("read %v bytes. read error for %064x: %v", n, key, err) + go func() { + time.Sleep(30 * time.Second) + key := ethutil.Hex2Bytes(pull) + reader := s.dpa.Retrieve(key) + logger.Debugf("retrieved reader for %064x", key) + b := make([]byte, 0x1000) + var length int64 + for { + n, err := reader.Read(b) + if err != nil { + if err != io.EOF { + logger.Debugf("read %v bytes. read error for %064x: %v", n, key, err) + } + return + } + length += int64(n) } - } + logger.Debugf("read %v bytes from %064x: %v", length, key) + }() } // TODO: read peers here From b8bd17de28b3ce837e5883a52a4e546d8a0b9219 Mon Sep 17 00:00:00 2001 From: zelig Date: Fri, 13 Feb 2015 13:21:13 +0100 Subject: [PATCH 06/12] get rid of ping message logging --- p2p/message.go | 13 +++++++++---- 1 file changed, 9 insertions(+), 4 deletions(-) diff --git a/p2p/message.go b/p2p/message.go index 7eaa884069..cabd1940bd 100644 --- a/p2p/message.go +++ b/p2p/message.go @@ -102,17 +102,22 @@ func writeMsg(w io.Writer, msg Msg) error { copy(start, magicToken) binary.BigEndian.PutUint32(start[4:], payloadLen) - srvlog.Debugf("Sending message:") - + if msg.Size > 1 { + srvlog.Debugf("Sending message (size %v):", msg.Size) + } for _, b := range [][]byte{start, listhdr, code} { - srvlog.Debugf(" %x", b) + if msg.Size > 1 { + srvlog.Debugf(" %x", b) + } if _, err := w.Write(b); err != nil { return err } } b := make([]byte, msg.Size) msg.Payload.Read(b) - srvlog.Debugf(" %x", b) + if msg.Size > 1 { + srvlog.Debugf(" %x", b) + } _, err := w.Write(b) return err } From cd9c7bea60197ed74b88c975b269681a929aac3c Mon Sep 17 00:00:00 2001 From: "Daniel A. Nagy" Date: Fri, 13 Feb 2015 13:29:46 +0100 Subject: [PATCH 07/12] Additional error handling. --- bzz/httpaccess.go | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/bzz/httpaccess.go b/bzz/httpaccess.go index 9563713d7d..d5c3f8027d 100644 --- a/bzz/httpaccess.go +++ b/bzz/httpaccess.go @@ -136,7 +136,12 @@ func handler(w http.ResponseWriter, r *http.Request, dpa *DPA) { size, err := manifestReader.Read(manifest) if int64(size) < manifestReader.Size() { dpaLogger.Debugf("Swarm: Manifest %s not found.", name) - http.Error(w, err.Error(), http.StatusNotFound) + if err == nil { + http.Error(w, "Manifest retrieval cut short: "+string(size)+" "+string(manifestReader.Size()), + http.StatusNotFound) + } else { + http.Error(w, err.Error(), http.StatusNotFound) + } return } dpaLogger.Debugf("Swarm: Manifest %s retrieved.", name) From 1797997697487251ae91880aceb8cd6b146d9750 Mon Sep 17 00:00:00 2001 From: "Daniel A. Nagy" Date: Fri, 13 Feb 2015 13:30:13 +0100 Subject: [PATCH 08/12] Error message fixed. --- bzz/httpaccess.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/bzz/httpaccess.go b/bzz/httpaccess.go index d5c3f8027d..1b95b0c335 100644 --- a/bzz/httpaccess.go +++ b/bzz/httpaccess.go @@ -137,7 +137,7 @@ func handler(w http.ResponseWriter, r *http.Request, dpa *DPA) { if int64(size) < manifestReader.Size() { dpaLogger.Debugf("Swarm: Manifest %s not found.", name) if err == nil { - http.Error(w, "Manifest retrieval cut short: "+string(size)+" "+string(manifestReader.Size()), + http.Error(w, "Manifest retrieval cut short: "+string(size)+"<"+string(manifestReader.Size()), http.StatusNotFound) } else { http.Error(w, err.Error(), http.StatusNotFound) From 5e56f338d4eeb1960c9ebf604eeb1e4c864ee1f3 Mon Sep 17 00:00:00 2001 From: zelig Date: Fri, 13 Feb 2015 13:38:24 +0100 Subject: [PATCH 09/12] improve join benchmarks --- bzz/chunker_test.go | 18 ++++++++++-------- 1 file changed, 10 insertions(+), 8 deletions(-) diff --git a/bzz/chunker_test.go b/bzz/chunker_test.go index 5886dfd6f5..cf0d210571 100644 --- a/bzz/chunker_test.go +++ b/bzz/chunker_test.go @@ -177,23 +177,25 @@ func readAll(reader SectionReader, result []byte) { } } +func benchReadAll(reader SectionReader) { + size := reader.Size() + output := make([]byte, 1000) + for pos := int64(0); pos < size; pos += 1000 { + reader.ReadAt(output, pos) + } +} + func benchmarkJoinRandomData(n int, chunks int, t *testing.B) { t.StopTimer() for i := 0; i < t.N; i++ { // fmt.Printf("round %v\n", i) chunker, tester := chunkerAndTester() - key, slice := tester.Split(chunker, n) + key, _ := tester.Split(chunker, n) // fmt.Printf("split done %v, joining...\n", i) t.StartTimer() reader := tester.Join(chunker, key, i) // fmt.Printf("join done %v, reading...\n", i) - result := make([]byte, n) - readAll(reader, result) - // fmt.Printf("read done %v\n", i) - t.StopTimer() - if !bytes.Equal(slice, result) { - t.Errorf("input output mismatch") - } + benchReadAll(reader) } } From e574002c3ef5b264e2151b6ab4947fa5a3d5de4e Mon Sep 17 00:00:00 2001 From: zelig Date: Fri, 13 Feb 2015 13:49:23 +0100 Subject: [PATCH 10/12] add log to netstore.put --- bzz/netstore.go | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/bzz/netstore.go b/bzz/netstore.go index edb37bce96..0587063cc2 100644 --- a/bzz/netstore.go +++ b/bzz/netstore.go @@ -69,8 +69,9 @@ func (self *NetStore) Put(entry *Chunk) { func (self *NetStore) put(entry *Chunk) { self.localStore.Put(entry) dpaLogger.Debugf("NetStore.put: localStore.Put of %064x completed.", entry.Key) - self.store(entry) + go self.store(entry) // only send responses once + dpaLogger.Debugf("NetStore.put: req: %#v", entry.Key) if entry.req != nil && entry.req.status == reqSearching { entry.req.status = reqFound close(entry.req.C) @@ -157,8 +158,6 @@ func (self *NetStore) addRetrieveRequest(req *retrieveRequestMsgData) { send, timeout := self.strategyUpdateRequest(chunk.req, req) // may change req status - dpaLogger.Debugf("Is %v == %v?", send, storeRequestMsg) - if send == storeRequestMsg { dpaLogger.Debugf("NetStore.addRetrieveRequest: %064x - content found, delivering...", req.Key) self.deliver(req, chunk) From e787f4c88549872164615f9427a68d65da679bf2 Mon Sep 17 00:00:00 2001 From: "Daniel A. Nagy" Date: Fri, 13 Feb 2015 13:58:15 +0100 Subject: [PATCH 11/12] req logged in putback --- bzz/netstore.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/bzz/netstore.go b/bzz/netstore.go index 0587063cc2..d4c3e6fe37 100644 --- a/bzz/netstore.go +++ b/bzz/netstore.go @@ -71,7 +71,7 @@ func (self *NetStore) put(entry *Chunk) { dpaLogger.Debugf("NetStore.put: localStore.Put of %064x completed.", entry.Key) go self.store(entry) // only send responses once - dpaLogger.Debugf("NetStore.put: req: %#v", entry.Key) + dpaLogger.Debugf("NetStore.put: req: %#v", entry.req) if entry.req != nil && entry.req.status == reqSearching { entry.req.status = reqFound close(entry.req.C) From b0d7fa2c7a5523f37b824f73fee7b440460b7792 Mon Sep 17 00:00:00 2001 From: "Daniel A. Nagy" Date: Fri, 13 Feb 2015 16:01:53 +0100 Subject: [PATCH 12/12] memStore consistently returns the same Chunk object. --- bzz/chunker.go | 2 +- bzz/memstore.go | 19 +++++++++---------- bzz/netstore.go | 8 ++++++-- 3 files changed, 16 insertions(+), 13 deletions(-) diff --git a/bzz/chunker.go b/bzz/chunker.go index 1838242d94..5d543c68bc 100644 --- a/bzz/chunker.go +++ b/bzz/chunker.go @@ -318,7 +318,7 @@ func (self *LazyChunkReader) ReadAt(b []byte, off int64) (read int, err error) { // dpaLogger.Debugf("chunk data received for %x", chunk.Key[:4]) } if len(chunk.Data) == 0 { - // dpaLogger.Debugf("No payload.") + dpaLogger.Debugf("No payload in %x.", chunk.Key) return 0, notFound } self.size = chunk.Size diff --git a/bzz/memstore.go b/bzz/memstore.go index ce012222d4..0fab5462e1 100644 --- a/bzz/memstore.go +++ b/bzz/memstore.go @@ -159,6 +159,7 @@ func (s *memStore) getEntryCnt() uint { } +// entry (not its copy) is going to be in memStore func (s *memStore) Put(entry *Chunk) { if s.capacity == 0 { @@ -193,13 +194,15 @@ func (s *memStore) Put(entry *Chunk) { if node.entry.Key.isEqual(entry.Key) { node.updateAccess(s.accessCnt) - if node.entry.Data == nil { - node.entry.Size = entry.Size - node.entry.Data = entry.Data + if entry.Data == nil { + entry.Size = node.entry.Size + entry.Data = node.entry.Data } - if node.entry.req == nil { - node.entry.req = entry.req + if entry.req == nil { + entry.req = node.entry.req } + entry.C = node.entry.C + node.entry = entry return } @@ -253,11 +256,7 @@ func (s *memStore) Get(hash Key) (chunk *Chunk, err error) { if node.entry.Key.isEqual(hash) { s.accessCnt++ node.updateAccess(s.accessCnt) - chunk = &Chunk{ - Key: hash, - Data: node.entry.Data, - Size: node.entry.Size, - } + chunk = node.entry if s.dbAccessCnt-node.lastDBaccess > dbForceUpdateAccessCnt { s.dbAccessCnt++ node.lastDBaccess = s.dbAccessCnt diff --git a/bzz/netstore.go b/bzz/netstore.go index d4c3e6fe37..65381b091c 100644 --- a/bzz/netstore.go +++ b/bzz/netstore.go @@ -68,7 +68,7 @@ func (self *NetStore) Put(entry *Chunk) { func (self *NetStore) put(entry *Chunk) { self.localStore.Put(entry) - dpaLogger.Debugf("NetStore.put: localStore.Put of %064x completed.", entry.Key) + dpaLogger.Debugf("NetStore.put: localStore.Put of %064x completed, %d bytes (%p).", entry.Key, len(entry.Data), entry) go self.store(entry) // only send responses once dpaLogger.Debugf("NetStore.put: req: %#v", entry.req) @@ -82,7 +82,9 @@ func (self *NetStore) put(entry *Chunk) { func (self *NetStore) addStoreRequest(req *storeRequestMsgData) { self.lock.Lock() defer self.lock.Unlock() + dpaLogger.Debugf("NetStore.addStoreRequest: req = %#v", req) chunk, err := self.localStore.Get(req.Key) + dpaLogger.Debugf("NetStore.addStoreRequest: chunk reference %p", chunk) // we assume that a returned chunk is the one stored in the memory cache if err != nil { chunk = &Chunk{ @@ -91,6 +93,7 @@ func (self *NetStore) addStoreRequest(req *storeRequestMsgData) { Size: int64(req.Size), } } else if chunk.Data == nil { + // response to a search request chunk.Data = req.Data chunk.Size = int64(req.Size) } else { @@ -116,7 +119,7 @@ func (self *NetStore) Get(key Key) (chunk *Chunk, err error) { dpaLogger.Debugf("NetStore.Get: %064x request time out ", key) err = notFound case <-chunk.req.C: - dpaLogger.Debugf("NetStore.get: %064x retrieved", key) + dpaLogger.Debugf("NetStore.Get: %064x retrieved, %d bytes (%p)", key, len(chunk.Data), chunk) } return @@ -137,6 +140,7 @@ func (self *NetStore) get(key Key) (chunk *Chunk) { if chunk.req == nil { chunk.req = new(requestStatus) + chunk.req.C = make(chan bool) } return }