From fc66054d6695bf119c537b051c5c2a78306746d4 Mon Sep 17 00:00:00 2001 From: Fabian Groffen Date: Tue, 1 Jul 2008 17:22:29 +0000 Subject: [PATCH] Merged from trunk 10853:10869 | 10854 | Avoid python-2.6 deprecation warnings for md5 and sha | | zmedico | modules by trying to import hashlib first and then falling | | | back to the deprecated modules if necessary. Thanks to | | | ColdWind for reporting. | | 10855 | Reimplement parallel-fetch by spawning the `ebuild fetch` | | zmedico | command for each ebuild. The benefit of using this approach | | | is that it can be integrated together with parallel build | | | scheduling that's planned. Parallel-fetch support for | | | binhost is not implemented yet, though it worked previously. | | 10856 | Clear the self._task_queue to avoid duplicate parallel-fetch | | zmedico | tasks in --keep-going mode. | | 10857 | Add "(no inline comments)" to qualify "comments begin with | | zmedico | #" statements. | | 10858 | Bug #230245 - Pass the correct directory when calling `snv | | zmedico | list` and `svn status` since repoman supports category-level | | | and repo-level commits. | | 10859 | Bug #230245 - Use os.path.basename() on paths returned from | | zmedico | `svn list` and `svn status`. | | 10860 | Bug #230249 - Disable the "ebuild.notadded" check when not | | zmedico | in commit mode and running `svn list` and `svn status` calls | | | in every package dir will be too expensive. | | 10861 | Fix typo. | | zmedico | | | 10862 | add a call to pruneNonExisting() at the end of | | zmedico | dbapi.vartree.PreservedLibsRegistry.__init__() | | 10864 | Split out a write_contents() function and a | | zmedico | vardbapi.removeFromContents() function. This is refactoring | | | of code from the blocker file collision contents handling in | | | dblink.treewalk(). Also, there is a new | | | dblink._match_contents() method derived from isowner(). It | | | returns the exact path from the contents file that matches | | | the given path, regardless of path differences due to things | | | such as symlinks. | | 10865 | Handle potential errors in PreservedLibsRegistry.store() now | | zmedico | that it can be called via pruneNonExisting(), due to things | | | such as portageq calls where the user may not have write | | | permission to the registry. | | 10866 | Also avoid sandbox violations in | | zmedico | PreservedLibsRegistry.store(), for running portage inside | | | ebuild phases. | | 10867 | Never do realpath() on an empty string for | | zmedico | portdbapi.porttree_root since otherwise it can evaluate to | | | $CWD which leads to undesireable results. | | 10868 | Add a new BinpkgFetcherAsync class and use it to implement | | zmedico | parellel-fetch for --getbinpkg. | | 10869 | Add a "prefix" keyword parameter to slot_dict_class() which | | zmedico | controls the prefix used when mapping attribute names from | | | keys. Use this to change the syntax from files["foo"] to | | | files.foo (it's fewer characters to look at). | svn path=/main/branches/prefix/; revision=10880 --- bin/repoman | 18 +- man/portage.5 | 28 +- pym/_emerge/__init__.py | 578 +++++++++++++++++++++++++++------- pym/portage/__init__.py | 5 +- pym/portage/cache/mappings.py | 24 +- pym/portage/checksum.py | 16 +- pym/portage/dbapi/porttree.py | 4 +- pym/portage/dbapi/vartree.py | 113 +++++-- 8 files changed, 610 insertions(+), 176 deletions(-) diff --git a/bin/repoman b/bin/repoman index 0003acda4..58f048725 100755 --- a/bin/repoman +++ b/bin/repoman @@ -769,6 +769,12 @@ arch_caches={} arch_xmatch_caches = {} shared_xmatch_caches = {"cp-list":{}} +# Disable the "ebuild.notadded" check when not in commit mode and +# running `svn list` and `svn status` calls in every package dir +# will be too expensive. +check_ebuild_notadded = not \ + (vcs == "svn" and repolevel < 3 and options.mode != "commit") + for x in scanlist: #ebuilds and digests added to cvs respectively. logging.info("checking package %s" % x) @@ -865,12 +871,12 @@ for x in scanlist: if not os.path.isdir(os.path.join(checkdir, "files")): has_filesdir = False - if vcs: + if vcs and check_ebuild_notadded: try: if vcs == "cvs": myf=open(checkdir+"/CVS/Entries","r") if vcs == "svn": - myf=os.popen("svn list") + myf = os.popen("svn list " + checkdir) myl=myf.readlines() myf.close() for l in myl: @@ -887,16 +893,16 @@ for x in scanlist: if l[-1:] == "/": continue if l[-7:] == ".ebuild": - eadded.append(l[:-7]) + eadded.append(os.path.basename(l[:-7])) if vcs == "svn": - myf=os.popen("svn status") + myf = os.popen("svn status " + checkdir) myl=myf.readlines() myf.close() for l in myl: if l[0] == "A": l = l.rstrip().split(' ')[-1] if l[-7:] == ".ebuild": - eadded.append(l[:-7]) + eadded.append(os.path.basename(l[:-7])) except IOError: if options.mode == 'commit' and vcs == "cvs": stats["CVS/Entries.IO_error"] += 1 @@ -1068,7 +1074,7 @@ for x in scanlist: if stat.S_IMODE(os.stat(full_path).st_mode) & 0111: stats["file.executable"] += 1 fails["file.executable"].append(x+"/"+y+".ebuild") - if vcs and y not in eadded: + if vcs and check_ebuild_notadded and y not in eadded: #ebuild not added to vcs stats["ebuild.notadded"]=stats["ebuild.notadded"]+1 fails["ebuild.notadded"].append(x+"/"+y+".ebuild") diff --git a/man/portage.5 b/man/portage.5 index d6eb83cb6..d6a29678a 100644 --- a/man/portage.5 +++ b/man/portage.5 @@ -187,7 +187,7 @@ Provides the list of packages that compose the special \fIsystem\fR set. .I Format: .nf -\- comments begin with # +\- comments begin with # (no inline comments) \- one DEPEND atom per line \- packages to be added to the system set begin with a * .fi @@ -234,7 +234,7 @@ package.provided. .I Format: .nf -\- comments begin with # +\- comments begin with # (no inline comments) \- one DEPEND atom per line \- relational operators are not allowed \- must include a version @@ -262,7 +262,7 @@ a '\-'. .I Format: .nf -\- comments begin with # +\- comments begin with # (no inline comments) \- one DEPEND atom per line with space-delimited USE flags .fi @@ -284,7 +284,7 @@ a '\-'. .I Format: .nf -\- comments begin with # +\- comments begin with # (no inline comments) \- one DEPEND atom per line with space-delimited USE flags .fi @@ -318,7 +318,7 @@ a '\-'. .I Format: .nf -\- comments begin with # +\- comments begin with # (no inline comments) \- one USE flag per line .fi .TP @@ -334,7 +334,7 @@ a '\-'. .I Format: .nf -\- comments begin with # +\- comments begin with # (no inline comments) \- one USE flag per line .fi .TP @@ -348,7 +348,7 @@ the package that does the very bare minimum to send e\-mail. .I Format: .nf -\- comments begin with # +\- comments begin with # (no inline comments) \- one virtual and DEPEND atom base pair per line .fi @@ -478,7 +478,7 @@ documentation for QT. Easy as pie my friend! .I Format: .nf -\- comments begin with # +\- comments begin with # (no inline comments) \- one DEPEND atom per line with space-delimited USE flags .fi @@ -500,7 +500,7 @@ RESTRICT="mirror" or RESTRICT="fetch". .I Format: .nf -\- comments begin with # +\- comments begin with # (no inline comments) \- mirror type followed by a list of hosts .fi @@ -585,7 +585,7 @@ package has been masked and WHO is doing the masking. .I Format: .nf -\- comments begin with # +\- comments begin with # (no inline comments) \- one DEPEND atom per line .fi @@ -606,7 +606,7 @@ allowed per stable/dev/KEYWORD; the last one found is the last one used. .I Format: .nf -\- comments begin with # +\- comments begin with # (no inline comments) \- one profile list per line in format: arch dir status \- arch must be listed in arch.list \- dir is relative to profiles.desc @@ -631,7 +631,7 @@ mirrors. Keeps us from overloading a single server. .I Format: .nf -\- comments begin with # +\- comments begin with # (no inline comments) \- mirror type followed by a list of hosts .fi @@ -649,7 +649,7 @@ All global USE flags must be listed here with a description of what they do. .I Format: .nf -\- comments begin with # +\- comments begin with # (no inline comments) \- use flag \- some description .fi @@ -666,7 +666,7 @@ description. .nf .I Format: -\- comments begin with # +\- comments begin with # (no inline comments) \- package:use flag \- description .I Example: diff --git a/pym/_emerge/__init__.py b/pym/_emerge/__init__.py index 90162d7e9..65c01a6d8 100644 --- a/pym/_emerge/__init__.py +++ b/pym/_emerge/__init__.py @@ -21,7 +21,11 @@ except KeyboardInterrupt: sys.exit(1) import array +import fcntl import select +import shlex +import urlparse +import weakref import gc import os, stat import platform @@ -1460,26 +1464,141 @@ class _PackageMetadataWrapper(_PackageMetadataWrapperBase): v = 0 self._pkg.mtime = v -class EbuildFetcher(Task): +class EbuildFetcher(SlotObject): __slots__ = ("fetch_all", "pkg", "pretend", "settings") - def _get_hash_key(self): - hash_key = getattr(self, "_hash_key", None) - if hash_key is None: - self._hash_key = ("EbuildFetcher", self.pkg._get_hash_key()) - return self._hash_key - def execute(self): portdb = self.pkg.root_config.trees["porttree"].dbapi ebuild_path = portdb.findname(self.pkg.cpv) debug = self.settings.get("PORTAGE_DEBUG") == "1" + retval = portage.doebuild(ebuild_path, "fetch", - self.settings["ROOT"], self.settings, debug, - self.pretend, fetchonly=1, fetchall=self.fetch_all, + self.settings["ROOT"], self.settings, debug=debug, + listonly=self.pretend, fetchonly=1, fetchall=self.fetch_all, mydbapi=portdb, tree="porttree") return retval +class EbuildFetcherAsync(SlotObject): + + __slots__ = ("log_file", "fd_pipes", "pkg", + "register", "unregister", + "pid", "returncode", "files") + + _file_names = ("fetcher", "out") + _files_dict = slot_dict_class(_file_names, prefix="") + _bufsize = 4096 + + def start(self): + # flush any pending output + fd_pipes = self.fd_pipes + if fd_pipes is None: + fd_pipes = { + 0 : sys.stdin.fileno(), + 1 : sys.stdout.fileno(), + 2 : sys.stderr.fileno(), + } + + log_file = self.log_file + self.files = self._files_dict() + files = self.files + + if log_file is not None: + files.out = open(log_file, "a") + portage.util.apply_secpass_permissions(log_file, + uid=portage.portage_uid, gid=portage.portage_gid, + mode=0660) + else: + 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() + + import fcntl + fcntl.fcntl(master_fd, fcntl.F_SETFL, + fcntl.fcntl(master_fd, fcntl.F_GETFL) | os.O_NONBLOCK) + + 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 + + root_config = self.pkg.root_config + portdb = root_config.trees["porttree"].dbapi + ebuild_path = portdb.findname(self.pkg.cpv) + settings = root_config.settings + + fetch_env = dict((k, settings[k]) for k in settings) + fetch_env["FEATURES"] = fetch_env.get("FEATURES", "") + " -cvs" + fetch_env["PORTAGE_NICENESS"] = "0" + fetch_env["PORTAGE_PARALLEL_FETCHONLY"] = "1" + + ebuild_binary = os.path.join( + settings["EBUILD_BIN_PATH"], "ebuild") + + fetch_args = [ebuild_binary, ebuild_path, "fetch"] + debug = settings.get("PORTAGE_DEBUG") == "1" + if debug: + fetch_args.append("--debug") + + retval = portage.process.spawn(fetch_args, env=fetch_env, + fd_pipes=fd_pipes, returnpid=True) + + self.pid = retval[0] + + os.close(slave_fd) + files.fetcher = os.fdopen(master_fd, 'r') + self.register(files.fetcher.fileno(), + select.POLLIN, self._output_handler) + + def _output_handler(self, fd, event): + files = self.files + buf = array.array('B') + try: + buf.fromfile(files.fetcher, self._bufsize) + except EOFError: + pass + if buf: + buf.tofile(files.out) + files.out.flush() + else: + self.unregister(files.fetcher.fileno()) + for f in files.values(): + f.close() + + def poll(self): + if self.returncode is not None: + return self.returncode + retval = os.waitpid(self.pid, os.WNOHANG) + if retval == (0, 0): + return None + self._set_returncode(retval) + return self.returncode + + def wait(self): + if self.returncode is not None: + return self.returncode + self._set_returncode(os.waitpid(self.pid, 0)) + return self.returncode + + def _set_returncode(self, wait_retval): + + retval = wait_retval[1] + portage.process.spawned_pids.remove(self.pid) + if retval != os.EX_OK: + if retval & 0xff: + retval = (retval & 0xff) << 8 + else: + retval = retval >> 8 + + self.returncode = retval + class EbuildBuildDir(SlotObject): __slots__ = ("pkg", "settings", @@ -1593,9 +1712,12 @@ class EbuildBuild(Task): ebuild_phase = EbuildPhase(fd_pipes=fd_pipes, pkg=self.pkg, phase=mydo, register=self.register, settings=settings, unregister=self.unregister) + ebuild_phase.start() - self.schedule() - retval = ebuild_phase.wait() + retval = None + while retval is None: + self.schedule() + retval = ebuild_phase.poll() portage._post_phase_userpriv_perms(settings) if mydo == "install": @@ -1615,7 +1737,7 @@ class EbuildPhase(SlotObject): "pid", "returncode", "files") _file_names = ("log", "stdout", "ebuild") - _files_dict = slot_dict_class(_file_names) + _files_dict = slot_dict_class(_file_names, prefix="") _bufsize = 4096 def start(self): @@ -1690,33 +1812,48 @@ class EbuildPhase(SlotObject): 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.register(files["ebuild"].fileno(), + files.log = open(logfile, 'a') + files.stdout = os.fdopen(os.dup(fd_pipes_orig[1]), 'w') + files.ebuild = os.fdopen(master_fd, 'r') + self.register(files.ebuild.fileno(), select.POLLIN, self._output_handler) def _output_handler(self, fd, event): files = self.files buf = array.array('B') try: - buf.fromfile(files["ebuild"], self._bufsize) + buf.fromfile(files.ebuild, self._bufsize) except EOFError: pass if buf: - buf.tofile(files["stdout"]) - files["stdout"].flush() - buf.tofile(files["log"]) - files["log"].flush() + buf.tofile(files.stdout) + files.stdout.flush() + buf.tofile(files.log) + files.log.flush() else: - self.unregister(files["ebuild"].fileno()) + self.unregister(files.ebuild.fileno()) for f in files.values(): f.close() + def poll(self): + if self.returncode is not None: + return self.returncode + retval = os.waitpid(self.pid, os.WNOHANG) + if retval == (0, 0): + return None + self._set_returncode(retval) + return self.returncode + def wait(self): - pid = self.pid - retval = os.waitpid(pid, 0)[1] - portage.process.spawned_pids.remove(pid) + if self.returncode is not None: + return self.returncode + self._set_returncode(os.waitpid(self.pid, 0)) + return self.returncode + + def _set_returncode(self, wait_retval): + + retval = wait_retval[1] + portage.process.spawned_pids.remove(self.pid) if retval != os.EX_OK: if retval & 0xff: retval = (retval & 0xff) << 8 @@ -1733,7 +1870,6 @@ class EbuildPhase(SlotObject): eerror(l, phase=self.phase, key=self.pkg.cpv) self.returncode = retval - return self.returncode class EbuildBinpkg(Task): """ @@ -1886,6 +2022,195 @@ class BinpkgFetcher(Task): rval = 1 return rval +class BinpkgFetcherAsync(SlotObject): + + __slots__ = ("cancelled", "log_file", "fd_pipes", "pkg", + "register", "unregister", + "locked", "files", "pid", "pkg_path", "returncode", "_lock_obj") + + _file_names = ("fetcher", "out") + _files_dict = slot_dict_class(_file_names, prefix="") + _bufsize = 4096 + + def __init__(self, **kwargs): + SlotObject.__init__(self, **kwargs) + pkg = self.pkg + self.pkg_path = pkg.root_config.trees["bintree"].getname(pkg.cpv) + + def start(self): + + if self.cancelled: + self.pid = -1 + return + + fd_pipes = self.fd_pipes + if fd_pipes is None: + fd_pipes = { + 0 : sys.stdin.fileno(), + 1 : sys.stdout.fileno(), + 2 : sys.stderr.fileno(), + } + + log_file = self.log_file + self.files = self._files_dict() + files = self.files + + if log_file is not None: + files.out = open(log_file, "a") + portage.util.apply_secpass_permissions(log_file, + uid=portage.portage_uid, gid=portage.portage_gid, + mode=0660) + else: + # 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() + + files.out = os.fdopen(os.dup(fd_pipes[1]), 'w') + + 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.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 + + pkg = self.pkg + bintree = pkg.root_config.trees["bintree"] + settings = bintree.settings + use_locks = "distlocks" in settings.features + pkg_path = self.pkg_path + resume = os.path.exists(pkg_path) + + # urljoin doesn't work correctly with + # unrecognized protocols like sftp + if bintree._remote_has_index: + rel_uri = bintree._remotepkgs[pkg.cpv].get("PATH") + if not rel_uri: + rel_uri = pkg.cpv + ".tbz2" + uri = bintree._remote_base_uri.rstrip("/") + \ + "/" + rel_uri.lstrip("/") + else: + uri = settings["PORTAGE_BINHOST"].rstrip("/") + \ + "/" + pkg.pf + ".tbz2" + + protocol = urlparse.urlparse(uri)[0] + fcmd_prefix = "FETCHCOMMAND" + if resume: + fcmd_prefix = "RESUMECOMMAND" + fcmd = settings.get(fcmd_prefix + "_" + protocol.upper()) + if not fcmd: + fcmd = settings.get(fcmd_prefix) + + fcmd_vars = { + "DISTDIR" : os.path.dirname(pkg_path), + "URI" : uri, + "FILE" : os.path.basename(pkg_path) + } + + fetch_env = dict((k, settings[k]) for k in settings) + fetch_args = [portage.util.varexpand(x, mydict=fcmd_vars) \ + for x in shlex.split(fcmd)] + + portage.util.ensure_dirs(os.path.dirname(pkg_path)) + if use_locks: + self.lock() + + retval = portage.process.spawn(fetch_args, env=fetch_env, + fd_pipes=fd_pipes, returnpid=True) + + self.pid = retval[0] + + os.close(slave_fd) + files.fetcher = os.fdopen(master_fd, 'r') + self.register(files.fetcher.fileno(), + select.POLLIN, self._output_handler) + + def _output_handler(self, fd, event): + files = self.files + buf = array.array('B') + try: + buf.fromfile(files.fetcher, self._bufsize) + except EOFError: + pass + if buf: + buf.tofile(files.out) + files.out.flush() + else: + self.unregister(files.fetcher.fileno()) + for f in files.values(): + f.close() + if self.locked: + self.unlock() + + def lock(self): + """ + This raises an AlreadyLocked exception if lock() is called + while a lock is already held. In order to avoid this, call + unlock() or check whether the "locked" attribute is True + or False before calling lock(). + """ + if self._lock_obj is not None: + raise self.AlreadyLocked((self._lock_obj,)) + + self._lock_obj = portage.locks.lockfile( + self.pkg_path, wantnewlockfile=1) + self.locked = True + + class AlreadyLocked(portage.exception.PortageException): + pass + + def unlock(self): + if self._lock_obj is None: + return + portage.locks.unlockfile(self._lock_obj) + self._lock_obj = None + self.locked = False + + def poll(self): + if self.returncode is not None: + return self.returncode + retval = os.waitpid(self.pid, os.WNOHANG) + if retval == (0, 0): + return None + self._set_returncode(retval) + return self.returncode + + def cancel(self): + if self.isAlive(): + os.kill(self.pid, signal.SIGTERM) + self.cancelled = True + if self.pid is not None: + self.wait() + return self.returncode + + def isAlive(self): + return self.pid is not None and \ + self.returncode is None + + def wait(self): + if self.returncode is not None: + return self.returncode + self._set_returncode(os.waitpid(self.pid, 0)) + return self.returncode + + def _set_returncode(self, wait_retval): + + retval = wait_retval[1] + portage.process.spawned_pids.remove(self.pid) + if retval != os.EX_OK: + if retval & 0xff: + retval = (retval & 0xff) << 8 + else: + retval = retval >> 8 + + self.returncode = retval + class BinpkgMerge(Task): __slots__ = ("find_blockers", "ldpath_mtimes", @@ -6368,6 +6693,8 @@ class Scheduler(object): "--fetchonly", "--fetch-all-uri", "--nodeps", "--pretend"]) + _fetch_log = EPREFIX + "/var/log/emerge-fetch.log" + def __init__(self, settings, trees, mtimedb, myopts, spinner, mergelist, favorites, digraph): self.settings = settings @@ -6386,9 +6713,38 @@ class Scheduler(object): self.pkgsettings[root] = portage.config( clone=trees[root]["vartree"].settings) self.curval = 0 - self._spawned_pids = [] self._poll_event_handlers = {} self._poll = select.poll() + from collections import deque + self._task_queue = deque() + self._running_tasks = set() + self._max_jobs = 1 + self._parallel_fetch = False + features = self.settings.features + if "parallel-fetch" in features and \ + not ("--pretend" in self.myopts or \ + "--fetch-all-uri" in self.myopts or \ + "--fetchonly" in self.myopts): + if "distlocks" not in features: + portage.writemsg(red("!!!")+"\n", noiselevel=-1) + portage.writemsg(red("!!!")+" parallel-fetching " + \ + "requires the distlocks feature enabled"+"\n", + noiselevel=-1) + portage.writemsg(red("!!!")+" you have it disabled, " + \ + "thus parallel-fetching is being disabled"+"\n", + noiselevel=-1) + portage.writemsg(red("!!!")+"\n", noiselevel=-1) + elif len(mergelist) > 1: + self._parallel_fetch = True + + # clear out existing fetch log if it exists + try: + open(self._fetch_log, 'w') + except EnvironmentError: + pass + + def _add_task(self, task): + self._task_queue.append(task) class _pkg_failure(portage.exception.PortageException): """ @@ -6441,20 +6797,17 @@ class Scheduler(object): def merge(self): keep_going = "--keep-going" in self.myopts + running_tasks = self._running_tasks while True: try: rval = self._merge() finally: - spawned_pids = self._spawned_pids - while spawned_pids: - pid = spawned_pids.pop() - try: - if os.waitpid(pid, os.WNOHANG) == (0, 0): - os.kill(pid, signal.SIGTERM) - os.waitpid(pid, 0) - except OSError: - pass # cleaned up elsewhere. + # clean up child process if necessary + self._task_queue.clear() + while running_tasks: + task = running_tasks.pop() + task.cancel() if rval == os.EX_OK or not keep_going: break @@ -6534,25 +6887,6 @@ class Scheduler(object): mydepgraph.break_refs(dropped_tasks) return (mylist, dropped_tasks) - def _poll_child_processes(self): - """ - After each merge, collect status from child processes - in order to clean up zombies (such as the parallel-fetch - process). - """ - spawned_pids = self._spawned_pids - if not spawned_pids: - return - for pid in list(spawned_pids): - try: - if os.waitpid(pid, os.WNOHANG) == (0, 0): - continue - except OSError: - # This pid has been cleaned up elsewhere, - # so remove it from our list. - pass - spawned_pids.remove(pid) - def _register(self, f, eventmask, handler): self._poll_event_handlers[f] = handler self._poll.register(f, eventmask) @@ -6560,11 +6894,45 @@ class Scheduler(object): def _unregister(self, f): self._poll.unregister(f) del self._poll_event_handlers[f] + self._schedule_tasks() def _schedule(self): - while self._poll_event_handlers: - for f, event in self._poll.poll(): - self._poll_event_handlers[f](f, event) + event_handlers = self._poll_event_handlers + running_tasks = self._running_tasks + poll = self._poll.poll + + self._schedule_tasks() + + while event_handlers: + 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 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 _merge(self): mylist = self._mergelist @@ -6592,6 +6960,28 @@ class Scheduler(object): if isinstance(x, Package) and x.operation == "merge"] mtimedb.commit() + prefetchers = weakref.WeakValueDictionary() + getbinpkg = "--getbinpkg" in self.myopts + + if self._parallel_fetch: + for pkg in mylist: + if not isinstance(pkg, Package): + continue + if pkg.type_name == "ebuild": + self._add_task(EbuildFetcherAsync( + log_file=self._fetch_log, + pkg=pkg, register=self._register, + unregister=self._unregister)) + elif pkg.type_name == "binary" and getbinpkg and \ + pkg.root_config.trees["bintree"].isremote(pkg.cpv): + prefetcher = BinpkgFetcherAsync( + log_file=self._fetch_log, + pkg=pkg, register=self._register, + unregister=self._unregister) + prefetchers[pkg] = prefetcher + self._add_task(prefetcher) + del prefetcher + # Verify all the manifests now so that the user is notified of failure # as soon as possible. if "--fetchonly" not in self.myopts and \ @@ -6625,49 +7015,6 @@ class Scheduler(object): myfeat = self.settings.features[:] bad_resume_opts = set(["--ask", "--changelog", "--skipfirst", "--resume"]) - if "parallel-fetch" in myfeat and \ - not ("--pretend" in self.myopts or \ - "--fetch-all-uri" in self.myopts or \ - "--fetchonly" in self.myopts): - if "distlocks" not in myfeat: - print red("!!!") - print red("!!!")+" parallel-fetching requires the distlocks feature enabled" - print red("!!!")+" you have it disabled, thus parallel-fetching is being disabled" - print red("!!!") - elif len(mymergelist) > 1: - fetch_log = EPREFIX+"/var/log/emerge-fetch.log" - logfile = open(fetch_log, "w") - fd_pipes = {1:logfile.fileno(), 2:logfile.fileno()} - portage.util.apply_secpass_permissions(fetch_log, - uid=portage.portage_uid, gid=portage.portage_gid, - mode=0660) - fetch_env = os.environ.copy() - fetch_env["FEATURES"] = fetch_env.get("FEATURES", "") + " -cvs" - fetch_env["PORTAGE_NICENESS"] = "0" - fetch_env["PORTAGE_PARALLEL_FETCHONLY"] = "1" - fetch_args = [sys.argv[0], "--resume", - "--fetchonly", "--nodeps"] - resume_opts = self.myopts.copy() - # For automatic resume, we need to prevent - # any of bad_resume_opts from leaking in - # via EMERGE_DEFAULT_OPTS. - resume_opts["--ignore-default-opts"] = True - for myopt, myarg in resume_opts.iteritems(): - if myopt not in bad_resume_opts: - if myarg is True: - fetch_args.append(myopt) - else: - fetch_args.append(myopt +"="+ myarg) - self._spawned_pids.extend( - portage.process.spawn( - fetch_args, env=fetch_env, - fd_pipes=fd_pipes, returnpid=True)) - logfile.close() # belongs to the spawned process - del fetch_log, logfile, fd_pipes, fetch_env, fetch_args, \ - resume_opts - print ">>> starting parallel fetching pid %d" % \ - self._spawned_pids[-1] - metadata_keys = [k for k in portage.auxdbkeys \ if not k.startswith("UNUSED_")] + ["USE"] @@ -6704,14 +7051,15 @@ class Scheduler(object): self._execute_task(bad_resume_opts, failed_fetches, mydbapi, mergecount, - myfeat, mymergelist, x, xterm_titles) + myfeat, mymergelist, x, + prefetchers, xterm_titles) except self._pkg_failure, e: return e.status return self._post_merge(mtimedb, xterm_titles, failed_fetches) def _execute_task(self, bad_resume_opts, failed_fetches, mydbapi, mergecount, myfeat, - mymergelist, pkg, xterm_titles): + mymergelist, pkg, prefetchers, xterm_titles): favorites = self._favorites mtimedb = self._mtimedb from portage.elog import elog_process @@ -6862,8 +7210,27 @@ class Scheduler(object): phasefilter=filter_mergephases) build_dir.unlock() - elif x[0]=="binary": - #merge the tbz2 + elif x.type_name == "binary": + # The prefetcher have already completed or it + # could be running now. If it's running now, + # wait for it to complete since it holds + # a lock on the file being fetched. The + # portage.locks functions are only designed + # to work between separate processes. Since + # the lock is held by the current process, + # use the scheduler and fetcher methods to + # synchronize with the fetcher. + prefetcher = prefetchers.get(pkg) + if prefetcher is not None: + if not prefetcher.isAlive(): + prefetcher.cancel() + else: + retval = None + while retval is None: + self._schedule() + retval = prefetcher.poll() + del prefetcher + fetcher = BinpkgFetcher(pkg=pkg, pretend=pretend, use_locks=("distlocks" in pkgsettings.features)) mytbz2 = fetcher.pkg_path @@ -6967,7 +7334,6 @@ class Scheduler(object): # due to power failure, SIGKILL, etc... mtimedb.commit() self.curval += 1 - self._poll_child_processes() def _post_merge(self, mtimedb, xterm_titles, failed_fetches): if "--pretend" not in self.myopts: diff --git a/pym/portage/__init__.py b/pym/portage/__init__.py index 101efdbdb..d1c453133 100644 --- a/pym/portage/__init__.py +++ b/pym/portage/__init__.py @@ -3294,8 +3294,9 @@ def fetch(myuris, mysettings, listonly=0, fetchonly=0, locks_in_subdir=".locks", # file size. The parent process will verify their checksums prior to # the unpack phase. - parallel_fetchonly = fetchonly and \ - "PORTAGE_PARALLEL_FETCHONLY" in mysettings + parallel_fetchonly = "PORTAGE_PARALLEL_FETCHONLY" in mysettings + if parallel_fetchonly: + fetchonly = 1 check_config_instance(mysettings) diff --git a/pym/portage/cache/mappings.py b/pym/portage/cache/mappings.py index 2cddd8147..2ccc96b05 100644 --- a/pym/portage/cache/mappings.py +++ b/pym/portage/cache/mappings.py @@ -104,14 +104,17 @@ class LazyLoad(UserDict.DictMixin): _slot_dict_classes = weakref.WeakValueDictionary() -def slot_dict_class(keys): +def slot_dict_class(keys, prefix="_val_"): """ Generates mapping classes that behave similar to a dict but store values as object attributes that are allocated via __slots__. Instances of these objects have a smaller memory footprint than a normal dict object. @param keys: Fixed set of allowed keys - @type keys: iterable + @type keys: Iterable + @param prefix: a prefix to use when mapping + attribute names from keys + @type prefix: String @rtype: SlotDict @returns: A class that constructs SlotDict instances having the specified keys. @@ -120,14 +123,15 @@ def slot_dict_class(keys): keys_set = keys else: keys_set = frozenset(keys) - v = _slot_dict_classes.get(keys_set) + v = _slot_dict_classes.get((keys_set, prefix)) if v is None: class SlotDict(object): allowed_keys = keys_set + _prefix = prefix __slots__ = ("__weakref__",) + \ - tuple("_val_" + k for k in allowed_keys) + tuple(prefix + k for k in allowed_keys) def __iter__(self): for k, v in self.iteritems(): @@ -145,7 +149,7 @@ def slot_dict_class(keys): def iteritems(self): for k in self.allowed_keys: try: - yield (k, getattr(self, "_val_" + k)) + yield (k, getattr(self, self._prefix + k)) except AttributeError: pass @@ -161,12 +165,12 @@ def slot_dict_class(keys): def __delitem__(self, k): try: - delattr(self, "_val_" + k) + delattr(self, self._prefix + k) except AttributeError: raise KeyError(k) def __setitem__(self, k, v): - setattr(self, "_val_" + k, v) + setattr(self, self._prefix + k, v) def setdefault(self, key, default=None): try: @@ -186,7 +190,7 @@ def slot_dict_class(keys): def __getitem__(self, k): try: - return getattr(self, "_val_" + k) + return getattr(self, self._prefix + k) except AttributeError: raise KeyError(k) @@ -197,7 +201,7 @@ def slot_dict_class(keys): return default def __contains__(self, k): - return hasattr(self, "_val_" + k) + return hasattr(self, self._prefix + k) def has_key(self, k): return k in self @@ -232,7 +236,7 @@ def slot_dict_class(keys): def clear(self): for k in self.allowed_keys: try: - delattr(self, "_val_" + k) + delattr(self, self._prefix + k) except AttributeError: pass diff --git a/pym/portage/checksum.py b/pym/portage/checksum.py index 77716aefc..52ce59148 100644 --- a/pym/portage/checksum.py +++ b/pym/portage/checksum.py @@ -11,7 +11,6 @@ import tempfile import portage.exception import portage.process import commands -import md5, sha #dict of all available hash functions hashfunc_map = {} @@ -46,8 +45,19 @@ def _generate_hash_function(hashtype, hashobject, origin="unknown"): # override earlier ones # Use the internal modules as last fallback -md5hash = _generate_hash_function("MD5", md5.new, origin="internal") -sha1hash = _generate_hash_function("SHA1", sha.new, origin="internal") +try: + from hashlib import md5 as _new_md5 +except ImportError: + from md5 import new as _new_md5 + +md5hash = _generate_hash_function("MD5", _new_md5, origin="internal") + +try: + from hashlib import sha1 as _new_sha1 +except ImportError: + from sha import new as _new_sha1 + +sha1hash = _generate_hash_function("SHA1", _new_sha1, origin="internal") # Use pycrypto when available, prefer it over the internal fallbacks try: diff --git a/pym/portage/dbapi/porttree.py b/pym/portage/dbapi/porttree.py index b6e39f63b..e2a53aac4 100644 --- a/pym/portage/dbapi/porttree.py +++ b/pym/portage/dbapi/porttree.py @@ -63,7 +63,9 @@ class portdbapi(dbapi): self.manifestVerifier = portage.gpg.FileChecker(self.mysettings["PORTAGE_GPG_DIR"], "gentoo.gpg", minimumTrust=self.manifestVerifyLevel) #self.root=settings["PORTDIR"] - self.porttree_root = os.path.realpath(porttree_root) + self.porttree_root = porttree_root + if porttree_root: + self.porttree_root = os.path.realpath(porttree_root) self.depcachedir = os.path.realpath(self.mysettings.depcachedir) diff --git a/pym/portage/dbapi/vartree.py b/pym/portage/dbapi/vartree.py index 8b4093a87..760eeb986 100644 --- a/pym/portage/dbapi/vartree.py +++ b/pym/portage/dbapi/vartree.py @@ -46,7 +46,8 @@ class PreservedLibsRegistry(object): self._filename = filename self._autocommit = autocommit self.load() - + self.pruneNonExisting() + def load(self): """ Reload the registry data from file """ try: @@ -63,9 +64,15 @@ class PreservedLibsRegistry(object): """ Store the registry data to file. No need to call this if autocommit was enabled. """ - f = atomic_ofstream(self._filename) - cPickle.dump(self._data, f) - f.close() + if os.environ.get("SANDBOX_ON") == "1": + return + try: + f = atomic_ofstream(self._filename) + cPickle.dump(self._data, f) + f.close() + except EnvironmentError, e: + if e.errno != PermissionDenied.errno: + writemsg("!!! %s %s\n" % (e, self._filename), noiselevel=-1) def register(self, cpv, slot, counter, paths): """ Register new objects in the registry. If there is a record with the @@ -1030,6 +1037,35 @@ class vardbapi(dbapi): return dblink(category, pf, self.root, self.settings, vartree=self.vartree) + def removeFromContents(self, pkg, paths, relative_paths=True): + """ + @param pkg: cpv for an installed package + @type pkg: string + @param paths: paths of files to remove from contents + @type paths: iterable + """ + if not hasattr(pkg, "getcontents"): + pkg = self._dblink(pkg) + root = self.root + root_len = len(root) - 1 + new_contents = pkg.getcontents().copy() + contents_key = None + + for filename in paths: + filename = normalize_path(filename) + if relative_paths: + relative_filename = filename + else: + relative_filename = filename[root_len:] + contents_key = pkg._match_contents(relative_filename, root) + if contents_key: + del new_contents[contents_key] + + if contents_key: + f = atomic_ofstream(os.path.join(pkg.dbdir, "CONTENTS")) + write_contents(new_contents, root, f) + f.close() + class _owners_cache(object): """ This class maintains an hash table that serves to index package @@ -2061,7 +2097,7 @@ class dblink(object): #remove self from vartree database so that our own virtual gets zapped if we're the last node self.vartree.zap(self.mycpv) - def isowner(self,filename, destroot): + def isowner(self, filename, destroot): """ Check if a file belongs to this package. This may result in a stat call for the parent directory of @@ -2080,12 +2116,25 @@ class dblink(object): 1. True if this package owns the file. 2. False if this package does not own the file. """ + return bool(self._match_contents(filename, destroot)) + + def _match_contents(self, filename, destroot): + """ + The matching contents entry is returned, which is useful + since the path may differ from the one given by the caller, + due to symlinks. + + @rtype: String + @return: the contents entry corresponding to the given path, or False + if the file is not owned by this package. + """ + destfile = normalize_path( os.path.join(destroot, filename.lstrip(os.path.sep))) pkgfiles = self.getcontents() if pkgfiles and destfile in pkgfiles: - return True + return destfile if pkgfiles: basename = os.path.basename(destfile) if self._contents_basenames is None: @@ -2135,7 +2184,7 @@ class dblink(object): for p_path in p_path_list: x = os.path.join(p_path, basename) if x in pkgfiles: - return True + return x return False @@ -2803,33 +2852,8 @@ class dblink(object): contents = self.getcontents() destroot_len = len(destroot) - 1 for blocker in blockers: - blocker_contents = blocker.getcontents() - collisions = [] - for filename in blocker_contents: - relative_filename = filename[destroot_len:] - if self.isowner(relative_filename, destroot): - collisions.append(filename) - if not collisions: - continue - for filename in collisions: - del blocker_contents[filename] - f = atomic_ofstream(os.path.join(blocker.dbdir, "CONTENTS")) - for filename in sorted(blocker_contents): - entry_data = blocker_contents[filename] - entry_type = entry_data[0] - relative_filename = filename[destroot_len:] - if entry_type == "obj": - entry_type, mtime, md5sum = entry_data - line = "%s %s %s %s\n" % \ - (entry_type, relative_filename, md5sum, mtime) - elif entry_type == "sym": - entry_type, mtime, link = entry_data - line = "%s %s -> %s %s\n" % \ - (entry_type, relative_filename, link, mtime) - else: # dir, dev, fif - line = "%s %s\n" % (entry_type, relative_filename) - f.write(line) - f.close() + self.vartree.dbapi.removeFromContents(blocker, iter(contents), + relative_paths=False) self.vartree.dbapi._add(self) contents = self.getcontents() @@ -3240,6 +3264,27 @@ class dblink(object): "Is this a regular package (does it have a CATEGORY file? A dblink can be virtual *and* regular)" return os.path.exists(os.path.join(self.dbdir, "CATEGORY")) +def write_contents(contents, root, f): + """ + Write contents to any file like object. The file will be left open. + """ + root_len = len(root) - 1 + for filename in sorted(contents): + entry_data = contents[filename] + entry_type = entry_data[0] + relative_filename = filename[root_len:] + if entry_type == "obj": + entry_type, mtime, md5sum = entry_data + line = "%s %s %s %s\n" % \ + (entry_type, relative_filename, md5sum, mtime) + elif entry_type == "sym": + entry_type, mtime, link = entry_data + line = "%s %s -> %s %s\n" % \ + (entry_type, relative_filename, link, mtime) + else: # dir, dev, fif + line = "%s %s\n" % (entry_type, relative_filename) + f.write(line) + def tar_contents(contents, root, tar, protect=None, onProgress=None): from portage.util import normalize_path import tarfile -- 2.26.2