Merged from trunk 10928:10943
authorFabian Groffen <grobian@gentoo.org>
Mon, 7 Jul 2008 09:11:52 +0000 (09:11 -0000)
committerFabian Groffen <grobian@gentoo.org>
Mon, 7 Jul 2008 09:11:52 +0000 (09:11 -0000)
   | 10929   | Bug #230525 - Work around ObjectProxy breakage in `portageq  |
   | zmedico | vdb_path`.                                                   |

   | 10930   | Add a PollConstants class to serve as a layer of indirection |
   | zmedico | for accessing select.POLL* attributes which do not exist on  |
   |         | interix. Thanks to grobian for reporting.                    |

   | 10931   | TODO: For parallel scheduling, everything in MergeListItem   |
   | zmedico | needs asynchronous execution support (start, poll, and wait  |
   |         | methods).                                                    |

   | 10932   | Convert EbuildFetcher to an asynchronous interface: * Rename |
   | zmedico | EbuildFetcher to EbuildFetchPretend. * Rename                |
   |         | EbuildFetcherAsync to EbuildFetcher.                         |

   | 10933   | Make EbuildBuild inherit from EbuildBuildDir since it's      |
   | zmedico | going to hold the build dir lock while it runs               |
   |         | asynchronously.                                              |

   | 10934   | Make Binpkg inherit from EbuildBuildDir since it's going to  |
   | zmedico | hold the build dir lock while it runs asynchronously.        |

   | 10935   | * Implement a new binarytree.digestCheck() method and use it |
   | zmedico | to check digests for binary packages. * Split out a          |
   |         | AsynchronousTask base class from SubProcess. * Derive a new  |
   |         | BinpkgVerifier class from AsynchronousTask. Even though this |
   |         | is not really asynchronous yet, it can fake the interface by |
   |         | doing everything in the start() method.                      |

   | 10936   | Rename BinpkgFetcherAsync class to BinpkgFetcher and use it  |
   | zmedico | to replace the older class with the same name                |

   | 10937   | Remove unused variable.                                      |
   | zmedico |                                                              |

   | 10938   | Use the EbuildPhase class to execute the "clean" phase       |
   | zmedico | asynchronously.                                              |

   | 10939   | Add async support to EbuildBinpkg by making it inherit from  |
   | zmedico | EbuildPhase.                                                 |

   | 10940   | Add a new CompositeTask class which can be used to combine   |
   | zmedico | separate AsynchronousTask instances into a single instance.  |
   |         | The CompositeTask instance used task exit listeners as a     |
   |         | means to (asynchronously) trigger progression from one       |
   |         | subtask to the next. This technique is used to group         |
   |         | together all the ebuild phases executed by EbuildExecuter,   |
   |         | and should be useful for grouping many more sets of tasks    |
   |         | into similar composite tasks.                                |

   | 10941   | Fix CompositeTask.cancel() so that it's safe to call when    |
   | zmedico | there is no running task.                                    |

   | 10942   | Make EbuildExecuter inherit from AsynchronousTask.           |
   | zmedico |                                                              |

   | 10943   | Split out common code from CompositeTask and EbuildExecuter  |
   | zmedico | into a new CompositeTask class and rename the old            |
   |         | CompositeTask class to TaskSequence.                         |

svn path=/main/branches/prefix/; revision=10967

bin/portageq
pym/_emerge/__init__.py
pym/portage/__init__.py
pym/portage/dbapi/bintree.py
pym/portage/dbapi/vartree.py
pym/portage/exception.py
pym/portage/locks.py

index 3acba5079e08eb2d1a2ed673a7d5567db93b47df..2a0cf6f3f713d419bdf869f75afac4f0c0f445f9 100755 (executable)
@@ -294,8 +294,10 @@ def vdb_path(argv):
        Returns the path used for the var(installed) package database for the
        set environment/configuration options.
        """
-       print portage.root+portage.VDB_PATH
-
+       out = sys.stdout
+       out.write(os.path.join(portage.settings["ROOT"], portage.VDB_PATH) + "\n")
+       out.flush()
+       return os.EX_OK
 
 def gentoo_mirrors(argv):
        """
@@ -475,18 +477,18 @@ def main():
                os.environ["ROOT"] = sys.argv[2]
 
        global portage
+       try:
+               import portage
+       except ImportError:
+               from os import path as osp
+               # for an explanation on this logic, see pym/_emerge/__init__.py
+               if os.environ.__contains__("PORTAGE_PYTHONPATH"):
+                       sys.path.insert(0, os.environ["PORTAGE_PYTHONPATH"])
+               else:
+                       sys.path.insert(0, osp.join(osp.dirname(osp.dirname(osp.realpath(__file__))), "pym"))
+               import portage
 
        try:
-               try:
-                       import portage
-               except ImportError:
-                       from os import path as osp
-                       # for an explanation on this logic, see pym/_emerge/__init__.py
-                       if os.environ.__contains__("PORTAGE_PYTHONPATH"):
-                               sys.path.insert(0, os.environ["PORTAGE_PYTHONPATH"])
-                       else:
-                               sys.path.insert(0, os.path.join(os.path.dirname(os.path.dirname(os.path.realpath(__file__))), "pym"))
-                       import portage
                if uses_root:
                        sys.argv[2] = portage.settings["ROOT"]
                retval = function(sys.argv[2:])
index 5380c2ad9a5b93d2d07cdac1691d6cffd62942c4..dd43dc54977fdda51c0620d87f82b5edef32d31e 100644 (file)
@@ -21,6 +21,7 @@ except KeyboardInterrupt:
        sys.exit(1)
 
 import array
+from collections import deque
 import fcntl
 import select
 import shlex
@@ -1465,9 +1466,9 @@ class _PackageMetadataWrapper(_PackageMetadataWrapperBase):
                                v = 0
                self._pkg.mtime = v
 
-class EbuildFetcher(SlotObject):
+class EbuildFetchPretend(SlotObject):
 
-       __slots__ = ("cancelled", "fetch_all", "pkg", "pretend", "settings")
+       __slots__ = ("fetch_all", "pkg", "settings")
 
        def execute(self):
                portdb = self.pkg.root_config.trees["porttree"].dbapi
@@ -1476,12 +1477,122 @@ class EbuildFetcher(SlotObject):
 
                retval = portage.doebuild(ebuild_path, "fetch",
                        self.settings["ROOT"], self.settings, debug=debug,
-                       listonly=self.pretend, fetchonly=1, fetchall=self.fetch_all,
+                       listonly=1, fetchonly=1, fetchall=self.fetch_all,
                        mydbapi=portdb, tree="porttree")
                return retval
 
-class SubProcess(SlotObject):
-       __slots__ = ("cancelled", "pid", "returncode")
+class AsynchronousTask(SlotObject):
+       __slots__ = ("cancelled", "returncode") + ("_exit_listeners",)
+
+       def start(self):
+               """
+               Start an asynchronous task and then return as soon as possible.
+               """
+               pass
+
+       def isAlive(self):
+               return self.returncode is None
+
+       def poll(self):
+               return self.returncode
+
+       def wait(self):
+               self._wait_hook()
+               return self.returncode
+
+       def cancel(self):
+               pass
+
+       def addExitListener(self, f):
+               """
+               The function will be called with one argument, a reference to self.
+               """
+               if self._exit_listeners is None:
+                       self._exit_listeners = []
+               self._exit_listeners.append(f)
+
+       def removeExitListener(self, f):
+               self._exit_listeners.remove(f)
+
+       def _wait_hook(self):
+               """
+               Call this method before returning from wait. This hook is
+               used to trigger exit listeners when the returncode first
+               becomes available.
+               """
+               if self._exit_listeners is not None:
+                       for f in self._exit_listeners:
+                               f(self)
+                       self._exit_listeners = None
+
+class CompositeTask(AsynchronousTask):
+
+       __slots__ = ("scheduler",) + ("_current_task",)
+
+       def isAlive(self):
+               return self._current_task is not None
+
+       def cancel(self):
+               self.cancelled = True
+               if self._current_task is not None:
+                       self._current_task.cancel()
+
+       def wait(self):
+
+               while True:
+                       task = self._current_task
+                       if task is None:
+                               break
+                       if hasattr(task, "reg_id"):
+                               self.scheduler.schedule(task.reg_id)
+                       task.wait()
+
+               self._wait_hook()
+               return self.returncode
+
+class TaskSequence(CompositeTask):
+       """
+       A collection of tasks that executes sequentially. Each task
+       must have a _set_returncode() method that can be wrapped as
+       a means to trigger movement from one task to the next.
+       """
+
+       __slots__ = ("_task_queue",)
+
+       def __init__(self, **kwargs):
+               AsynchronousTask.__init__(self, **kwargs)
+               self._task_queue = deque()
+
+       def add(self, task):
+               self._task_queue.append(task)
+
+       def start(self):
+               self._start_next_task()
+
+       def cancel(self):
+               self._task_queue.clear()
+               CompositeTask.cancel(self)
+
+       def _start_next_task(self):
+               self._current_task = self._task_queue.popleft()
+               task = self._current_task
+               task.addExitListener(self._task_exit_handler)
+               task.start()
+
+       def _task_exit_handler(self, task):
+               if task is not self._current_task:
+                       raise AssertionError("Unrecognized task: %s" % (task,))
+
+               if self._task_queue and \
+                       task.returncode == os.EX_OK:
+                       self._start_next_task()
+                       return
+
+               self._current_task = None
+               self.returncode = task.returncode
+
+class SubProcess(AsynchronousTask):
+       __slots__ = ("pid",)
 
        def poll(self):
                if self.returncode is not None:
@@ -1508,6 +1619,7 @@ class SubProcess(SlotObject):
                if self.returncode is not None:
                        return self.returncode
                self._set_returncode(os.waitpid(self.pid, 0))
+               self._wait_hook()
                return self.returncode
 
        def _set_returncode(self, wait_retval):
@@ -1534,7 +1646,7 @@ class SpawnProcess(SubProcess):
                "uid", "gid", "groups", "umask", "logfile",
                "path_lookup", "pre_exec")
 
-       __slots__ = ("args", "files", "register", "unregister", "registered") + \
+       __slots__ = ("args", "files", "registered", "reg_id", "scheduler") + \
                _spawn_kwarg_names
 
        _file_names = ("process", "out")
@@ -1600,9 +1712,9 @@ class SpawnProcess(SubProcess):
 
                os.close(slave_fd)
                files.process = os.fdopen(master_fd, 'r')
+               self.reg_id = self.scheduler.register(files.process.fileno(),
+                       PollConstants.POLLIN, self._output_handler)
                self.registered = True
-               self.register(files.process.fileno(),
-                       select.POLLIN, self._output_handler)
 
        def _output_handler(self, fd, event):
                files = self.files
@@ -1620,9 +1732,9 @@ class SpawnProcess(SubProcess):
                                f.flush()
                                f.close()
                        self.registered = False
-                       self.unregister(fd)
+               return self.registered
 
-class EbuildFetcherAsync(SpawnProcess):
+class EbuildFetcher(SpawnProcess):
 
        __slots__ = ("pkg",)
 
@@ -1724,7 +1836,7 @@ class EbuildBuildDir(SlotObject):
        class AlreadyLocked(portage.exception.PortageException):
                pass
 
-class EbuildBuild(SlotObject):
+class EbuildBuild(EbuildBuildDir):
 
        __slots__ = ("args_set", "find_blockers",
                "ldpath_mtimes", "logger", "opts",
@@ -1763,18 +1875,29 @@ class EbuildBuild(SlotObject):
                                not opts.buildpkg
 
                if opts.fetchonly:
-                       fetcher = EbuildFetcher(fetch_all=opts.fetch_all_uri,
-                               pkg=pkg, pretend=opts.pretend, settings=settings)
-                       retval = fetcher.execute()
+                       if opts.pretend:
+
+                               fetcher = EbuildFetchPretend(
+                                       fetch_all=opts.fetch_all_uri,
+                                       pkg=pkg, settings=settings)
+
+                               retval = fetcher.execute()
+
+                       else:
+
+                               fetcher = EbuildFetcher(pkg=pkg, scheduler=scheduler)
+                               fetcher.start()
+                               scheduler.schedule(fetcher.reg_id)
+                               retval = fetcher.wait()
+
                        if retval != os.EX_OK:
                                from portage.elog.messages import eerror
                                eerror("!!! Fetch for %s failed, continuing..." % pkg.cpv,
                                        phase="unpack", key=pkg.cpv)
                        return retval
 
-               build_dir = EbuildBuildDir(pkg=pkg, settings=settings)
                try:
-                       build_dir.lock()
+                       self.lock()
                        # Cleaning is triggered before the setup
                        # phase, in portage.doebuild().
                        msg = " === (%s of %s) Cleaning (%s::%s)" % \
@@ -1793,15 +1916,20 @@ class EbuildBuild(SlotObject):
                                        (pkg_count.curval, pkg_count.maxval, pkg.cpv)
                                logger.log(msg, short_msg=short_msg)
 
-                               build = EbuildExecuter(pkg=pkg, register=scheduler.register,
-                                       schedule=scheduler.schedule, settings=settings,
-                                       unregister=scheduler.unregister)
-                               retval = build.execute()
+                               build = EbuildExecuter(pkg=pkg, scheduler=scheduler,
+                                       settings=settings)
+                               build.start()
+                               retval = build.wait()
                                if retval != os.EX_OK:
                                        return retval
 
-                               build = EbuildBinpkg(pkg=pkg, settings=settings)
-                               retval = build.execute()
+                               build = EbuildBinpkg(pkg=pkg,
+                                       scheduler=scheduler, settings=settings)
+
+                               build.start()
+                               scheduler.schedule(build.reg_id)
+                               retval = build.wait()
+
                                if retval != os.EX_OK:
                                        return retval
 
@@ -1831,10 +1959,10 @@ class EbuildBuild(SlotObject):
                                        (pkg_count.curval, pkg_count.curval, pkg.cpv)
                                logger.log(msg, short_msg=short_msg)
 
-                               build = EbuildExecuter(pkg=pkg, register=scheduler.register,
-                                       schedule=scheduler.schedule, settings=settings,
-                                       unregister=scheduler.unregister)
-                               retval = build.execute()
+                               build = EbuildExecuter(pkg=pkg, scheduler=scheduler,
+                                       settings=settings)
+                               build.start()
+                               retval = build.wait()
                                if retval != os.EX_OK:
                                        return retval
 
@@ -1847,34 +1975,45 @@ class EbuildBuild(SlotObject):
                                if retval != os.EX_OK:
                                        return retval
                finally:
-                       if build_dir.locked:
+                       if self.locked:
                                portage.elog.elog_process(pkg.cpv, settings)
-                               build_dir.unlock()
+                               self.unlock()
                return os.EX_OK
 
-class EbuildExecuter(SlotObject):
+class EbuildExecuter(CompositeTask):
 
-       __slots__ = ("pkg", "register", "schedule", "settings", "unregister")
+       __slots__ = ("pkg", "scheduler", "settings")
 
        _phases = ("setup", "unpack", "compile", "test", "install")
 
-       def execute(self):
-               root_config = self.pkg.root_config
+       def start(self):
+               pkg = self.pkg
+               scheduler = self.scheduler
                tree = "porttree"
-               portdb = root_config.trees[tree].dbapi
-               ebuild_path = portdb.findname(self.pkg.cpv)
                settings = self.settings
-               debug = settings.get("PORTAGE_DEBUG") == "1"
-               cleanup = 1
 
-               retval = portage.doebuild(ebuild_path, "clean",
-                       root_config.root, settings, debug, cleanup=cleanup,
-                       mydbapi=portdb, tree="porttree")
-               if retval != os.EX_OK:
-                       return retval
+               phase = "clean"
+               clean_phase = EbuildPhase(pkg=pkg, phase=phase,
+                       scheduler=scheduler, settings=settings, tree=tree)
+               clean_phase.addExitListener(self._clean_phase_exit)
+               self._current_task = clean_phase
+               clean_phase.start()
+
+       def _clean_phase_exit(self, clean_phase):
+
+               if clean_phase.returncode != os.EX_OK:
+                       self.returncode = clean_phase.returncode
+                       self._current_task = None
+                       return
+
+               pkg = self.pkg
+               scheduler = self.scheduler
+               tree = "porttree"
+               settings = self.settings
+               cleanup = 1
 
                # This initializes PORTAGE_LOG_FILE.
-               portage.prepare_build_dirs(root_config.root, settings, cleanup)
+               portage.prepare_build_dirs(pkg.root, settings, cleanup)
 
                fd_pipes = {
                        0 : sys.stdin.fileno(),
@@ -1882,32 +2021,36 @@ class EbuildExecuter(SlotObject):
                        2 : sys.stderr.fileno(),
                }
 
-               for mydo in self._phases:
-                       ebuild_phase = EbuildPhase(fd_pipes=fd_pipes,
-                               pkg=self.pkg, phase=mydo, register=self.register,
-                               settings=settings, tree=tree, unregister=self.unregister)
+               ebuild_phases = TaskSequence(scheduler=scheduler)
 
-                       ebuild_phase.start()
-                       retval = None
-                       while retval is None:
-                               self.schedule()
-                               retval = ebuild_phase.poll()
+               for phase in self._phases:
+                       ebuild_phases.add(EbuildPhase(fd_pipes=fd_pipes,
+                               pkg=pkg, phase=phase, scheduler=scheduler,
+                               settings=settings, tree=tree))
 
-                       if retval != os.EX_OK:
-                               return retval
+               ebuild_phases.addExitListener(self._ebuild_phases_exit)
+               self._current_task = ebuild_phases
+               ebuild_phases.start()
 
-               return os.EX_OK
+       def _ebuild_phases_exit(self, ebuild_phases):
+               self.returncode = ebuild_phases.returncode
+               self._current_task = None
 
 class EbuildPhase(SubProcess):
 
        __slots__ = ("fd_pipes", "phase", "pkg",
-               "register", "settings", "tree", "unregister",
-               "files", "registered")
+               "scheduler", "settings", "tree",
+               "files", "registered", "reg_id")
 
        _file_names = ("log", "stdout", "ebuild")
        _files_dict = slot_dict_class(_file_names, prefix="")
        _bufsize = 4096
 
+       # A file descriptor is required for the scheduler to monitor changes from
+       # inside a poll() loop. When logging is not enabled, create a pipe just to
+       # serve this purpose alone.
+       _dummy_pipe_fd = 9
+
        def start(self):
                root_config = self.pkg.root_config
                tree = self.tree
@@ -1918,7 +2061,15 @@ class EbuildPhase(SubProcess):
                logfile = settings.get("PORTAGE_LOG_FILE")
                master_fd = None
                slave_fd = None
-               fd_pipes = self.fd_pipes.copy()
+               fd_pipes = None
+               if self.fd_pipes is not None:
+                       fd_pipes = self.fd_pipes.copy()
+               else:
+                       fd_pipes = {}
+
+               fd_pipes.setdefault(0, sys.stdin.fileno())
+               fd_pipes.setdefault(1, sys.stdout.fileno())
+               fd_pipes.setdefault(2, sys.stderr.fileno())
 
                # flush any pending output
                for fd in fd_pipes.itervalues():
@@ -1958,7 +2109,6 @@ class EbuildPhase(SubProcess):
                                mode[1] &= ~termios.OPOST
                                termios.tcsetattr(slave_fd, termios.TCSANOW, mode)
 
-                       import fcntl
                        fcntl.fcntl(master_fd, fcntl.F_SETFL,
                                fcntl.fcntl(master_fd, fcntl.F_GETFL) | os.O_NONBLOCK)
 
@@ -1972,6 +2122,14 @@ class EbuildPhase(SubProcess):
                        fd_pipes[1] = slave_fd
                        fd_pipes[2] = slave_fd
 
+               else:
+                       # Create a dummy pipe so the scheduler can monitor
+                       # the process from inside a poll() loop.
+                       master_fd, slave_fd = os.pipe()
+                       fcntl.fcntl(master_fd, fcntl.F_SETFL,
+                               fcntl.fcntl(master_fd, fcntl.F_GETFL) | os.O_NONBLOCK)
+                       fd_pipes[self._dummy_pipe_fd] = slave_fd
+
                retval = portage.doebuild(ebuild_path, self.phase,
                        root_config.root, settings, debug,
                        mydbapi=mydbapi, tree=tree,
@@ -1980,13 +2138,17 @@ class EbuildPhase(SubProcess):
                self.pid = retval[0]
 
                if logfile:
-                       os.close(slave_fd)
                        files.log = open(logfile, 'a')
                        files.stdout = os.fdopen(os.dup(fd_pipes_orig[1]), 'w')
-                       files.ebuild = os.fdopen(master_fd, 'r')
-                       self.registered = True
-                       self.register(files.ebuild.fileno(),
-                               select.POLLIN, self._output_handler)
+                       output_handler = self._output_handler
+               else:
+                       output_handler = self._dummy_handler
+
+               os.close(slave_fd)
+               files.ebuild = os.fdopen(master_fd, 'r')
+               self.reg_id = self.scheduler.register(files.ebuild.fileno(),
+                       PollConstants.POLLIN, output_handler)
+               self.registered = True
 
        def _output_handler(self, fd, event):
                files = self.files
@@ -2005,18 +2167,40 @@ class EbuildPhase(SubProcess):
                        for f in files.values():
                                f.close()
                        self.registered = False
-                       self.unregister(fd)
+               return self.registered
+
+       def _dummy_handler(self, fd, event):
+               """
+               This method is mainly interested in detecting EOF, since
+               the only purpose of the pipe is to allow the scheduler to
+               monitor the process from inside a poll() loop.
+               """
+               files = self.files
+               buf = array.array('B')
+               try:
+                       buf.fromfile(files.ebuild, self._bufsize)
+               except EOFError:
+                       pass
+               if buf:
+                       pass
+               else:
+                       fd = files.ebuild.fileno()
+                       for f in files.values():
+                               f.close()
+                       self.registered = False
+               return self.registered
 
        def _set_returncode(self, wait_retval):
                SubProcess._set_returncode(self, wait_retval)
-               msg = portage._doebuild_exit_status_check(
-                       self.phase, self.settings)
-               if msg:
-                       self.returncode = 1
-                       from textwrap import wrap
-                       from portage.elog.messages import eerror
-                       for l in wrap(msg, 72):
-                               eerror(l, phase=self.phase, key=self.pkg.cpv)
+               if self.phase != "clean":
+                       msg = portage._doebuild_exit_status_check(
+                               self.phase, self.settings)
+                       if msg:
+                               self.returncode = 1
+                               from textwrap import wrap
+                               from portage.elog.messages import eerror
+                               for l in wrap(msg, 72):
+                                       eerror(l, phase=self.phase, key=self.pkg.cpv)
 
                returncode = self.returncode
                settings = self.settings
@@ -2026,19 +2210,15 @@ class EbuildPhase(SubProcess):
                        if returncode == os.EX_OK:
                                returncode = portage._post_src_install_checks(settings)
 
-class EbuildBinpkg(Task):
+class EbuildBinpkg(EbuildPhase):
        """
        This assumes that src_install() has successfully completed.
        """
-       __slots__ = ("pkg", "settings")
-
-       def _get_hash_key(self):
-               hash_key = getattr(self, "_hash_key", None)
-               if hash_key is None:
-                       self._hash_key = ("EbuildBinpkg", self.pkg._get_hash_key())
-               return self._hash_key
+       __slots__ = ("_binpkg_tmpfile",)
 
-       def execute(self):
+       def start(self):
+               self.phase = "package"
+               self.tree = "porttree"
                pkg = self.pkg
                root_config = pkg.root_config
                portdb = root_config.trees["porttree"].dbapi
@@ -2050,29 +2230,23 @@ class EbuildBinpkg(Task):
                bintree.prevent_collision(pkg.cpv)
                binpkg_tmpfile = os.path.join(bintree.pkgdir,
                        pkg.cpv + ".tbz2." + str(os.getpid()))
+               self._binpkg_tmpfile = binpkg_tmpfile
                settings["PORTAGE_BINPKG_TMPFILE"] = binpkg_tmpfile
                settings.backup_changes("PORTAGE_BINPKG_TMPFILE")
 
-               # Earlier phases should already be done, so
-               # use "noauto" to quietly skip them.
-               settings.features.append("noauto")
-
                try:
-                       retval = portage.doebuild(ebuild_path,
-                               "package", root_config.root,
-                               settings, debug, mydbapi=portdb,
-                               tree="porttree")
+                       EbuildPhase.start(self)
                finally:
                        settings.pop("PORTAGE_BINPKG_TMPFILE", None)
-                       try:
-                               settings.features.remove("noauto")
-                       except ValueError:
-                               pass
 
-               if retval == os.EX_OK:
-                       bintree.inject(pkg.cpv, filename=binpkg_tmpfile)
+       def _set_returncode(self, wait_retval):
+               EbuildPhase._set_returncode(self, wait_retval)
 
-               return retval
+               pkg = self.pkg
+               bintree = pkg.root_config.trees["bintree"]
+               binpkg_tmpfile = self._binpkg_tmpfile
+               if self.returncode == os.EX_OK:
+                       bintree.inject(pkg.cpv, filename=binpkg_tmpfile)
 
 class EbuildMerge(SlotObject):
 
@@ -2133,7 +2307,7 @@ class PackageUninstall(Task):
                        return e.status
                return os.EX_OK
 
-class Binpkg(SlotObject):
+class Binpkg(EbuildBuildDir):
 
        __slots__ = ("find_blockers",
                "ldpath_mtimes", "logger", "opts",
@@ -2152,8 +2326,12 @@ class Binpkg(SlotObject):
                settings = self.settings
                world_atom = self.world_atom
                tree = "bintree"
+               root_config = pkg.root_config
+               bintree = root_config.trees[tree]
                settings.setcpv(pkg)
                debug = settings.get("PORTAGE_DEBUG") == "1"
+               verify = "strict" in settings.features and \
+                       not opts.pretend
 
                # The prefetcher has already completed or it
                # could be running now. If it's running now,
@@ -2183,23 +2361,24 @@ class Binpkg(SlotObject):
                                                for line in wrap(waiting_msg, 65))
                                        writemsg(waiting_msg, noiselevel=-1)
 
-                               while retval is None:
-                                       scheduler.schedule()
-                                       retval = prefetcher.poll()
+                                       scheduler.schedule(prefetcher.reg_id)
+                                       retval = prefetcher.wait()
                        del prefetcher
 
-               fetcher = BinpkgFetcher(pkg=pkg, pretend=opts.pretend,
-                       use_locks=("distlocks" in settings.features))
+               fetcher = BinpkgFetcher(pkg=pkg, scheduler=scheduler)
                pkg_path = fetcher.pkg_path
 
-               if opts.getbinpkg:
-                       retval = fetcher.execute()
-                       if fetcher.remote:
-                               msg = " --- (%s of %s) Fetching Binary (%s::%s)" %\
-                                       (pkg_count.curval, pkg_count.maxval, pkg.cpv, pkg_path)
-                               short_msg = "emerge: (%s of %s) %s Fetch" % \
-                                       (pkg_count.curval, pkg_count.maxval, pkg.cpv)
-                               logger.log(msg, short_msg=short_msg)
+               if opts.getbinpkg and bintree.isremote(pkg.cpv):
+
+                       msg = " --- (%s of %s) Fetching Binary (%s::%s)" %\
+                               (pkg_count.curval, pkg_count.maxval, pkg.cpv, pkg_path)
+                       short_msg = "emerge: (%s of %s) %s Fetch" % \
+                               (pkg_count.curval, pkg_count.maxval, pkg.cpv)
+                       logger.log(msg, short_msg=short_msg)
+
+                       fetcher.start()
+                       scheduler.schedule(fetcher.reg_id)
+                       retval = fetcher.wait()
 
                        if retval != os.EX_OK:
                                return retval
@@ -2207,6 +2386,13 @@ class Binpkg(SlotObject):
                if opts.fetchonly:
                        return os.EX_OK
 
+               if verify:
+                       verifier = BinpkgVerifier(pkg=pkg)
+                       verifier.start()
+                       retval = verifier.wait()
+                       if retval != os.EX_OK:
+                               return retval
+
                msg = " === (%s of %s) Merging Binary (%s::%s)" % \
                        (pkg_count.curval, pkg_count.maxval, pkg.cpv, pkg_path)
                short_msg = "emerge: (%s of %s) %s Merge Binary" % \
@@ -2215,8 +2401,6 @@ class Binpkg(SlotObject):
 
                dir_path = os.path.join(settings["PORTAGE_TMPDIR"],
                        "portage", pkg.category, pkg.pf)
-               build_dir = EbuildBuildDir(dir_path=dir_path,
-                       pkg=pkg, settings=settings)
                image_dir = os.path.join(dir_path, "image")
                infloc = os.path.join(dir_path, "build-info")
 
@@ -2227,16 +2411,22 @@ class Binpkg(SlotObject):
                }
 
                try:
-                       build_dir.lock()
+                       self.lock()
 
                        root_config = self.pkg.root_config
                        ebuild_path = os.path.join(infloc, pkg.pf + ".ebuild")
                        cleanup = 1
                        mydbapi = root_config.trees[tree].dbapi
 
-                       retval = portage.doebuild(ebuild_path, "clean",
-                               root_config.root, settings, debug, cleanup=cleanup,
-                               mydbapi=mydbapi, tree=tree)
+                       phase = "clean"
+                       ebuild_phase = EbuildPhase(fd_pipes=fd_pipes,
+                               pkg=pkg, phase=phase, scheduler=scheduler,
+                               settings=settings, tree=tree)
+
+                       ebuild_phase.start()
+                       scheduler.schedule(ebuild_phase.reg_id)
+                       retval = ebuild_phase.wait()
+
                        if retval != os.EX_OK:
                                return retval
 
@@ -2294,27 +2484,22 @@ class Binpkg(SlotObject):
 
                        phase = "setup"
                        ebuild_phase = EbuildPhase(fd_pipes=fd_pipes,
-                               pkg=pkg, phase=phase, register=scheduler.register,
-                               settings=settings, tree=tree, unregister=scheduler.unregister)
+                               pkg=pkg, phase=phase, scheduler=scheduler,
+                               settings=settings, tree=tree)
 
                        ebuild_phase.start()
-                       retval = None
-                       while retval is None:
-                               scheduler.schedule()
-                               retval = ebuild_phase.poll()
+                       scheduler.schedule(ebuild_phase.reg_id)
+                       retval = ebuild_phase.wait()
 
                        if retval != os.EX_OK:
                                return retval
 
                        extractor = BinpkgExtractorAsync(image_dir=image_dir,
-                               pkg=pkg, pkg_path=pkg_path, register=scheduler.register,
-                               unregister=scheduler.unregister)
+                               pkg=pkg, pkg_path=pkg_path, scheduler=scheduler)
                        portage.writemsg_stdout(">>> Extracting %s\n" % pkg.cpv)
                        extractor.start()
-                       retval = None
-                       while retval is None:
-                               scheduler.schedule()
-                               retval = extractor.poll()
+                       scheduler.schedule(extractor.reg_id)
+                       retval = extractor.wait()
 
                        if retval != os.EX_OK:
                                writemsg("!!! Error Extracting '%s'\n" % pkg_path,
@@ -2332,69 +2517,10 @@ class Binpkg(SlotObject):
 
                finally:
                        settings.pop("PORTAGE_BINPKG_FILE", None)
-                       build_dir.unlock()
+                       self.unlock()
                return os.EX_OK
 
-class BinpkgFetcher(Task):
-
-       __slots__ = ("use_locks", "pkg", "pretend",
-       "pkg_path", "remote")
-
-       def __init__(self, **kwargs):
-               Task.__init__(self, **kwargs)
-               pkg = self.pkg
-               self.pkg_path = pkg.root_config.trees["bintree"].getname(pkg.cpv)
-
-       def _get_hash_key(self):
-               hash_key = getattr(self, "_hash_key", None)
-               if hash_key is None:
-                       self._hash_key = ("BinpkgFetcher", self.pkg._get_hash_key())
-               return self._hash_key
-
-       def execute(self):
-               tbz2_lock = None
-               use_locks = self.use_locks
-               pkg = self.pkg
-               pretend = self.pretend
-               bintree = pkg.root_config.trees["bintree"]
-               pkgdir = bintree.pkgdir
-               pkg_path = self.pkg_path
-               rval = os.EX_OK
-
-               try:
-                       try:
-                               if not pretend and use_locks and os.access(pkgdir, os.W_OK):
-                                       portage.util.ensure_dirs(os.path.dirname(pkg_path))
-                                       tbz2_lock = portage.locks.lockfile(pkg_path,
-                                               wantnewlockfile=1)
-                               if bintree.isremote(pkg.cpv):
-                                       self.remote = True
-                                       if not pretend:
-                                               bintree.gettbz2(pkg.cpv)
-                       finally:
-                               if tbz2_lock is not None:
-                                       portage.locks.unlockfile(tbz2_lock)
-               except portage.exception.FileNotFound:
-                       writemsg("!!! Fetching Binary failed " + \
-                               "for '%s'\n" % pkg.cpv, noiselevel=-1)
-                       rval = 1
-               except portage.exception.DigestException, e:
-                       writemsg("\n!!! Digest verification failed:\n",
-                               noiselevel=-1)
-                       writemsg("!!! %s\n" % e.value[0],
-                               noiselevel=-1)
-                       writemsg("!!! Reason: %s\n" % e.value[1],
-                               noiselevel=-1)
-                       writemsg("!!! Got: %s\n" % e.value[2],
-                               noiselevel=-1)
-                       writemsg("!!! Expected: %s\n" % e.value[3],
-                               noiselevel=-1)
-                       if not pretend:
-                               os.unlink(pkg_path)
-                       rval = 1
-               return rval
-
-class BinpkgFetcherAsync(SpawnProcess):
+class BinpkgFetcher(SpawnProcess):
 
        __slots__ = ("pkg",
                "locked", "pkg_path", "_lock_obj")
@@ -2483,6 +2609,47 @@ class BinpkgFetcherAsync(SpawnProcess):
                self._lock_obj = None
                self.locked = False
 
+class BinpkgVerifier(AsynchronousTask):
+       __slots__ = ("pkg",)
+
+       def start(self):
+               """
+               Note: Unlike a normal AsynchronousTask.start() method,
+               this one does all work is synchronously. The returncode
+               attribute will be set before it returns.
+               """
+
+               pkg = self.pkg
+               root_config = pkg.root_config
+               bintree = root_config.trees["bintree"]
+               rval = os.EX_OK
+               try:
+                       bintree.digestCheck(pkg)
+               except portage.exception.FileNotFound:
+                       writemsg("!!! Fetching Binary failed " + \
+                               "for '%s'\n" % pkg.cpv, noiselevel=-1)
+                       rval = 1
+               except portage.exception.DigestException, e:
+                       writemsg("\n!!! Digest verification failed:\n",
+                               noiselevel=-1)
+                       writemsg("!!! %s\n" % e.value[0],
+                               noiselevel=-1)
+                       writemsg("!!! Reason: %s\n" % e.value[1],
+                               noiselevel=-1)
+                       writemsg("!!! Got: %s\n" % e.value[2],
+                               noiselevel=-1)
+                       writemsg("!!! Expected: %s\n" % e.value[3],
+                               noiselevel=-1)
+                       rval = 1
+
+               self.returncode = rval
+
+       def cancel(self):
+               self.cancelled = True
+
+       def poll(self):
+               return self.returncode
+
 class BinpkgExtractorAsync(SpawnProcess):
 
        __slots__ = ("image_dir", "pkg", "pkg_path")
@@ -2498,6 +2665,88 @@ class BinpkgExtractorAsync(SpawnProcess):
                self.env = self.pkg.root_config.settings.environ()
                SpawnProcess.start(self)
 
+class MergeListItem(SlotObject):
+
+       """
+       TODO: For parallel scheduling, everything here needs asynchronous
+       execution support (start, poll, and wait methods).
+       """
+
+       __slots__ = ("args_set", "binpkg_opts", "build_opts", "emerge_opts",
+               "failed_fetches", "find_blockers", "logger", "mtimedb", "pkg",
+               "pkg_count", "prefetcher", "scheduler", "settings", "world_atom")
+
+       def execute(self):
+
+               args_set = self.args_set
+               binpkg_opts = self.binpkg_opts
+               build_opts = self.build_opts
+               emerge_opts = self.emerge_opts
+               failed_fetches = self.failed_fetches
+               find_blockers = self.find_blockers
+               logger = self.logger
+               mtimedb = self.mtimedb
+               pkg = self.pkg
+               pkg_count = self.pkg_count
+               prefetcher = self.prefetcher
+               scheduler = self.scheduler
+               settings = self.settings
+               world_atom = self.world_atom
+               ldpath_mtimes = mtimedb["ldpath"]
+
+               if pkg.installed:
+                       if not (build_opts.buildpkgonly or \
+                               build_opts.fetchonly or build_opts.pretend):
+
+                               uninstall = PackageUninstall(ldpath_mtimes=ldpath_mtimes,
+                                       opts=emerge_opts, pkg=pkg, settings=settings)
+
+                               retval = uninstall.execute()
+                               if retval != os.EX_OK:
+                                       return retval
+
+                       return os.EX_OK
+
+               if not build_opts.pretend:
+                       portage.writemsg_stdout(
+                               "\n>>> Emerging (%s of %s) %s to %s\n" % \
+                               (colorize("MERGE_LIST_PROGRESS", str(pkg_count.curval)),
+                               colorize("MERGE_LIST_PROGRESS", str(pkg_count.maxval)),
+                               colorize("GOOD", pkg.cpv), pkg.root), noiselevel=-1)
+                       logger.log(" >>> emerge (%s of %s) %s to %s" % \
+                               (pkg_count.curval, pkg_count.maxval, pkg.cpv, pkg.root))
+
+               if pkg.type_name == "ebuild":
+
+                       build = EbuildBuild(args_set=args_set,
+                               find_blockers=find_blockers,
+                               ldpath_mtimes=ldpath_mtimes, logger=logger,
+                               opts=build_opts, pkg=pkg, pkg_count=pkg_count,
+                               settings=settings, scheduler=scheduler,
+                               world_atom=world_atom)
+
+                       retval = build.execute()
+
+                       if retval != os.EX_OK:
+                               if build_opts.fetchonly:
+                                       failed_fetches.append(pkg.cpv)
+                               return retval
+
+               elif pkg.type_name == "binary":
+
+                       binpkg = Binpkg(find_blockers=find_blockers,
+                               ldpath_mtimes=ldpath_mtimes, logger=logger,
+                               opts=binpkg_opts, pkg=pkg, pkg_count=pkg_count,
+                               prefetcher=prefetcher, settings=settings,
+                               scheduler=scheduler, world_atom=world_atom)
+
+                       retval = binpkg.execute()
+
+                       if retval != os.EX_OK:
+                               return retval
+
+               return os.EX_OK
+
 class DependencyArg(object):
        def __init__(self, arg=None, root_config=None):
                self.arg = arg
@@ -6952,7 +7201,21 @@ class PackageCounters(object):
                                        (self.blocks - self.blocks_satisfied))
                return "".join(myoutput)
 
-class PollSelectFallback(object):
+class PollConstants(object):
+
+       """
+       Provides POLL* constants that are equivalent to those from the
+       select module, for use by PollSelectAdapter.
+       """
+
+       names = ("POLLIN", "POLLPRI", "POLLOUT", "POLLERR", "POLLHUP", "POLLNVAL")
+       v = 1
+       for k in names:
+               locals()[k] = getattr(select, k, v)
+               v *= 2
+       del k, v
+
+class PollSelectAdapter(PollConstants):
 
        """
        Use select to emulate a poll object, for
@@ -6965,14 +7228,15 @@ class PollSelectFallback(object):
 
        def register(self, fd, *args):
                """
-               Only select.POLLIN is currently supported!
+               Only POLLIN is currently supported!
                """
                if len(args) > 1:
                        raise TypeError(
                                "register expected at most 2 arguments, got " + \
                                repr(1 + len(args)))
 
-               eventmask = select.POLLIN | select.POLLPRI | select.POLLOUT
+               eventmask = PollConstants.POLLIN | \
+                       PollConstants.POLLPRI | PollConstants.POLLOUT
                if args:
                        eventmask = args[0]
 
@@ -7004,9 +7268,51 @@ class PollSelectFallback(object):
                select_events = select.select(*select_args)
                poll_events = []
                for fd in select_events[0]:
-                       poll_events.append((fd, select.POLLIN))
+                       poll_events.append((fd, PollConstants.POLLIN))
                return poll_events
 
+class SequentialTaskQueue(SlotObject):
+
+       __slots__ = ("max_jobs", "running_tasks", "_task_queue")
+
+       def __init__(self, **kwargs):
+               SlotObject.__init__(self, **kwargs)
+               self._task_queue = deque()
+               self.running_tasks = set()
+               if self.max_jobs is None:
+                       self.max_jobs = 1
+
+       def add(self, task):
+               self._task_queue.append(task)
+
+       def schedule(self):
+               task_queue = self._task_queue
+               running_tasks = self.running_tasks
+               max_jobs = self.max_jobs
+               state_changed = False
+
+               for task in list(running_tasks):
+                       if not task.registered and task.poll() is not None:
+                               running_tasks.remove(task)
+                               state_changed = True
+
+               while task_queue and (len(running_tasks) < max_jobs):
+                       task = task_queue.popleft()
+                       cancelled = getattr(task, "cancelled", None)
+                       if not cancelled:
+                               task.start()
+                               running_tasks.add(task)
+                       state_changed = True
+
+               return state_changed
+
+       def clear(self):
+               self._task_queue.clear()
+               running_tasks = self.running_tasks
+               while running_tasks:
+                       task = running_tasks.pop()
+                       task.cancel()
+
 class Scheduler(object):
 
        _opts_ignore_blockers = \
@@ -7020,7 +7326,7 @@ class Scheduler(object):
        _fetch_log = EPREFIX + "/var/log/emerge-fetch.log"
 
        class _iface_class(SlotObject):
-               __slots__ = ("register", "schedule", "unregister")
+               __slots__ = ("register", "schedule")
 
        class _build_opts_class(SlotObject):
                __slots__ = ("buildpkg", "buildpkgonly",
@@ -7070,24 +7376,32 @@ class Scheduler(object):
                self._logger = self._emerge_log_class(
                        xterm_titles=("notitles" not in settings.features))
                self._sched_iface = self._iface_class(
-                       register=self._register, schedule=self._schedule,
-                               unregister=self._unregister)
+                       register=self._register, schedule=self._schedule)
                self._poll_event_handlers = {}
+               self._poll_event_handler_ids = {}
+               # Increment id for each new handler.
+               self._event_handler_id = 0
 
                try:
                        self._poll = select.poll()
                except AttributeError:
-                       self._poll = PollSelectFallback()
+                       self._poll = PollSelectAdapter()
 
-               from collections import deque
-               self._task_queue = deque()
-               self._running_tasks = set()
-               self._max_jobs = 1
+               self._task_queues = slot_dict_class(("build", "prefetch"), prefix="")
+               for k in self._task_queues.allowed_keys:
+                       setattr(self._task_queues, k, SequentialTaskQueue())
+
+               self._add_task = self._task_queues.prefetch.add
+               self._schedule_tasks = self._task_queues.prefetch.schedule
                self._prefetchers = weakref.WeakValueDictionary()
+               self._pkg_queue = deque()
+               self._failed_pkgs = []
                self._failed_fetches = []
                self._parallel_fetch = False
+               merge_count = len([x for x in mergelist \
+                       if isinstance(x, Package) and x.operation == "merge"])
                self._pkg_count = self._pkg_count_class(
-                       curval=0, maxval=len(mergelist))
+                       curval=0, maxval=merge_count)
 
                features = self.settings.features
                if "parallel-fetch" in features and \
@@ -7112,9 +7426,6 @@ class Scheduler(object):
                                except EnvironmentError:
                                        pass
 
-       def _add_task(self, task):
-               self._task_queue.append(task)
-
        class _pkg_failure(portage.exception.PortageException):
                """
                An instance of this class is raised by unmerge() when
@@ -7181,7 +7492,8 @@ class Scheduler(object):
                        del quiet_config
 
                for x in self._mergelist:
-                       if x.type_name != "ebuild":
+                       if not isinstance(x, Package) or \
+                               x.type_name != "ebuild":
                                continue
 
                        if not shown_verifying_msg:
@@ -7225,15 +7537,15 @@ class Scheduler(object):
 
                elif pkg.type_name == "ebuild":
 
-                       prefetcher = EbuildFetcherAsync(logfile=self._fetch_log, pkg=pkg,
-                               register=self._register, unregister=self._unregister)
+                       prefetcher = EbuildFetcher(logfile=self._fetch_log, pkg=pkg,
+                               scheduler=self._sched_iface)
 
                elif pkg.type_name == "binary" and \
                        "--getbinpkg" in self.myopts and \
                        pkg.root_config.trees["bintree"].isremote(pkg.cpv):
 
-                       prefetcher = BinpkgFetcherAsync(logfile=self._fetch_log,
-                               pkg=pkg, register=self._register, unregister=self._unregister)
+                       prefetcher = BinpkgFetcher(logfile=self._fetch_log,
+                               pkg=pkg, scheduler=self._sched_iface)
 
                return prefetcher
 
@@ -7324,7 +7636,7 @@ class Scheduler(object):
                mtimedb = self._mtimedb
 
                while True:
-                       self._merge()
+                       rval = self._merge()
                        self._show_failed_fetches()
                        del self._failed_fetches[:]
 
@@ -7374,30 +7686,100 @@ class Scheduler(object):
                        self._mergelist = mylist
                        self._save_resume_list()
                        self._pkg_count.curval = 0
-                       self._pkg_count.maxval = len(mylist)
+                       self._pkg_count.maxval = len([x for x in mylist \
+                               if isinstance(x, Package) and x.operation == "merge"])
 
                self._logger.log(" *** Finished. Cleaning up...")
 
                return rval
 
+       def _add_packages(self):
+               pkg_queue = self._pkg_queue
+               for pkg in self._mergelist:
+                       if isinstance(pkg, Package):
+                               pkg_queue.append(pkg)
+                       elif isinstance(pkg, Blocker):
+                               pass
+
+       def _choose_pkg(self):
+               return self._pkg_queue.popleft()
+
+       def _main_loop(self):
+
+               pkg_queue = self._pkg_queue
+
+               while pkg_queue:
+                       pkg = self._choose_pkg()
+                       retval = self._execute_pkg(pkg)
+
+                       if retval != os.EX_OK:
+                               self._failed_pkgs.append((pkg, retval))
+                               if not self._build_opts.fetchonly:
+                                       return
+
+                       if pkg.installed:
+                               continue
+
+                       self._restart_if_necessary(pkg)
+
+                       # Call mtimedb.commit() after each merge so that
+                       # --resume still works after being interrupted
+                       # by reboot, sigkill or similar.
+                       mtimedb = self._mtimedb
+                       del mtimedb["resume"]["mergelist"][0]
+                       if not mtimedb["resume"]["mergelist"]:
+                               del mtimedb["resume"]
+                       mtimedb.commit()
+
        def _merge(self):
 
                self._add_prefetchers()
+               self._add_packages()
+               pkg_queue = self._pkg_queue
+               failed_pkgs = self._failed_pkgs
+               rval = os.EX_OK
 
                try:
-                       for task in self._mergelist:
-                               try:
-                                       self._execute_task(task)
-                               except self._pkg_failure, e:
-                                       return e.status
+                       self._main_loop()
                finally:
+                       # discard remaining packages if necessary
+                       pkg_queue.clear()
+
                        # clean up child process if necessary
-                       self._task_queue.clear()
-                       running_tasks = self._running_tasks
-                       while running_tasks:
-                               task = running_tasks.pop()
-                               task.cancel()
-               return os.EX_OK
+                       self._task_queues.prefetch.clear()
+
+                       # discard any failures and return the
+                       # exist status of the last one
+                       if failed_pkgs:
+                               pkg, rval = failed_pkgs[-1]
+
+                       del failed_pkgs[:]
+
+               return rval
+
+       def _execute_pkg(self, pkg):
+
+               if not pkg.installed:
+                       self._pkg_count.curval += 1
+
+               merge = MergeListItem(args_set=self._args_set,
+                       binpkg_opts=self._binpkg_opts,
+                       build_opts=self._build_opts,
+                       emerge_opts=self.myopts,
+                       failed_fetches=self._failed_fetches,
+                       find_blockers=self._find_blockers(pkg), logger=self._logger,
+                       mtimedb=self._mtimedb, pkg=pkg, pkg_count=self._pkg_count,
+                       prefetcher=self._prefetchers.get(pkg),
+                       scheduler=self._sched_iface,
+                       settings=self.pkgsettings[pkg.root],
+                       world_atom=self._world_atom)
+
+               retval = merge.execute()
+
+               if retval == os.EX_OK:
+                       self.curval += 1
+
+               return retval
 
        def _save_resume_list(self):
                """
@@ -7419,6 +7801,17 @@ class Scheduler(object):
                """
                print colorize("GOOD", "*** Resuming merge...")
 
+               if self._show_list():
+                       if "--tree" in self.myopts:
+                               portage.writemsg_stdout("\n" + \
+                                       darkgreen("These are the packages that " + \
+                                       "would be merged, in reverse order:\n\n"))
+
+                       else:
+                               portage.writemsg_stdout("\n" + \
+                                       darkgreen("These are the packages that " + \
+                                       "would be merged, in order:\n\n"))
+
                show_spinner = "--quiet" not in self.myopts and \
                        "--nodeps" not in self.myopts
 
@@ -7433,8 +7826,14 @@ class Scheduler(object):
                if show_spinner:
                        print "\b\b... done!"
 
+               if self._show_list():
+                       mylist = mydepgraph.altlist()
+                       if "--tree" in self.myopts:
+                               mylist.reverse()
+                       mydepgraph.display(mylist, favorites=self._favorites)
+
+               mydepgraph.display_problems()
                if not success:
-                       mydepgraph.display_problems()
                        return (None, None)
 
                mylist = mydepgraph.altlist()
@@ -7442,58 +7841,65 @@ class Scheduler(object):
                mydepgraph.break_refs(dropped_tasks)
                return (mylist, dropped_tasks)
 
+       def _show_list(self):
+               myopts = self.myopts
+               if "--quiet" not in myopts and \
+                       ("--ask" in myopts or "--tree" in myopts or \
+                       "--verbose" in myopts):
+                       return True
+               return False
+
        def _register(self, f, eventmask, handler):
-               self._poll_event_handlers[f] = handler
+               """
+               @rtype: Integer
+               @return: A unique registration id, for use in schedule() or
+                       unregister() calls.
+               """
+               self._event_handler_id += 1
+               reg_id = self._event_handler_id
+               self._poll_event_handler_ids[reg_id] = f
+               self._poll_event_handlers[f] = (handler, reg_id)
                self._poll.register(f, eventmask)
+               return reg_id
 
-       def _unregister(self, f):
+       def _unregister(self, reg_id):
+               f = self._poll_event_handler_ids[reg_id]
                self._poll.unregister(f)
                del self._poll_event_handlers[f]
+               del self._poll_event_handler_ids[reg_id]
                self._schedule_tasks()
 
-       def _schedule(self):
+       def _schedule(self, wait_id):
+               """
+               Schedule until wait_id is not longer registered
+               for poll() events.
+               @type wait_id: int
+               @param wait_id: a task id to wait for
+               """
                event_handlers = self._poll_event_handlers
-               running_tasks = self._running_tasks
+               handler_ids = self._poll_event_handler_ids
                poll = self._poll.poll
 
                self._schedule_tasks()
 
-               while event_handlers:
+               while wait_id in handler_ids:
                        for f, event in poll():
-                               event_handlers[f](f, event)
-
-                       if len(event_handlers) <= len(running_tasks):
-                               # Assuming one handler per task, this
-                               # means the caller has unregistered it's
-                               # handler, so it's time to yield.
-                               break
-
-       def _schedule_tasks(self):
-               task_queue = self._task_queue
-               running_tasks = self._running_tasks
-               max_jobs = self._max_jobs
-               state_changed = False
-
-               for task in list(running_tasks):
-                       if not task.registered and task.poll() is not None:
-                               running_tasks.remove(task)
-                               state_changed = True
-
-               while task_queue and (len(running_tasks) < max_jobs):
-                       task = task_queue.popleft()
-                       cancelled = getattr(task, "cancelled", None)
-                       if not cancelled:
-                               task.start()
-                               running_tasks.add(task)
-                       state_changed = True
-
-               return state_changed
+                               handler, reg_id = event_handlers[f]
+                               if not handler(f, event):
+                                       self._unregister(reg_id)
 
        def _world_atom(self, pkg):
                """
                Add the package to the world file, but only if
                it's supposed to be added. Otherwise, do nothing.
                """
+
+               if set(("--buildpkgonly", "--fetchonly",
+                       "--fetch-all-uri",
+                       "--oneshot", "--onlydeps",
+                       "--pretend")).intersection(self.myopts):
+                       return
+
                if pkg.root != self.target_root:
                        return
 
@@ -7518,80 +7924,6 @@ class Scheduler(object):
                finally:
                        world_set.unlock()
 
-       def _execute_task(self, pkg):
-
-                       buildpkgonly = "--buildpkgonly" in self.myopts
-                       fetch_all = "--fetch-all-uri" in self.myopts
-                       fetchonly = fetch_all or "--fetchonly" in self.myopts
-                       pretend = "--pretend" in self.myopts
-
-                       pkgsettings = self.pkgsettings[pkg.root]
-                       mtimedb = self._mtimedb
-                       ldpath_mtimes = mtimedb["ldpath"]
-                       failed_fetches = self._failed_fetches
-                       pkg_count = self._pkg_count
-                       prefetchers = self._prefetchers
-
-                       if not pkg.installed:
-                               pkg_count.curval += 1
-                               mergecount = pkg_count.curval
-                       else:
-                               if not (buildpkgonly or fetchonly or pretend):
-                                       uninstall = PackageUninstall(ldpath_mtimes=ldpath_mtimes,
-                                               opts=self.myopts, pkg=pkg, settings=pkgsettings)
-                                       retval = uninstall.execute()
-                                       if retval != os.EX_OK:
-                                               raise self._pkg_failure(retval)
-                               return
-
-                       if not pretend:
-                               portage.writemsg_stdout(
-                                       "\n>>> Emerging (%s of %s) %s to %s\n" % \
-                                       (colorize("MERGE_LIST_PROGRESS", str(pkg_count.curval)),
-                                       colorize("MERGE_LIST_PROGRESS", str(pkg_count.maxval)),
-                                       colorize("GOOD", pkg.cpv), pkg.root), noiselevel=-1)
-                               self._logger.log(" >>> emerge (%s of %s) %s to %s" % \
-                                       (pkg_count.curval, pkg_count.maxval, pkg.cpv, pkg.root))
-
-                       self._schedule()
-
-                       if pkg.type_name == "ebuild":
-                               build = EbuildBuild(args_set=self._args_set,
-                                       find_blockers=self._find_blockers(pkg),
-                                       ldpath_mtimes=ldpath_mtimes, logger=self._logger,
-                                       opts=self._build_opts, pkg=pkg, pkg_count=pkg_count,
-                                       settings=pkgsettings, scheduler=self._sched_iface,
-                                       world_atom=self._world_atom)
-                               retval = build.execute()
-                               if retval != os.EX_OK:
-                                       if fetchonly:
-                                               failed_fetches.append(pkg.cpv)
-                                       else:
-                                               raise self._pkg_failure(retval)
-
-                       elif pkg.type_name == "binary":
-                               binpkg = Binpkg(find_blockers=self._find_blockers(pkg),
-                                       ldpath_mtimes=ldpath_mtimes, logger=self._logger,
-                                       opts=self._binpkg_opts, pkg=pkg, pkg_count=pkg_count,
-                                       prefetcher=prefetchers.get(pkg), settings=pkgsettings,
-                                       scheduler=self._sched_iface, world_atom=self._world_atom)
-                               retval = binpkg.execute()
-                               if retval != os.EX_OK:
-                                       if fetchonly:
-                                               failed_fetches.append(pkg.cpv)
-                                       else:
-                                               raise self._pkg_failure(retval)
-
-                       self._restart_if_necessary(pkg)
-                       del mtimedb["resume"]["mergelist"][0]
-                       if not mtimedb["resume"]["mergelist"]:
-                               del mtimedb["resume"]
-                       # Commit after each merge so that --resume may still work in
-                       # in the event that portage is not allowed to exit normally
-                       # due to power failure, SIGKILL, etc...
-                       mtimedb.commit()
-                       self.curval += 1
-
 class UninstallFailure(portage.exception.PortageException):
        """
        An instance of this class is raised by unmerge() when
index c47a58eca783b241ecb3efdfc336ae6c719c890c..6263d749dfbd30fe2f84d34f5d1d80c3d0e3c847 100644 (file)
@@ -3163,24 +3163,25 @@ def _checksum_failure_temp_file(distdir, basename):
        os.rename(filename, temp_filename)
        return temp_filename
 
-def _check_digests(filename, digests):
+def _check_digests(filename, digests, show_errors=1):
        """
        Check digests and display a message if an error occurs.
        @return True if all digests match, False otherwise.
        """
        verified_ok, reason = portage.checksum.verify_all(filename, digests)
        if not verified_ok:
-               writemsg("!!! Previously fetched" + \
-                       " file: '%s'\n" % filename, noiselevel=-1)
-               writemsg("!!! Reason: %s\n" % reason[0],
-                       noiselevel=-1)
-               writemsg(("!!! Got:      %s\n" + \
-                       "!!! Expected: %s\n") % \
-                       (reason[1], reason[2]), noiselevel=-1)
+               if show_errors:
+                       writemsg("!!! Previously fetched" + \
+                               " file: '%s'\n" % filename, noiselevel=-1)
+                       writemsg("!!! Reason: %s\n" % reason[0],
+                               noiselevel=-1)
+                       writemsg(("!!! Got:      %s\n" + \
+                               "!!! Expected: %s\n") % \
+                               (reason[1], reason[2]), noiselevel=-1)
                return False
        return True
 
-def _check_distfile(filename, digests, eout):
+def _check_distfile(filename, digests, eout, show_errors=1):
        """
        @return a tuple of (match, stat_obj) where match is True if filename
        matches all given digests (if any) and stat_obj is a stat result, or
@@ -3203,7 +3204,7 @@ def _check_distfile(filename, digests, eout):
                        eout.ebegin("%s %s ;-)" % (os.path.basename(filename), "size"))
                        eout.eend(0)
        else:
-               if _check_digests(filename, digests):
+               if _check_digests(filename, digests, show_errors=show_errors):
                        eout.ebegin("%s %s ;-)" % (os.path.basename(filename),
                                " ".join(sorted(digests))))
                        eout.eend(0)
@@ -3532,15 +3533,27 @@ def fetch(myuris, mysettings, listonly=0, fetchonly=0, locks_in_subdir=".locks",
                                        from textwrap import wrap
                                        waiting_msg = "\n".join(msg_prefix + line \
                                                for line in wrap(waiting_msg, 65))
+
                                if locks_in_subdir:
-                                       file_lock = portage.locks.lockfile(
-                                               os.path.join(mysettings["DISTDIR"],
-                                               locks_in_subdir, myfile), wantnewlockfile=1,
-                                               waiting_msg=waiting_msg)
+                                       lock_file = os.path.join(mysettings["DISTDIR"],
+                                               locks_in_subdir, myfile)
+                               else:
+                                       lock_file = myfile_path
+
+                               lock_kwargs = {}
+                               if fetchonly:
+                                       lock_kwargs["flags"] = os.O_NONBLOCK
                                else:
-                                       file_lock = portage.locks.lockfile(
-                                               myfile_path, wantnewlockfile=1,
-                                               waiting_msg=waiting_msg)
+                                       lock_kwargs["waiting_msg"] = waiting_msg
+
+                               try:
+                                       file_lock = portage.locks.lockfile(myfile_path,
+                                               wantnewlockfile=1, **lock_kwargs)
+                               except portage.exception.TryAgain:
+                                       writemsg((">>> File '%s' is already locked by " + \
+                                               "another fetcher. Continuing...\n") % myfile,
+                                               noiselevel=-1)
+                                       continue
                try:
                        if not listonly:
 
@@ -5039,7 +5052,8 @@ def doebuild(myebuild, mydo, myroot, mysettings, debug=0, listonly=0,
 
                if mydo in clean_phases:
                        retval = spawn(_shell_quote(ebuild_sh_binary) + " clean",
-                               mysettings, debug=debug, free=1, logfile=None)
+                               mysettings, debug=debug, fd_pipes=fd_pipes, free=1,
+                               logfile=None, returnpid=returnpid)
                        return retval
 
                # get possible slot information from the deps file
index dd6ab5aadf42e541c086ce48c97e7dcd1d72bf70..6b3c392bb82394c0912fa1cb0b15e9f3076b5931 100644 (file)
@@ -7,13 +7,13 @@ from portage.dep import isvalidatom, isjustname, dep_getkey, match_from_list
 from portage.dbapi.virtual import fakedbapi
 from portage.exception import InvalidPackageName, InvalidAtom, \
        PermissionDenied, PortageException
-from portage.output import green
+from portage.output import green, EOutput
 from portage.util import ensure_dirs, normalize_path, writemsg, writemsg_stdout
 from portage.versions import best, catpkgsplit, catsplit
 from portage.update import update_dbentries
 from portage.const import EAPI, EAPIPREFIX
 
-from portage import dep_expand, listdir, _movefile
+from portage import dep_expand, listdir, _check_distfile, _movefile
 
 import portage.xpak, portage.getbinpkg
 
@@ -1001,26 +1001,6 @@ class binarytree(object):
                if not fcmd:
                        fcmd = self.settings.get(fcmd_prefix)
                success = portage.getbinpkg.file_get(url, mydest, fcmd=fcmd)
-               if success and "strict" in self.settings.features:
-                       metadata = self._remotepkgs[pkgname]
-                       digests = {}
-                       from portage.checksum import hashfunc_map, verify_all
-                       for k in hashfunc_map:
-                               v = metadata.get(k)
-                               if not v:
-                                       continue
-                               digests[k] = v
-                       if "SIZE" in metadata:
-                               try:
-                                       digests["size"] = long(self._remotepkgs[pkgname]["SIZE"])
-                               except ValueError:
-                                       writemsg("!!! Malformed SIZE attribute in remote " + \
-                                       "metadata for '%s'\n" % pkgname)
-                       if digests:
-                               ok, reason = verify_all(tbz2_path, digests)
-                               if not ok:
-                                       raise portage.exception.DigestException(
-                                               tuple([tbz2_path]+list(reason)))
                if not success:
                        try:
                                os.unlink(self.getname(pkgname))
@@ -1029,6 +1009,72 @@ class binarytree(object):
                        raise portage.exception.FileNotFound(mydest)
                self.inject(pkgname)
 
+       def _load_pkgindex(self):
+               pkgindex = self._new_pkgindex()
+               try:
+                       f = open(self._pkgindex_file)
+               except EnvironmentError:
+                       pass
+               else:
+                       try:
+                               pkgindex.read(f)
+                       finally:
+                               f.close()
+               return pkgindex
+
+       def digestCheck(self, pkg):
+               """
+               Verify digests for the given package and raise DigestException
+               if verification fails.
+               @rtype: bool
+               @returns: True if digests could be located, False otherwise.
+               """
+               cpv = pkg
+               if not isinstance(cpv, basestring):
+                       cpv = pkg.cpv
+                       pkg = None
+
+               pkg_path = self.getname(cpv)
+               metadata = None
+               if self._remotepkgs is None or cpv not in self._remotepkgs:
+                       for d in self._load_pkgindex().packages:
+                               if d["CPV"] == cpv:
+                                       metadata = d
+                                       break
+               else:
+                       metadata = self._remotepkgs[cpv]
+               if metadata is None:
+                       return False
+
+               digests = {}
+               from portage.checksum import hashfunc_map, verify_all
+               for k in hashfunc_map:
+                       v = metadata.get(k)
+                       if not v:
+                               continue
+                       digests[k] = v
+
+               if "SIZE" in metadata:
+                       try:
+                               digests["size"] = int(metadata["SIZE"])
+                       except ValueError:
+                               writemsg("!!! Malformed SIZE attribute in remote " + \
+                               "metadata for '%s'\n" % cpv)
+
+               if not digests:
+                       return False
+
+               eout = EOutput()
+               eout.quiet = self.settings.get("PORTAGE_QUIET") == "1"
+               ok, st = _check_distfile(pkg_path, digests, eout, show_errors=0)
+               if not ok:
+                       ok, reason = verify_all(pkg_path, digests)
+                       if not ok:
+                               raise portage.exception.DigestException(
+                                       (pkg_path,) + tuple(reason))
+
+               return True
+
        def getslot(self, mycatpkg):
                "Get a slot for a catpkg; assume it exists."
                myslot = ""
index bd016cedff1f625f3f431ba1a214b254b51fee32..79d49705e8ee81f28cf33d52ec07d915512c6fb4 100644 (file)
@@ -52,8 +52,8 @@ class PreservedLibsRegistry(object):
                """ Reload the registry data from file """
                try:
                        self._data = cPickle.load(open(self._filename, "r"))
-               except IOError, e:
-                       if e.errno == errno.ENOENT:
+               except (EOFError, IOError), e:
+                       if isinstance(e, EOFError) or e.errno == errno.ENOENT:
                                self._data = {}
                        elif e.errno == PermissionDenied.errno:
                                raise PermissionDenied(self._filename)
@@ -611,7 +611,6 @@ class vardbapi(dbapi):
                try:
                        dir_list = os.listdir(cat_dir)
                except EnvironmentError, e:
-                       from portage.exception import PermissionDenied
                        if e.errno == PermissionDenied.errno:
                                raise PermissionDenied(cat_dir)
                        del e
@@ -672,7 +671,7 @@ class vardbapi(dbapi):
                                        if catpkgsplit(subpath) is None:
                                                self.invalidentry(os.path.join(self.root, subpath))
                                                continue
-                               except portage.exception.InvalidData:
+                               except InvalidData:
                                        self.invalidentry(os.path.join(self.root, subpath))
                                        continue
                                returnme.append(subpath)
@@ -686,7 +685,7 @@ class vardbapi(dbapi):
                                y = y[1:]
                        try:
                                mysplit = catpkgsplit(y)
-                       except portage.exception.InvalidData:
+                       except InvalidData:
                                self.invalidentry(self.getpath(y))
                                continue
                        if not mysplit:
index d91584b7366c55b391ecf341292a5f7e93112b72..66cd16244135944e4ec45ad173be528b53997403 100644 (file)
@@ -57,6 +57,10 @@ class PermissionDenied(PortageException):
        from errno import EACCES as errno
        """Permission denied"""
 
+class TryAgain(PortageException):
+       from errno import EAGAIN as errno
+       """Try again"""
+
 class ReadOnlyFileSystem(PortageException):
        """Read-only file system"""
 
index 37fae639e6acbf1d336f5a93b5ad70bfcd45d062..5c88bfcf6a57973be4051306c40792f8378e124d 100644 (file)
@@ -5,7 +5,8 @@
 
 
 import errno, os, stat, time, types
-from portage.exception import InvalidData, DirectoryNotFound, FileNotFound
+from portage.exception import DirectoryNotFound, FileNotFound, \
+       InvalidData, TryAgain
 from portage.data import portage_gid
 from portage.util import writemsg
 from portage.localization import _
@@ -18,7 +19,8 @@ def lockdir(mydir):
 def unlockdir(mylock):
        return unlockfile(mylock)
 
-def lockfile(mypath, wantnewlockfile=0, unlinkfile=0, waiting_msg=None):
+def lockfile(mypath, wantnewlockfile=0, unlinkfile=0,
+       waiting_msg=None, flags=0):
        """Creates all dirs upto, the given dir. Creates a lockfile
        for the given directory as the file: directoryname+'.portage_lockfile'."""
        import fcntl
@@ -55,7 +57,8 @@ def lockfile(mypath, wantnewlockfile=0, unlinkfile=0, waiting_msg=None):
                        except OSError, e:
                                if e[0] == 2: # No such file or directory
                                        return lockfile(mypath, wantnewlockfile=wantnewlockfile,
-                                               unlinkfile=unlinkfile, waiting_msg=waiting_msg)
+                                               unlinkfile=unlinkfile, waiting_msg=waiting_msg,
+                                               flags=flags)
                                else:
                                        writemsg("Cannot chown a lockfile. This could cause inconvenience later.\n");
                        os.umask(old_mask)
@@ -79,6 +82,8 @@ def lockfile(mypath, wantnewlockfile=0, unlinkfile=0, waiting_msg=None):
                        raise
                if e.errno in (errno.EACCES, errno.EAGAIN):
                        # resource temp unavailable; eg, someone beat us to the lock.
+                       if flags & os.O_NONBLOCK:
+                               raise TryAgain(mypath)
                        if waiting_msg is None:
                                if isinstance(mypath, int):
                                        print "waiting for lock on fd %i" % myfd
@@ -115,7 +120,7 @@ def lockfile(mypath, wantnewlockfile=0, unlinkfile=0, waiting_msg=None):
                writemsg("lockfile recurse\n",1)
                lockfilename, myfd, unlinkfile, locking_method = lockfile(
                        mypath, wantnewlockfile=wantnewlockfile, unlinkfile=unlinkfile,
-                       waiting_msg=waiting_msg)
+                       waiting_msg=waiting_msg, flags=flags)
 
        writemsg(str((lockfilename,myfd,unlinkfile))+"\n",1)
        return (lockfilename,myfd,unlinkfile,locking_method)