Merged from trunk 10853:10869
authorFabian Groffen <grobian@gentoo.org>
Tue, 1 Jul 2008 17:22:29 +0000 (17:22 -0000)
committerFabian Groffen <grobian@gentoo.org>
Tue, 1 Jul 2008 17:22:29 +0000 (17:22 -0000)
   | 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
man/portage.5
pym/_emerge/__init__.py
pym/portage/__init__.py
pym/portage/cache/mappings.py
pym/portage/checksum.py
pym/portage/dbapi/porttree.py
pym/portage/dbapi/vartree.py

index 0003acda47855d38b72fd2790cb1f120bcbd6beb..58f048725e72a1e9f3e0b75cbcf2a5cd2dd0a694 100755 (executable)
@@ -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")
index d6eb83cb643de0cb5a3fd32910aa9cfae28007e7..d6a29678a4b19edc69e771c5081258c9b6b0d5ce 100644 (file)
@@ -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:
index 90162d7e90bdfd9af0f95b3bad25e672a6244434..65c01a6d80a936f42251b63dde829b030270a935 100644 (file)
@@ -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:
index 101efdbdb51b9768d4ca43ccf741e802b45f517c..d1c453133575c302d786bcb26ac5f708df685e4a 100644 (file)
@@ -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)
 
index 2cddd8147a3ca702cf7eeb139ba035cbc164c0bb..2ccc96b05c87844bc4ea69bf4c1856f84e20ccc8 100644 (file)
@@ -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
 
index 77716aefc44623c697508f63ab378f79bc441f85..52ce59148f9d9b11139b48beecbf1b2895d26c92 100644 (file)
@@ -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:
index b6e39f63b44d4c4c74610eb9b4d35a04d9c3d748..e2a53aac4443375f02957ecdee3cea3460bef560 100644 (file)
@@ -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)
 
index 8b4093a87736e7d11c98feb7080adbb023b12474..760eeb98697af0e8b0791c81fcaca764dae2f3cd 100644 (file)
@@ -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