From 40132af80ca8b485fbbcd201b4897f52e41b8766 Mon Sep 17 00:00:00 2001 From: Ashlesha Atrey Date: Thu, 27 Feb 2020 14:01:51 -0800 Subject: [PATCH] Improve performance of process bundles Improve the performance of processing bundles by parallelizing readIncludes() and detectManifestChanges() operations. Fixes #725 Signed-off-by: Ashlesha Atrey --- swupd/create_manifests.go | 213 ++++++++++++++++++++++++++++---------- 1 file changed, 159 insertions(+), 54 deletions(-) diff --git a/swupd/create_manifests.go b/swupd/create_manifests.go index 95efbd2..ed57c6d 100644 --- a/swupd/create_manifests.go +++ b/swupd/create_manifests.go @@ -147,8 +147,159 @@ func initBundles(ui UpdateInfo, c config, numWorkers int) ([]*Manifest, error) { return tmpManifests, err } +func readIncludes(tmpManifests []*Manifest, numWorkers int) (*Manifest, error) { + var wg sync.WaitGroup + if numWorkers <= 0 { + numWorkers = runtime.NumCPU() + } + if numWorkers > len(tmpManifests) { + numWorkers = len(tmpManifests) + } + wg.Add(numWorkers) -func processBundles(ui UpdateInfo, c config, numWorkers int) ([]*Manifest, error) { + manifestChan := make(chan *Manifest) + errorChan := make(chan error, numWorkers) + defer close(errorChan) + var newFull *Manifest + // read includes for subtraction processing + fmt.Println("Reading bundle includes...") + bundleWorker := func() { + defer wg.Done() + var err error + + for manifest := range manifestChan { + if manifest.Name == "full" { + newFull = manifest + continue + } + if manifest.Name != "os-core" { + // read in bundle includes + if err = manifest.ReadIncludesFromBundleInfo(tmpManifests); err != nil { + errorChan <- err + return + } + } + } + } + for i := 0; i < numWorkers; i++ { + go bundleWorker() + } + + var err error + for _, bn := range tmpManifests { + select { + case manifestChan <- bn: + case err = <-errorChan: + // If there is an error,we break + // immediately and return error + break + } + if err != nil { + break + } + } + + close(manifestChan) + wg.Wait() + + if err != nil { + return nil, err + } + + if len(errorChan) > 0 { + return nil, <-errorChan + } + + return newFull, err +} + +func detectManifestChanges(oldMoM *Manifest, tmpManifests []*Manifest, numWorkers int, ui UpdateInfo, c config) ([]*Manifest, error) { + fmt.Println("Detecting manifest changes...") + + newManifests := []*Manifest{} + taskCh := make(chan *Manifest) + var wg sync.WaitGroup + mux := &sync.Mutex{} + + if numWorkers <= 0 { + numWorkers = runtime.NumCPU() + } + if numWorkers > len(tmpManifests) { + numWorkers = len(tmpManifests) + } + + wg.Add(numWorkers) + errorCh := make(chan error, numWorkers) + + defer close(errorCh) + + taskRunner := func() { + defer wg.Done() + for manifest := range taskCh { + var tErr error + // Check for changed includes, changed or added or deleted files + // must be done after subtractManifests because the oldM is a subtracted + // manifest + ver := getManifestVerFromMoM(oldMoM, manifest) + if ver == 0 { + ver = ui.previous + } + + oldMPath := filepath.Join(c.outputDir, fmt.Sprint(ver), "Manifest."+manifest.Name) + oldM, err := getOldManifest(oldMPath) + if err != nil { + errorCh <- tErr + return + } + changedIncludes := includesChanged(manifest, oldM) + oldM.sortFilesName() + changedFiles, added, deleted := manifest.linkPeersAndChange(oldM, ui.minVersion) + // if nothing changed, skip + if changedFiles == 0 && added == 0 && deleted == 0 && !changedIncludes { + continue + } + + // detect modifier flag for all files in the manifest + // must happen after finding newDeleted files to catch ghosted files. + manifest.applyHeuristics() + // Assign final FileCount based on the files that made it this far + manifest.Header.FileCount = uint32(len(manifest.Files)) + // If we made it this far, this bundle has a change and should be written + mux.Lock() + newManifests = append(newManifests, manifest) + mux.Unlock() + } + } + + for i := 0; i < numWorkers; i++ { + go taskRunner() + } + var err error + for _, tmpMan := range tmpManifests { + select { + case taskCh <- tmpMan: + case err = <-errorCh: + // break as soon as there is a failure. + break + } + if err != nil { + break + } + } + close(taskCh) + wg.Wait() + if err != nil { + return nil, err + } + // Sending loop might finish before any goroutine could send an error back, so check for + // error again after they are all done. + if len(errorCh) > 0 { + return nil, <-errorCh + } + return newManifests, nil +} + +func processBundles(ui UpdateInfo, c config, numWorkers int, oldMoM *Manifest) ([]*Manifest, error) { var newFull *Manifest var err error // initialize bundles with with their info files @@ -157,21 +308,11 @@ func processBundles(ui UpdateInfo, c config, numWorkers int) ([]*Manifest, error return nil, err } - // read includes for subtraction processing - fmt.Println("Reading bundle includes...") - for _, bundle := range tmpManifests { - if bundle.Name == "full" { - newFull = bundle - continue - } - if bundle.Name != "os-core" { - // read in bundle includes - if err = bundle.ReadIncludesFromBundleInfo(tmpManifests); err != nil { - return nil, err - } - } + //read includes from bundleinfo files + newFull, err = readIncludes(tmpManifests, numWorkers) + if err != nil { + return nil, err } - // Add manifest file records. Important this is done after all includes // have been read so nested subtraction works. fmt.Println("Adding manifest file records...") @@ -179,47 +320,11 @@ func processBundles(ui UpdateInfo, c config, numWorkers int) ([]*Manifest, error return nil, err } - // Need old MoM to get version of last bundle manifest - oldMoMPath := filepath.Join(c.outputDir, fmt.Sprint(ui.previous), "Manifest.MoM") - oldMoM, err := getOldManifest(oldMoMPath) + // final loop detects changes, applies heuristics to files, and sorts the file lists + newManifests, err := detectManifestChanges(oldMoM, tmpManifests, numWorkers, ui, c) if err != nil { return nil, err } - - // final loop detects changes, applies heuristics to files, and sorts the file lists - fmt.Println("Detecting manifest changes...") - newManifests := []*Manifest{} - for _, bundle := range tmpManifests { - // Check for changed includes, changed or added or deleted files - // must be done after subtractManifests because the oldM is a subtracted - // manifest - ver := getManifestVerFromMoM(oldMoM, bundle) - if ver == 0 { - ver = ui.previous - } - - oldMPath := filepath.Join(c.outputDir, fmt.Sprint(ver), "Manifest."+bundle.Name) - oldM, err := getOldManifest(oldMPath) - if err != nil { - return nil, err - } - changedIncludes := includesChanged(bundle, oldM) - oldM.sortFilesName() - changedFiles, added, deleted := bundle.linkPeersAndChange(oldM, ui.minVersion) - // if nothing changed, skip - if changedFiles == 0 && added == 0 && deleted == 0 && !changedIncludes { - continue - } - - // detect modifier flag for all files in the manifest - // must happen after finding newDeleted files to catch ghosted files. - bundle.applyHeuristics() - // Assign final FileCount based on the files that made it this far - bundle.Header.FileCount = uint32(len(bundle.Files)) - // If we made it this far, this bundle has a change and should be written - newManifests = append(newManifests, bundle) - } - // maximize full manifest while all the manifests are still sorted by name maximizeFull(newFull, newManifests) @@ -327,7 +432,7 @@ func CreateManifests(version, previous, minVersion uint32, format uint, statedir timeStamp: timeStamp, } var newManifests []*Manifest - if newManifests, err = processBundles(ui, c, numWorkers); err != nil { + if newManifests, err = processBundles(ui, c, numWorkers, oldMoM); err != nil { return nil, err }