From a59a9b5fb25cae3e1076f6e0123c25da0b0648c9 Mon Sep 17 00:00:00 2001 From: zsfelfoldi Date: Sat, 30 May 2015 08:06:37 +0200 Subject: [PATCH] implemented manifest trie with add and remove Api.Upload also creates a trie now --- bzz/api.go | 57 +++++++------- bzz/manifest_trie.go | 178 +++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 205 insertions(+), 30 deletions(-) create mode 100644 bzz/manifest_trie.go diff --git a/bzz/api.go b/bzz/api.go index 38fcf6933e..8c8a247433 100644 --- a/bzz/api.go +++ b/bzz/api.go @@ -135,7 +135,7 @@ const maxParallelFiles = 5 // using dpa store // TODO: localpath should point to a manifest func (self *Api) Upload(lpath string) (string, error) { - var files []string + var list []*manifestTrieEntry localpath, err1 := filepath.Abs(filepath.Clean(lpath)) if err1 != nil { return "", err1 @@ -154,7 +154,10 @@ func (self *Api) Upload(lpath string) (string, error) { if path[:len(localpath)] != localpath { return fmt.Errorf("Path prefix of '%s' does not match localpath '%s'", path, localpath) } - files = append(files, path) + entry := &manifestTrieEntry{ + Path: path, + } + list = append(list, entry) } return err }) @@ -162,68 +165,62 @@ func (self *Api) Upload(lpath string) (string, error) { return "", err } - cnt := len(files) - hashes := make([]Key, cnt) + cnt := len(list) errors := make([]error, cnt) - ctypes := make([]string, cnt) done := make(chan bool, maxParallelFiles) dcnt := 0 - for i, path := range files { + for i, entry := range list { if i >= dcnt+maxParallelFiles { <-done dcnt++ } - go func(i int, path string, done chan bool) { - f, err := os.Open(path) + go func(i int, entry *manifestTrieEntry, done chan bool) { + f, err := os.Open(entry.Path) if err == nil { stat, _ := f.Stat() sr := io.NewSectionReader(f, 0, stat.Size()) wg := &sync.WaitGroup{} - hashes[i], err = self.dpa.Store(sr, wg) + var hash Key + hash, err = self.dpa.Store(sr, wg) + if hash != nil { + list[i].Hash = fmt.Sprintf("%064x", hash) + } wg.Wait() } if err == nil { - cmd := exec.Command("file", "--mime-type", "-b", path) + cmd := exec.Command("file", "--mime-type", "-b", entry.Path) var out bytes.Buffer cmd.Stdout = &out err = cmd.Run() if err == nil { - ctypes[i] = strings.TrimSuffix(out.String(), "\n") + list[i].ContentType = strings.TrimSuffix(out.String(), "\n") } } errors[i] = err done <- true - }(i, path, done) + }(i, entry, done) } for dcnt < cnt { <-done dcnt++ } - var buffer bytes.Buffer - buffer.WriteString(`{"entries":[`) - sc := "," - if err != nil { - return "", err - } - - for i, path := range files { + trie := &manifestTrie{} + for i, entry := range list { if errors[i] != nil { return "", errors[i] } - if i == cnt-1 { - sc = "]}" - } - buffer.WriteString(fmt.Sprintf(`{"hash":"%064x","path":"%s","contentType":"%s"}%s`, hashes[i], path[start:], ctypes[i], sc)) + entry.Path = entry.Path[start:] + trie.addEntry(entry) } - manifest := buffer.Bytes() - sr := io.NewSectionReader(bytes.NewReader(manifest), 0, int64(len(manifest))) - wg := &sync.WaitGroup{} - key, err2 := self.dpa.Store(sr, wg) - wg.Wait() - return fmt.Sprintf("%064x", key), err2 + err2 := trie.recalcAndStore(self.dpa) + var hs string + if err2 == nil { + hs = fmt.Sprintf("%064x", trie.hash) + } + return hs, err2 } func (self *Api) Register(sender common.Address, hash common.Hash, domain string) (err error) { diff --git a/bzz/manifest_trie.go b/bzz/manifest_trie.go new file mode 100644 index 0000000000..524b40c6ee --- /dev/null +++ b/bzz/manifest_trie.go @@ -0,0 +1,178 @@ +package bzz + +import ( + "bytes" + "encoding/json" + "fmt" + "io" + "sync" +) + +type manifestTrie struct { + entries [257]*manifestTrieEntry // indexed by first character of path, entries[256] is the empty path entry + hash Key // if hash != nil, it is stored +} + +type manifestJSON struct { + Entries []*manifestTrieEntry `json:"entries"` +} + +type manifestTrieEntry struct { + Path string `json:"path"` + Hash string `json:"hash"` // for manifest content type, empty until subtrie is evaluated + ContentType string `json:"contentType"` + subtrie *manifestTrie +} + +func loadManifestTrie(dpa *DPA, hash Key) (trie *manifestTrie, err error) { + + dpaLogger.Debugf("Swarm: manifest lookup key: '%064x'.", hash) + // retrieve manifest via DPA + manifestReader := dpa.Retrieve(hash) + // TODO check size for oversized manifests + manifestData := make([]byte, manifestReader.Size()) + var size int + size, err = manifestReader.Read(manifestData) + if int64(size) < manifestReader.Size() { + dpaLogger.Debugf("Swarm: Manifest %064x not found.", hash) + if err == nil { + err = fmt.Errorf("Manifest retrieval cut short: %v < %v", size, manifestReader.Size()) + } + return + } + + dpaLogger.Debugf("Swarm: Manifest %064x retrieved", hash) + man := manifestJSON{} + err = json.Unmarshal(manifestData, &man) + if err != nil { + err = fmt.Errorf("Manifest %064x is malformed: %v", hash, err) + dpaLogger.Debugf("Swarm: %v", err) + return + } + + dpaLogger.Debugf("Swarm: Manifest %064x has %d entries.", hash, len(man.Entries)) + + trie = &manifestTrie{} + for _, entry := range man.Entries { + trie.addEntry(entry) + } + return +} + +func (self *manifestTrie) addEntry(entry *manifestTrieEntry) { + self.hash = nil // trie modified, hash needs to be re-calculated on demand + + if len(entry.Path) == 0 { + self.entries[256] = entry + return + } + + b := byte(entry.Path[0]) + if (self.entries[b] == nil) || (self.entries[b].Path == entry.Path) { + self.entries[b] = entry + return + } + + oldentry := self.entries[b] + cpl := 0 + for (len(entry.Path) > cpl) && (len(oldentry.Path) > cpl) && (entry.Path[cpl] == oldentry.Path[cpl]) { + cpl++ + } + + if (oldentry.ContentType == manifestType) && (cpl == len(oldentry.Path)) { + entry.Path = entry.Path[cpl:] + oldentry.subtrie.addEntry(entry) + oldentry.Hash = "" + return + } + + commonPrefix := entry.Path[:cpl] + + subtrie := &manifestTrie{} + entry.Path = entry.Path[cpl:] + oldentry.Path = oldentry.Path[cpl:] + subtrie.addEntry(entry) + subtrie.addEntry(oldentry) + + self.entries[b] = &manifestTrieEntry{ + Path: commonPrefix, + Hash: "", + ContentType: manifestType, + subtrie: subtrie, + } +} + +func (self *manifestTrie) getCountLast() (cnt int, entry *manifestTrieEntry) { + for _, e := range self.entries { + if e != nil { + cnt++ + entry = e + } + } + return +} + +func (self *manifestTrie) deleteEntry(path string) { + self.hash = nil // trie modified, hash needs to be re-calculated on demand + + if len(path) == 0 { + self.entries[256] = nil + return + } + + b := byte(path[0]) + if (self.entries[b] != nil) && (self.entries[b].Path == path) { + self.entries[b] = nil + return + } + + entry := self.entries[b] + epl := len(entry.Path) + if (entry.ContentType == manifestType) && (len(path) >= epl) && (path[:epl] == entry.Path) { + entry.subtrie.deleteEntry(path[epl:]) + entry.Hash = "" + // remove subtree if it has less than 2 elements + cnt, lastentry := entry.subtrie.getCountLast() + if cnt < 2 { + if lastentry != nil { + lastentry.Path = entry.Path + lastentry.Path + } + self.entries[b] = lastentry + } + } +} + +func (self *manifestTrie) recalcAndStore(dpa *DPA) error { + if self.hash != nil { + return nil + } + + var buffer bytes.Buffer + buffer.WriteString(`{"entries":[`) + + list := &manifestJSON{} + for _, entry := range self.entries { + if entry != nil { + if entry.Hash == "" { // TODO: paralellize + err := entry.subtrie.recalcAndStore(dpa) + if err != nil { + return err + } + entry.Hash = fmt.Sprintf("%064x", entry.subtrie.hash) + } + list.Entries = append(list.Entries, entry) + } + } + + manifest, err := json.Marshal(list) + if err != nil { + return err + } + + sr := io.NewSectionReader(bytes.NewReader(manifest), 0, int64(len(manifest))) + wg := &sync.WaitGroup{} + key, err2 := dpa.Store(sr, wg) + wg.Wait() + self.hash = key + return err2 +}