From b3a20725cc146dbc38f1be037130945a63990307 Mon Sep 17 00:00:00 2001 From: Fabian Groffen Date: Mon, 7 Jul 2008 09:18:40 +0000 Subject: [PATCH] Merged from trunk 10944:10966 | 10945 | Correct TaskSequence docstring to refer to the | | zmedico | addExitListener() method. | | 10946 | Call _wait_hook() from poll() if the wait call occurs there. | | zmedico | | | 10947 | * Fix AsynchronousTask.poll() to call _wait_hook() when | | zmedico | necessary. * Use the default poll() and cancel() | | | implementations for BinpkgVerifier. | | 10948 | Fix typo in xterm titles total package count which causes it | | zmedico | to show the current package instead. Thanks to Arfrever for | | | this patch. | | 10949 | Add a CompositeTask._assert_current() method that | | zmedico | asynchronous callbacks can use detect possible bugs. | | 10950 | Split out a reusable CompositeTask._default_exit() method | | zmedico | that subclasses can use as a generic task exit callback. | | 10951 | Add CompositeTask._final_exit() method and use it to fix | | zmedico | breakage from the previous commit. | | 10952 | Split out a CompositeTask._start_task() for subclasses to | | zmedico | use as a generic way to start a task. | | 10953 | Add async support to EbuildBuild, and an synchronous | | zmedico | install() method. | | 10954 | * Fix broke return code handling from previous commit, in | | zmedico | MergeListItem.execute(). * Fix | | | TaskSequence._task_exit_handler() so it won't call | | | final_exit() if _default_exit() has already set | | | self._current_task to None. | | 10955 | Fix --getbinpkg to inject downloaded packages into the local | | zmedico | bintree. | | 10956 | Make AsynchronousTask subclasses override _wait() and | | zmedico | _poll() so that calls to public methods can be wrapped for | | | implementing hooks such as exit listener notification. | | 10957 | Fix parent class constructor call in the BinpkgFetcher | | zmedico | constructor. | | 10958 | Make BinpkgFetcher send output directly to stdout when | | zmedico | appropriate, so that wget's progress bar works normally. | | 10959 | Add async support to the Binpkg class. | | zmedico | | | 10960 | Add async support to MergeListItem. | | zmedico | | | 10961 | Add a PackageMerge class to serve as an asynchronous | | zmedico | interface to package merges. For now it executes | | | synchronously inside the start() method. | | 10962 | * Implement MergeListItem._poll() and _wait(). * Fix | | zmedico | BinpkgVerifier.start() to call wait() since it's not | | | asynchronous. | | 10963 | Fix typo in Binpkg.start() which prevents --genbinpkg | | zmedico | prefetcher sync from working properly in some cases. | | 10964 | Fix EbuildPhase._set_returncode() so that it correctly | | zmedico | updates the returncode attrbute instead of just a local | | | variable. | | 10965 | Fix broken code in AsynchronousTask.poll(). | | zmedico | | | 10966 | * Implement CompositeTask._poll(). * Make AsynchronousTask | | zmedico | classes call self.wait() to notify exit listeners. * Rewrite | | | Scheduler._main_loop() to bring it closer to allowing | | | parallel build scheduling. | svn path=/main/branches/prefix/; revision=10968 --- pym/_emerge/__init__.py | 1073 +++++++++++++++++++++++++-------------- 1 file changed, 689 insertions(+), 384 deletions(-) diff --git a/pym/_emerge/__init__.py b/pym/_emerge/__init__.py index dd43dc549..8a7cf4992 100644 --- a/pym/_emerge/__init__.py +++ b/pym/_emerge/__init__.py @@ -1482,6 +1482,15 @@ class EbuildFetchPretend(SlotObject): return retval class AsynchronousTask(SlotObject): + """ + Subclasses override _wait() and _poll() so that calls + to public methods can be wrapped for implementing + hooks such as exit listener notification. + + Sublasses should call self.wait() to notify exit listeners after + the task is complete and self.returncode has been set. + """ + __slots__ = ("cancelled", "returncode") + ("_exit_listeners",) def start(self): @@ -1494,14 +1503,24 @@ class AsynchronousTask(SlotObject): return self.returncode is None def poll(self): + self._wait_hook() + return self._poll() + + def _poll(self): return self.returncode def wait(self): + if self.returncode is None: + self._wait() self._wait_hook() return self.returncode + def _wait(self): + return self.returncode + def cancel(self): - pass + self.cancelled = True + self.wait() def addExitListener(self, f): """ @@ -1516,11 +1535,13 @@ class AsynchronousTask(SlotObject): def _wait_hook(self): """ - Call this method before returning from wait. This hook is + Call this method after the task completes, just before returning + the returncode from wait() or poll(). This hook is used to trigger exit listeners when the returncode first becomes available. """ - if self._exit_listeners is not None: + if self.returncode is not None and \ + self._exit_listeners is not None: for f in self._exit_listeners: f(self) self._exit_listeners = None @@ -1537,7 +1558,30 @@ class CompositeTask(AsynchronousTask): if self._current_task is not None: self._current_task.cancel() - def wait(self): + def _poll(self): + """ + This does a loop calling self._current_task.poll() + repeatedly as long as the value of self._current_task + keeps changing. It calls poll() a maximum of one time + for a given self._current_task instance. This is useful + since calling poll() on a task can trigger advance to + the next task could eventually lead to the returncode + being set in cases when polling only a single task would + not have the same effect. + """ + + prev = None + while True: + task = self._current_task + if task is None or task is prev: + # don't poll the same task more than once + break + task.poll() + prev = task + + return self.returncode + + def _wait(self): while True: task = self._current_task @@ -1547,13 +1591,64 @@ class CompositeTask(AsynchronousTask): self.scheduler.schedule(task.reg_id) task.wait() - self._wait_hook() return self.returncode + def _assert_current(self, task): + """ + Raises an AssertionError if the given task is not the + same one as self._current_task. This can be useful + for detecting bugs. + """ + if task is not self._current_task: + raise AssertionError("Unrecognized task: %s" % (task,)) + + def _default_exit(self, task): + """ + Calls _assert_current() on the given task and then sets the + composite returncode attribute if task.returncode != os.EX_OK. + If the task failed then self._current_task will be set to None. + Subclasses can use this as a generic task exit callback. + + @rtype: int + @returns: The task.returncode attribute. + """ + self._assert_current(task) + if task.returncode != os.EX_OK: + self.returncode = task.returncode + self._current_task = None + return task.returncode + + def _final_exit(self, task): + """ + Assumes that task is the final task of this composite task. + Calls _default_exit() and sets self.returncode to the task's + returncode and sets self._current_task to None. + + Subclasses can use this as a generic final task exit callback. + + """ + self._default_exit(task) + self._current_task = None + self.returncode = task.returncode + return self.returncode + + def _start_task(self, task, exit_handler): + """ + Register exit handler for the given task, set it + as self._current_task, and call task.start(). + + Subclasses can use this as a generic way to start + a task. + + """ + task.addExitListener(exit_handler) + self._current_task = task + task.start() + class TaskSequence(CompositeTask): """ A collection of tasks that executes sequentially. Each task - must have a _set_returncode() method that can be wrapped as + must have a addExitListener() method that can be used as a means to trigger movement from one task to the next. """ @@ -1574,27 +1669,27 @@ class TaskSequence(CompositeTask): 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() + self._start_task(self._task_queue.popleft(), + self._task_exit_handler) 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: + if self._default_exit(task) != os.EX_OK: + pass + elif self._task_queue: self._start_next_task() - return - - self._current_task = None - self.returncode = task.returncode + else: + self._final_exit(task) + self.wait() class SubProcess(AsynchronousTask): __slots__ = ("pid",) - def poll(self): + # 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 _poll(self): if self.returncode is not None: return self.returncode retval = os.waitpid(self.pid, os.WNOHANG) @@ -1615,11 +1710,10 @@ class SubProcess(AsynchronousTask): return self.pid is not None and \ self.returncode is None - def wait(self): + def _wait(self): 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): @@ -1658,43 +1752,48 @@ class SpawnProcess(SubProcess): if self.cancelled: return - # flush any pending output + if self.fd_pipes is None: + self.fd_pipes = {} fd_pipes = self.fd_pipes - if fd_pipes is None: - fd_pipes = { - 0 : sys.stdin.fileno(), - 1 : sys.stdout.fileno(), - 2 : sys.stderr.fileno(), - } + 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(): + if fd == sys.stdout.fileno(): + sys.stdout.flush() + if fd == sys.stderr.fileno(): + sys.stderr.flush() logfile = self.logfile self.files = self._files_dict() files = self.files + master_fd, slave_fd = os.pipe() + fcntl.fcntl(master_fd, fcntl.F_SETFL, + fcntl.fcntl(master_fd, fcntl.F_GETFL) | os.O_NONBLOCK) + if logfile is not None: + + fd_pipes_orig = fd_pipes.copy() + fd_pipes[0] = fd_pipes_orig[0] + fd_pipes[1] = slave_fd + fd_pipes[2] = slave_fd + files.out = open(logfile, "a") portage.util.apply_secpass_permissions(logfile, uid=portage.portage_uid, gid=portage.portage_gid, mode=0660) - else: - fd_pipes.setdefault(1, sys.stdout.fileno()) - for fd in fd_pipes.itervalues(): - if fd == sys.stdout.fileno(): - sys.stdout.flush() - if fd == sys.stderr.fileno(): - sys.stderr.flush() - files.out = os.fdopen(os.dup(fd_pipes[1]), 'w') - master_fd, slave_fd = os.pipe() + output_handler = self._output_handler - fcntl.fcntl(master_fd, fcntl.F_SETFL, - fcntl.fcntl(master_fd, fcntl.F_GETFL) | os.O_NONBLOCK) + else: - fd_pipes.setdefault(0, sys.stdin.fileno()) - fd_pipes_orig = fd_pipes.copy() - fd_pipes[0] = fd_pipes_orig[0] - fd_pipes[1] = slave_fd - fd_pipes[2] = slave_fd + # Create a dummy pipe so the scheduler can monitor + # the process from inside a poll() loop. + fd_pipes[self._dummy_pipe_fd] = slave_fd + output_handler = self._dummy_handler kwargs = {} for k in self._spawn_kwarg_names: @@ -1713,7 +1812,7 @@ 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) + PollConstants.POLLIN, output_handler) self.registered = True def _output_handler(self, fd, event): @@ -1734,6 +1833,27 @@ class SpawnProcess(SubProcess): self.registered = False 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.process, self._bufsize) + except EOFError: + pass + if buf: + pass + else: + fd = files.process.fileno() + for f in files.values(): + f.close() + self.registered = False + return self.registered + class EbuildFetcher(SpawnProcess): __slots__ = ("pkg",) @@ -1836,14 +1956,14 @@ class EbuildBuildDir(SlotObject): class AlreadyLocked(portage.exception.PortageException): pass -class EbuildBuild(EbuildBuildDir): +class EbuildBuild(CompositeTask): __slots__ = ("args_set", "find_blockers", - "ldpath_mtimes", "logger", "opts", - "pkg", "pkg_count", "scheduler", - "settings", "world_atom") + "ldpath_mtimes", "logger", "opts", "pkg", "pkg_count", + "settings", "world_atom") + \ + ("_build_dir", "_buildpkg", "_ebuild_path", "_tree") - def execute(self): + def start(self): args_set = self.args_set find_blockers = self.find_blockers @@ -1861,6 +1981,7 @@ class EbuildBuild(EbuildBuildDir): world_set = root_config.sets["world"] vartree = root_config.trees["vartree"] tree = "porttree" + self._tree = tree portdb = root_config.trees[tree].dbapi debug = settings.get("PORTAGE_DEBUG") == "1" features = self.settings.features @@ -1868,6 +1989,7 @@ class EbuildBuild(EbuildBuildDir): settings.backup_changes("EMERGE_FROM") settings.reset() ebuild_path = portdb.findname(self.pkg.cpv) + self._ebuild_path = ebuild_path #buildsyspkg: Check if we need to _force_ binary package creation issyspkg = "buildsyspkg" in features and \ @@ -1876,109 +1998,146 @@ class EbuildBuild(EbuildBuildDir): if opts.fetchonly: if opts.pretend: - fetcher = EbuildFetchPretend( fetch_all=opts.fetch_all_uri, pkg=pkg, settings=settings) - retval = fetcher.execute() + self.returncode = retval + self.wait() else: - fetcher = EbuildFetcher(pkg=pkg, scheduler=scheduler) - fetcher.start() - scheduler.schedule(fetcher.reg_id) - retval = fetcher.wait() + self._start_task(fetcher, self._fetchonly_exit) - 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 + return - try: - self.lock() - # Cleaning is triggered before the setup - # phase, in portage.doebuild(). - msg = " === (%s of %s) Cleaning (%s::%s)" % \ + self._build_dir = EbuildBuildDir(pkg=pkg, settings=settings) + self._build_dir.lock() + + # Cleaning is triggered before the setup + # phase, in portage.doebuild(). + msg = " === (%s of %s) Cleaning (%s::%s)" % \ + (pkg_count.curval, pkg_count.maxval, pkg.cpv, ebuild_path) + short_msg = "emerge: (%s of %s) %s Clean" % \ + (pkg_count.curval, pkg_count.maxval, pkg.cpv) + logger.log(msg, short_msg=short_msg) + + if opts.buildpkg or issyspkg: + + self._buildpkg = True + if issyspkg: + portage.writemsg_stdout(">>> This is a system package, " + \ + "let's pack a rescue tarball.\n", noiselevel=-1) + + msg = " === (%s of %s) Compiling/Packaging (%s::%s)" % \ (pkg_count.curval, pkg_count.maxval, pkg.cpv, ebuild_path) - short_msg = "emerge: (%s of %s) %s Clean" % \ + short_msg = "emerge: (%s of %s) %s Compile" % \ (pkg_count.curval, pkg_count.maxval, pkg.cpv) logger.log(msg, short_msg=short_msg) - if opts.buildpkg or issyspkg: - if issyspkg: - portage.writemsg(">>> This is a system package, " + \ - "let's pack a rescue tarball.\n", noiselevel=-1) - msg = " === (%s of %s) Compiling/Packaging (%s::%s)" % \ - (pkg_count.curval, pkg_count.maxval, pkg.cpv, ebuild_path) - short_msg = "emerge: (%s of %s) %s Compile" % \ - (pkg_count.curval, pkg_count.maxval, pkg.cpv) - logger.log(msg, short_msg=short_msg) - - build = EbuildExecuter(pkg=pkg, scheduler=scheduler, - settings=settings) - build.start() - retval = build.wait() - if retval != os.EX_OK: - return retval + else: + msg = " === (%s of %s) Compiling/Merging (%s::%s)" % \ + (pkg_count.curval, pkg_count.maxval, pkg.cpv, ebuild_path) + short_msg = "emerge: (%s of %s) %s Compile" % \ + (pkg_count.curval, pkg_count.maxval, pkg.cpv) + logger.log(msg, short_msg=short_msg) - build = EbuildBinpkg(pkg=pkg, - scheduler=scheduler, settings=settings) + build = EbuildExecuter(pkg=pkg, scheduler=scheduler, + settings=settings) + self._start_task(build, self._build_exit) - build.start() - scheduler.schedule(build.reg_id) - retval = build.wait() + def _fetchonly_exit(self, fetcher): + if self._final_exit(fetcher) != os.EX_OK: + pkg = self.pkg + eerror("!!! Fetch for %s failed, continuing..." % pkg.cpv, + phase="unpack", key=pkg.cpv) + self.wait() - if retval != os.EX_OK: - return retval + def _unlock_builddir(self): + portage.elog.elog_process(self.pkg.cpv, self.settings) + self._build_dir.unlock() - if not opts.buildpkgonly: - msg = " === (%s of %s) Merging (%s::%s)" % \ - (pkg_count.curval, pkg_count.maxval, - pkg.cpv, ebuild_path) - short_msg = "emerge: (%s of %s) %s Merge" % \ - (pkg_count.curval, pkg_count.maxval, pkg.cpv) - logger.log(msg, short_msg=short_msg) - - merge = EbuildMerge(find_blockers=find_blockers, - ldpath_mtimes=ldpath_mtimes, logger=logger, pkg=pkg, - pkg_count=pkg_count, pkg_path=ebuild_path, - settings=settings, tree=tree, world_atom=world_atom) - retval = merge.execute() - if retval != os.EX_OK: - return retval - elif "noclean" not in settings.features: - portage.doebuild(ebuild_path, "clean", root, - settings, debug=debug, mydbapi=portdb, - tree=tree) - else: - msg = " === (%s of %s) Compiling/Merging (%s::%s)" % \ - (pkg_count.curval, pkg_count.maxval, pkg.cpv, ebuild_path) - short_msg = "emerge: (%s of %s) %s Compile" % \ - (pkg_count.curval, pkg_count.curval, pkg.cpv) - logger.log(msg, short_msg=short_msg) - - build = EbuildExecuter(pkg=pkg, scheduler=scheduler, - settings=settings) - build.start() - retval = build.wait() - if retval != os.EX_OK: - return retval + def _build_exit(self, build): + if self._default_exit(build) != os.EX_OK: + self._unlock_builddir() + return - merge = EbuildMerge(find_blockers=self.find_blockers, - ldpath_mtimes=ldpath_mtimes, logger=logger, pkg=pkg, - pkg_count=pkg_count, pkg_path=ebuild_path, - settings=settings, tree=tree, world_atom=world_atom) - retval = merge.execute() + opts = self.opts + buildpkg = self._buildpkg - if retval != os.EX_OK: - return retval + if not buildpkg: + self._final_exit(build) + self.wait() + return + + packager = EbuildBinpkg(pkg=self.pkg, + scheduler=self.scheduler, settings=self.settings) + + self._start_task(packager, self._buildpkg_exit) + + def _buildpkg_exit(self, packager): + """ + Released build dir lock when there is a failure or + when in buildpkgonly mode. Otherwise, the lock will + be released when merge() is called. + """ + + if self._default_exit(packager) == os.EX_OK and \ + self.opts.buildpkgonly: + # Need to call "clean" phase for buildpkgonly mode + phase = "clean" + clean_phase = EbuildPhase(pkg=self.pkg, phase=phase, + scheduler=self.scheduler, settings=self.settings, + tree=self._tree) + self._start_task(clean_phase, self._clean_exit) + return + + if self._final_exit(packager) != os.EX_OK or \ + self.opts.buildpkgonly: + self._unlock_builddir() + self.wait() + + def _clean_exit(self, clean_phase): + if self._final_exit(clean_phase) != os.EX_OK or \ + self.opts.buildpkgonly: + self._unlock_builddir() + self.wait() + + def install(self): + """ + Install the package and then clean up and release locks. + Only call this after the build has completed successfully + and neither fetchonly nor buildpkgonly mode are enabled. + """ + + find_blockers = self.find_blockers + ldpath_mtimes = self.ldpath_mtimes + logger = self.logger + pkg = self.pkg + pkg_count = self.pkg_count + settings = self.settings + world_atom = self.world_atom + ebuild_path = self._ebuild_path + tree = self._tree + + merge = EbuildMerge(find_blockers=self.find_blockers, + ldpath_mtimes=ldpath_mtimes, logger=logger, pkg=pkg, + pkg_count=pkg_count, pkg_path=ebuild_path, + settings=settings, tree=tree, world_atom=world_atom) + + msg = " === (%s of %s) Merging (%s::%s)" % \ + (pkg_count.curval, pkg_count.maxval, + pkg.cpv, ebuild_path) + short_msg = "emerge: (%s of %s) %s Merge" % \ + (pkg_count.curval, pkg_count.maxval, pkg.cpv) + logger.log(msg, short_msg=short_msg) + + try: + rval = merge.execute() finally: - if self.locked: - portage.elog.elog_process(pkg.cpv, settings) - self.unlock() - return os.EX_OK + self._unlock_builddir() + + return rval class EbuildExecuter(CompositeTask): @@ -1995,15 +2154,11 @@ class EbuildExecuter(CompositeTask): 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() + self._start_task(clean_phase, self._clean_phase_exit) def _clean_phase_exit(self, clean_phase): - if clean_phase.returncode != os.EX_OK: - self.returncode = clean_phase.returncode - self._current_task = None + if self._default_exit(clean_phase) != os.EX_OK: return pkg = self.pkg @@ -2028,13 +2183,7 @@ class EbuildExecuter(CompositeTask): pkg=pkg, phase=phase, scheduler=scheduler, settings=settings, tree=tree)) - ebuild_phases.addExitListener(self._ebuild_phases_exit) - self._current_task = ebuild_phases - ebuild_phases.start() - - def _ebuild_phases_exit(self, ebuild_phases): - self.returncode = ebuild_phases.returncode - self._current_task = None + self._start_task(ebuild_phases, self._final_exit) class EbuildPhase(SubProcess): @@ -2046,11 +2195,6 @@ class EbuildPhase(SubProcess): _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 @@ -2202,13 +2346,12 @@ class EbuildPhase(SubProcess): for l in wrap(msg, 72): eerror(l, phase=self.phase, key=self.pkg.cpv) - returncode = self.returncode settings = self.settings portage._post_phase_userpriv_perms(settings) if self.phase == "install": portage._check_build_log(settings) - if returncode == os.EX_OK: - returncode = portage._post_src_install_checks(settings) + if self.returncode == os.EX_OK: + self.returncode = portage._post_src_install_checks(settings) class EbuildBinpkg(EbuildPhase): """ @@ -2307,31 +2450,31 @@ class PackageUninstall(Task): return e.status return os.EX_OK -class Binpkg(EbuildBuildDir): +class Binpkg(CompositeTask): __slots__ = ("find_blockers", "ldpath_mtimes", "logger", "opts", - "pkg", "pkg_count", "prefetcher", "scheduler", - "settings", "world_atom") + "pkg", "pkg_count", "prefetcher", "settings", "world_atom") + \ + ("_bintree", "_build_dir", "_ebuild_path", "_fetched_pkg", + "_image_dir", "_infloc", "_pkg_path", "_tree", "_verify") - def execute(self): + def start(self): - find_blockers = self.find_blockers - ldpath_mtimes = self.ldpath_mtimes - logger = self.logger - opts = self.opts pkg = self.pkg - pkg_count = self.pkg_count - scheduler = self.scheduler 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 + self._tree = "bintree" + self._bintree = self.pkg.root_config.trees[self._tree] + self._verify = "strict" in self.settings.features and \ + not self.opts.pretend + + dir_path = os.path.join(settings["PORTAGE_TMPDIR"], + "portage", pkg.category, pkg.pf) + self._build_dir = EbuildBuildDir(dir_path=dir_path, + pkg=pkg, settings=settings) + self._image_dir = os.path.join(dir_path, "image") + self._infloc = os.path.join(dir_path, "build-info") + self._ebuild_path = os.path.join(self._infloc, pkg.pf + ".ebuild") # The prefetcher has already completed or it # could be running now. If it's running now, @@ -2343,55 +2486,83 @@ class Binpkg(EbuildBuildDir): # use the scheduler and fetcher methods to # synchronize with the fetcher. prefetcher = self.prefetcher - if prefetcher is not None: - if not prefetcher.isAlive(): - prefetcher.cancel() - else: - retval = prefetcher.poll() - - if retval is None: - waiting_msg = ("Fetching '%s' " + \ - "in the background. " + \ - "To view fetch progress, run `tail -f " + \ - "/var/log/emerge-fetch.log` in another " + \ - "terminal.") % prefetcher.pkg_path - msg_prefix = colorize("GOOD", " * ") - from textwrap import wrap - waiting_msg = "".join("%s%s\n" % (msg_prefix, line) \ - for line in wrap(waiting_msg, 65)) - writemsg(waiting_msg, noiselevel=-1) - - scheduler.schedule(prefetcher.reg_id) - retval = prefetcher.wait() - del prefetcher - - fetcher = BinpkgFetcher(pkg=pkg, scheduler=scheduler) + if prefetcher is None: + pass + elif not prefetcher.isAlive(): + prefetcher.cancel() + elif prefetcher.poll() is None: + + waiting_msg = ("Fetching '%s' " + \ + "in the background. " + \ + "To view fetch progress, run `tail -f " + \ + "/var/log/emerge-fetch.log` in another " + \ + "terminal.") % prefetcher.pkg_path + msg_prefix = colorize("GOOD", " * ") + from textwrap import wrap + waiting_msg = "".join("%s%s\n" % (msg_prefix, line) \ + for line in wrap(waiting_msg, 65)) + writemsg(waiting_msg, noiselevel=-1) + + self._current_task = prefetcher + prefetcher.addExitListener(self._prefetch_exit) + return + + self._prefetch_exit(prefetcher) + + def _prefetch_exit(self, prefetcher): + + pkg = self.pkg + pkg_count = self.pkg_count + fetcher = BinpkgFetcher(pkg=self.pkg, scheduler=self.scheduler) pkg_path = fetcher.pkg_path + self._pkg_path = pkg_path - if opts.getbinpkg and bintree.isremote(pkg.cpv): + if self.opts.getbinpkg and self._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) + self.logger.log(msg, short_msg=short_msg) - fetcher.start() - scheduler.schedule(fetcher.reg_id) - retval = fetcher.wait() + self._start_task(fetcher, self._fetcher_exit) + return - if retval != os.EX_OK: - return retval + self._fetcher_exit(fetcher) - if opts.fetchonly: - return os.EX_OK + def _fetcher_exit(self, fetcher): - if verify: - verifier = BinpkgVerifier(pkg=pkg) - verifier.start() - retval = verifier.wait() - if retval != os.EX_OK: - return retval + # The fetcher only has a returncode when + # --getbinpkg is enabled. + if fetcher.returncode is not None: + self._fetched_pkg = True + if self.opts.fetchonly: + self._final_exit(fetcher) + self.wait() + return + elif self._default_exit(fetcher) != os.EX_OK: + return + + verifier = None + if self._verify: + verifier = BinpkgVerifier(pkg=self.pkg) + self._start_task(verifier, self._verifier_exit) + return + + self._verifier_exit(verifier) + + def _verifier_exit(self, verifier): + if verifier is not None and \ + self._default_exit(verifier) != os.EX_OK: + return + + logger = self.logger + pkg = self.pkg + pkg_count = self.pkg_count + pkg_path = self._pkg_path + + if self._fetched_pkg: + self._bintree.inject(pkg.cpv, filename=pkg_path) msg = " === (%s of %s) Merging Binary (%s::%s)" % \ (pkg_count.curval, pkg_count.maxval, pkg.cpv, pkg_path) @@ -2399,126 +2570,131 @@ class Binpkg(EbuildBuildDir): (pkg_count.curval, pkg_count.maxval, pkg.cpv) logger.log(msg, short_msg=short_msg) - dir_path = os.path.join(settings["PORTAGE_TMPDIR"], - "portage", pkg.category, pkg.pf) - image_dir = os.path.join(dir_path, "image") - infloc = os.path.join(dir_path, "build-info") + self._build_dir.lock() - fd_pipes = { - 0 : sys.stdin.fileno(), - 1 : sys.stdout.fileno(), - 2 : sys.stderr.fileno(), - } + phase = "clean" + settings = self.settings + settings.setcpv(pkg) + settings["EBUILD"] = self._ebuild_path + ebuild_phase = EbuildPhase( + pkg=pkg, phase=phase, scheduler=self.scheduler, + settings=settings, tree=self._tree) - try: - self.lock() + self._start_task(ebuild_phase, self._clean_exit) - root_config = self.pkg.root_config - ebuild_path = os.path.join(infloc, pkg.pf + ".ebuild") - cleanup = 1 - mydbapi = root_config.trees[tree].dbapi + def _clean_exit(self, clean_phase): + if self._default_exit(clean_phase) != os.EX_OK: + self._unlock_builddir() + return - phase = "clean" - ebuild_phase = EbuildPhase(fd_pipes=fd_pipes, - pkg=pkg, phase=phase, scheduler=scheduler, - settings=settings, tree=tree) + dir_path = self._build_dir.dir_path - ebuild_phase.start() - scheduler.schedule(ebuild_phase.reg_id) - retval = ebuild_phase.wait() + try: + shutil.rmtree(dir_path) + except (IOError, OSError), e: + if e.errno != errno.ENOENT: + raise + del e - if retval != os.EX_OK: - return retval + infloc = self._infloc + pkg = self.pkg + pkg_path = self._pkg_path - try: - shutil.rmtree(dir_path) - except (IOError, OSError), e: - if e.errno != errno.ENOENT: - raise - del e + dir_mode = 0755 + for mydir in (dir_path, self._image_dir, infloc): + portage.util.ensure_dirs(mydir, uid=portage.data.portage_uid, + gid=portage.data.portage_gid, mode=dir_mode) - # This initializes PORTAGE_LOG_FILE. - portage.prepare_build_dirs(root_config.root, settings, cleanup) - - dir_mode = 0755 - for mydir in (dir_path, image_dir, infloc): - portage.util.ensure_dirs(mydir, uid=portage.data.portage_uid, - gid=portage.data.portage_gid, mode=dir_mode) - - portage.writemsg_stdout(">>> Extracting info\n") - - pkg_xpak = portage.xpak.tbz2(pkg_path) - check_missing_metadata = ("CATEGORY", "PF") - missing_metadata = set() - for k in check_missing_metadata: - v = pkg_xpak.getfile(k) - if not v: - missing_metadata.add(k) - - pkg_xpak.unpackinfo(infloc) - for k in missing_metadata: - if k == "CATEGORY": - v = pkg.category - elif k == "PF": - v = pkg.pf - else: - continue + portage.writemsg_stdout(">>> Extracting info\n") - f = open(os.path.join(infloc, k), 'wb') - try: - f.write(v + "\n") - finally: - f.close() + # This initializes PORTAGE_LOG_FILE. + portage.prepare_build_dirs(self.settings["ROOT"], self.settings, 1) + + pkg_xpak = portage.xpak.tbz2(self._pkg_path) + check_missing_metadata = ("CATEGORY", "PF") + missing_metadata = set() + for k in check_missing_metadata: + v = pkg_xpak.getfile(k) + if not v: + missing_metadata.add(k) + + pkg_xpak.unpackinfo(infloc) + for k in missing_metadata: + if k == "CATEGORY": + v = pkg.category + elif k == "PF": + v = pkg.pf + else: + continue - # Store the md5sum in the vdb. - f = open(os.path.join(infloc, "BINPKGMD5"), "w") + f = open(os.path.join(infloc, k), 'wb') try: - f.write(str(portage.checksum.perform_md5(pkg_path)) + "\n") + f.write(v + "\n") finally: f.close() - # This gives bashrc users an opportunity to do various things - # such as remove binary packages after they're installed. - settings["PORTAGE_BINPKG_FILE"] = pkg_path - settings.backup_changes("PORTAGE_BINPKG_FILE") + # Store the md5sum in the vdb. + f = open(os.path.join(infloc, "BINPKGMD5"), "w") + try: + f.write(str(portage.checksum.perform_md5(pkg_path)) + "\n") + finally: + f.close() - phase = "setup" - ebuild_phase = EbuildPhase(fd_pipes=fd_pipes, - pkg=pkg, phase=phase, scheduler=scheduler, - settings=settings, tree=tree) + # This gives bashrc users an opportunity to do various things + # such as remove binary packages after they're installed. + settings = self.settings + settings.setcpv(self.pkg) + settings["PORTAGE_BINPKG_FILE"] = pkg_path + settings.backup_changes("PORTAGE_BINPKG_FILE") - ebuild_phase.start() - scheduler.schedule(ebuild_phase.reg_id) - retval = ebuild_phase.wait() + phase = "setup" + ebuild_phase = EbuildPhase( + pkg=self.pkg, phase=phase, scheduler=self.scheduler, + settings=settings, tree=self._tree) - if retval != os.EX_OK: - return retval + self._start_task(ebuild_phase, self._setup_exit) - extractor = BinpkgExtractorAsync(image_dir=image_dir, - pkg=pkg, pkg_path=pkg_path, scheduler=scheduler) - portage.writemsg_stdout(">>> Extracting %s\n" % pkg.cpv) - extractor.start() - scheduler.schedule(extractor.reg_id) - retval = extractor.wait() + def _setup_exit(self, setup_phase): + if self._default_exit(setup_phase) != os.EX_OK: + self._unlock_builddir() + return - if retval != os.EX_OK: - writemsg("!!! Error Extracting '%s'\n" % pkg_path, - noiselevel=-1) - return retval + extractor = BinpkgExtractorAsync(image_dir=self._image_dir, + pkg=self.pkg, pkg_path=self._pkg_path, scheduler=self.scheduler) + portage.writemsg_stdout(">>> Extracting %s\n" % self.pkg.cpv) + self._start_task(extractor, self._extractor_exit) - merge = EbuildMerge(find_blockers=find_blockers, - ldpath_mtimes=ldpath_mtimes, logger=logger, pkg=pkg, - pkg_count=pkg_count, pkg_path=pkg_path, - settings=settings, tree=tree, world_atom=world_atom) + def _extractor_exit(self, extractor): + if self._final_exit(extractor) != os.EX_OK: + self._unlock_builddir() + writemsg("!!! Error Extracting '%s'\n" % self._pkg_path, + noiselevel=-1) + self.wait() - retval = merge.execute() - if retval != os.EX_OK: - return retval + def _unlock_builddir(self): + portage.elog.elog_process(self.pkg.cpv, self.settings) + self._build_dir.unlock() + + def install(self): + + # This gives bashrc users an opportunity to do various things + # such as remove binary packages after they're installed. + settings = self.settings + settings["PORTAGE_BINPKG_FILE"] = self._pkg_path + settings.backup_changes("PORTAGE_BINPKG_FILE") + merge = EbuildMerge(find_blockers=self.find_blockers, + ldpath_mtimes=self.ldpath_mtimes, logger=self.logger, + pkg=self.pkg, pkg_count=self.pkg_count, + pkg_path=self._pkg_path, settings=settings, + tree=self._tree, world_atom=self.world_atom) + + try: + retval = merge.execute() finally: settings.pop("PORTAGE_BINPKG_FILE", None) - self.unlock() - return os.EX_OK + self._unlock_builddir() + return retval class BinpkgFetcher(SpawnProcess): @@ -2526,7 +2702,7 @@ class BinpkgFetcher(SpawnProcess): "locked", "pkg_path", "_lock_obj") def __init__(self, **kwargs): - SubProcess.__init__(self, **kwargs) + SpawnProcess.__init__(self, **kwargs) pkg = self.pkg self.pkg_path = pkg.root_config.trees["bintree"].getname(pkg.cpv) @@ -2576,6 +2752,17 @@ class BinpkgFetcher(SpawnProcess): if use_locks: self.lock() + if self.fd_pipes is None: + self.fd_pipes = {} + fd_pipes = self.fd_pipes + + # Redirect all output to stdout since some fetchers like + # wget pollute stderr (if portage detects a problem then it + # can send it's own message to stderr). + fd_pipes.setdefault(0, sys.stdin.fileno()) + fd_pipes.setdefault(1, sys.stdout.fileno()) + fd_pipes.setdefault(2, sys.stdout.fileno()) + self.args = fetch_args self.env = fetch_env SpawnProcess.start(self) @@ -2643,12 +2830,7 @@ class BinpkgVerifier(AsynchronousTask): rval = 1 self.returncode = rval - - def cancel(self): - self.cancelled = True - - def poll(self): - return self.returncode + self.wait() class BinpkgExtractorAsync(SpawnProcess): @@ -2665,7 +2847,7 @@ class BinpkgExtractorAsync(SpawnProcess): self.env = self.pkg.root_config.settings.environ() SpawnProcess.start(self) -class MergeListItem(SlotObject): +class MergeListItem(CompositeTask): """ TODO: For parallel scheduling, everything here needs asynchronous @@ -2674,39 +2856,30 @@ class MergeListItem(SlotObject): __slots__ = ("args_set", "binpkg_opts", "build_opts", "emerge_opts", "failed_fetches", "find_blockers", "logger", "mtimedb", "pkg", - "pkg_count", "prefetcher", "scheduler", "settings", "world_atom") + "pkg_count", "prefetcher", "settings", "world_atom") + \ + ("_install_task",) - def execute(self): + def start(self): - args_set = self.args_set - binpkg_opts = self.binpkg_opts + pkg = self.pkg build_opts = self.build_opts - emerge_opts = self.emerge_opts - failed_fetches = self.failed_fetches + + if pkg.installed: + # uninstall, executed by self.merge() + self.returncode = os.EX_OK + self.wait() + return + + args_set = self.args_set 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" % \ @@ -2725,27 +2898,83 @@ class MergeListItem(SlotObject): 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 + self._install_task = build + self._start_task(build, self._ebuild_exit) + self.wait() + return 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, + opts=self.binpkg_opts, pkg=pkg, pkg_count=pkg_count, + prefetcher=self.prefetcher, settings=settings, scheduler=scheduler, world_atom=world_atom) - retval = binpkg.execute() + self._install_task = binpkg + self._start_task(binpkg, self._final_exit) + self.wait() + return - if retval != os.EX_OK: - return retval + def _ebuild_exit(self, build): + if self._final_exit(build) != os.EX_OK: + if self.build_opts.fetchonly: + self.failed_fetches.append(self.pkg.cpv) + self.wait() - return os.EX_OK + def _poll(self): + self._install_task.poll() + return self.returncode + + def _wait(self): + self._install_task.wait() + return self.returncode + + def merge(self): + + pkg = self.pkg + build_opts = self.build_opts + failed_fetches = self.failed_fetches + find_blockers = self.find_blockers + logger = self.logger + mtimedb = self.mtimedb + 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=self.emerge_opts, pkg=pkg, settings=settings) + + retval = uninstall.execute() + if retval != os.EX_OK: + return retval + return os.EX_OK + + if build_opts.fetchonly or \ + build_opts.buildpkgonly: + return self.returncode + + retval = self._install_task.install() + return retval + +class PackageMerge(CompositeTask): + """ + TODO: Implement asynchronous merge so that the scheduler can + run while a merge is executing. + """ + + __slots__ = ("merge",) + + def start(self): + self.returncode = self.merge.merge() + self.wait() class DependencyArg(object): def __init__(self, arg=None, root_config=None): @@ -7286,13 +7515,19 @@ class SequentialTaskQueue(SlotObject): self._task_queue.append(task) def schedule(self): + + if not self: + return False + 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: + if hasattr(task, "registered") and task.registered: + continue + if task.poll() is not None: running_tasks.remove(task) state_changed = True @@ -7313,6 +7548,12 @@ class SequentialTaskQueue(SlotObject): task = running_tasks.pop() task.cancel() + def __nonzero__(self): + return bool(self._task_queue or self.running_tasks) + + def __len__(self): + return len(self._task_queue) + len(self.running_tasks) + class Scheduler(object): _opts_ignore_blockers = \ @@ -7328,6 +7569,9 @@ class Scheduler(object): class _iface_class(SlotObject): __slots__ = ("register", "schedule") + _task_queues_class = slot_dict_class( + ("build", "extract", "merge", "prefetch",), prefix="") + class _build_opts_class(SlotObject): __slots__ = ("buildpkg", "buildpkgonly", "fetch_all_uri", "fetchonly", "pretend") @@ -7387,12 +7631,11 @@ class Scheduler(object): except AttributeError: self._poll = PollSelectAdapter() - self._task_queues = slot_dict_class(("build", "prefetch"), prefix="") + self._task_queues = self._task_queues_class() 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 = [] @@ -7402,6 +7645,8 @@ class Scheduler(object): if isinstance(x, Package) and x.operation == "merge"]) self._pkg_count = self._pkg_count_class( curval=0, maxval=merge_count) + self._max_jobs = 1 + self._jobs = 0 features = self.settings.features if "parallel-fetch" in features and \ @@ -7701,35 +7946,40 @@ class Scheduler(object): elif isinstance(pkg, Blocker): pass - def _choose_pkg(self): - return self._pkg_queue.popleft() - - def _main_loop(self): - - pkg_queue = self._pkg_queue + def _merge_exit(self, merge): + self._jobs -= 1 + pkg = merge.merge.pkg + if merge.returncode != os.EX_OK: + self._failed_pkgs.append((pkg, retval)) + return - while pkg_queue: - pkg = self._choose_pkg() - retval = self._execute_pkg(pkg) + if pkg.installed: + return - if retval != os.EX_OK: - self._failed_pkgs.append((pkg, retval)) - if not self._build_opts.fetchonly: - return + self._restart_if_necessary(pkg) - if pkg.installed: - continue + # 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() - self._restart_if_necessary(pkg) + def _build_exit(self, build): + if build.returncode == os.EX_OK: + self.curval += 1 + merge = PackageMerge(merge=build) + merge.addExitListener(self._merge_exit) + self._task_queues.merge.add(merge) + self._task_queues.merge.schedule() + else: + self._failed_pkgs.append((build.pkg, build.returncode)) + self._jobs -= 1 - # 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 _extract_exit(self, build): + self._build_exit(build) def _merge(self): @@ -7757,12 +8007,68 @@ class Scheduler(object): return rval - def _execute_pkg(self, pkg): + def _choose_pkg(self): + return self._pkg_queue.popleft() + + def _main_loop(self): + + pkg_queue = self._pkg_queue + failed_pkgs = self._failed_pkgs + task_queues = self._task_queues - if not pkg.installed: - self._pkg_count.curval += 1 + while pkg_queue and not failed_pkgs: - merge = MergeListItem(args_set=self._args_set, + pkg = self._choose_pkg() + + if not pkg.installed: + self._pkg_count.curval += 1 + + task = self._task(pkg) + + self._jobs += 1 + if pkg.installed: + merge = PackageMerge(merge=task) + merge.addExitListener(self._merge_exit) + task_queues.merge.add(merge) + elif pkg.built: + task.addExitListener(self._extract_exit) + task_queues.extract.add(task) + else: + task.addExitListener(self._build_exit) + task_queues.build.add(task) + + self._schedule_main() + + while self._jobs: + self._schedule_main(wait=True) + + def _schedule_main(self, wait=False): + + event_handlers = self._poll_event_handlers + poll = self._poll.poll + max_jobs = self._max_jobs + + self._schedule_tasks() + + while event_handlers: + jobs = self._jobs + + for f, event in poll(): + handler, reg_id = event_handlers[f] + if not handler(f, event): + self._unregister(reg_id) + + if jobs == self._jobs: + continue + + self._schedule_tasks() + + if not wait and self._jobs < max_jobs: + break + + def _task(self, pkg): + + task = MergeListItem(args_set=self._args_set, binpkg_opts=self._binpkg_opts, build_opts=self._build_opts, emerge_opts=self.myopts, @@ -7774,12 +8080,7 @@ class Scheduler(object): settings=self.pkgsettings[pkg.root], world_atom=self._world_atom) - retval = merge.execute() - - if retval == os.EX_OK: - self.curval += 1 - - return retval + return task def _save_resume_list(self): """ @@ -7869,6 +8170,10 @@ class Scheduler(object): del self._poll_event_handler_ids[reg_id] self._schedule_tasks() + def _schedule_tasks(self): + for x in self._task_queues.values(): + x.schedule() + def _schedule(self, wait_id): """ Schedule until wait_id is not longer registered -- 2.26.2