From f3fdcb2064d8c48488a6b02492db3277c0680773 Mon Sep 17 00:00:00 2001 From: Balint Gabor Date: Thu, 4 Jan 2018 12:59:04 +0100 Subject: [PATCH] Fis storageWG in chunker storeWG did not work because it was possible to start to wait on earlier than the first Add happened --- swarm/storage/chunker.go | 32 ++++++++++++++++++-------------- 1 file changed, 18 insertions(+), 14 deletions(-) diff --git a/swarm/storage/chunker.go b/swarm/storage/chunker.go index dde975c1a0..be142f227b 100644 --- a/swarm/storage/chunker.go +++ b/swarm/storage/chunker.go @@ -130,7 +130,7 @@ func (self *TreeChunker) Split(data io.Reader, size int64, chunkC chan *Chunk) ( quitC := make(chan bool) self.incrementWorkerCount() - go self.hashWorker(jobC, chunkC, errC, quitC, storeWg) + self.runHashWorker(jobC, chunkC, errC, quitC, storeWg) depth := 0 treeSize := self.chunkSize @@ -230,7 +230,7 @@ func (self *TreeChunker) split(depth int, treeSize int64, key Key, data io.Reade worker := self.getWorkerCount() if int64(len(jobC)) > worker && worker < ChunkProcessors { self.incrementWorkerCount() - go self.hashWorker(jobC, chunkC, errC, quitC, storeWg) + self.runHashWorker(jobC, chunkC, errC, quitC, storeWg) } select { @@ -239,23 +239,27 @@ func (self *TreeChunker) split(depth int, treeSize int64, key Key, data io.Reade } } -func (self *TreeChunker) hashWorker(jobC chan *hashJob, chunkC chan *Chunk, errC chan error, quitC chan bool, storeWg *sync.WaitGroup) { - defer self.decrementWorkerCount() +func (self *TreeChunker) runHashWorker(jobC chan *hashJob, chunkC chan *Chunk, errC chan error, quitC chan bool, storeWg *sync.WaitGroup) { + storeWg.Add(1) - hasher := self.hashFunc() - for { - select { + go func() { + defer self.decrementWorkerCount() + defer storeWg.Done() + hasher := self.hashFunc() + for { + select { - case job, ok := <-jobC: - if !ok { + case job, ok := <-jobC: + if !ok { + return + } + // now we got the hashes in the chunk, then hash the chunks + self.hashChunk(hasher, job, chunkC, storeWg) + case <-quitC: return } - // now we got the hashes in the chunk, then hash the chunks - self.hashChunk(hasher, job, chunkC, storeWg) - case <-quitC: - return } - } + }() } // The treeChunkers own Hash hashes together