sys.path.insert(0, os.path.join(os.path.dirname(os.path.dirname(os.path.realpath(__file__))), "pym"))
import portage
-portage._disable_legacy_globals()
+
from portage import digraph, portdbapi
from portage.const import NEWS_LIB_PATH, CACHE_PATH, PRIVATE_PATH, USER_CONFIG_PATH, GLOBAL_CONFIG_PATH
(self.type_name, self.root, self.cpv, self.operation)
return self._hash_key
+ def __cmp__(self, other):
+ if self > other:
+ return 1
+ elif self < other:
+ return -1
+ return 0
+
def __lt__(self, other):
if other.cp != self.cp:
return False
"""
if self.returncode is not None and \
self._exit_listeners is not None:
- for f in self._exit_listeners:
+
+ # This prevents recursion, in case one of the
+ # exit handlers triggers this method again by
+ # calling wait().
+ exit_listeners = self._exit_listeners
+ self._exit_listeners = None
+
+ for f in exit_listeners:
f(self)
self._exit_listeners = None
return self.returncode
if self.pid is None:
return self.returncode
- retval = os.waitpid(self.pid, os.WNOHANG)
+
+ try:
+ retval = os.waitpid(self.pid, os.WNOHANG)
+ except OSError, e:
+ if e.errno != errno.ECHILD:
+ raise
+ del e
+ retval = (self.pid, 1)
+
if retval == (0, 0):
return None
self._set_returncode(retval)
def cancel(self):
if self.isAlive():
- os.kill(self.pid, signal.SIGTERM)
+ try:
+ os.kill(self.pid, signal.SIGTERM)
+ except OSError, e:
+ if e.errno != errno.ESRCH:
+ raise
+ del e
+
self.cancelled = True
if self.pid is not None:
self.wait()
return self.returncode
if self.registered:
self.scheduler.schedule(self._reg_id)
- self._set_returncode(os.waitpid(self.pid, 0))
+ try:
+ wait_retval = os.waitpid(self.pid, 0)
+ except OSError, e:
+ if e.errno != errno.ECHILD:
+ raise
+ del e
+ self._set_returncode((self.pid, 1))
+ else:
+ self._set_returncode(wait_retval)
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
retval = portage.process.spawn(self.args, **kwargs)
self.pid = retval[0]
+ portage.process.spawned_pids.remove(self.pid)
os.close(slave_fd)
files.process = os.fdopen(master_fd, 'r')
f.flush()
f.close()
self.registered = False
- self._wait()
+ self.scheduler.unregister(self._reg_id)
+ self.wait()
return self.registered
def _dummy_handler(self, fd, event):
for f in files.values():
f.close()
self.registered = False
- self._wait()
+ self.scheduler.unregister(self._reg_id)
+ self.wait()
return self.registered
class EbuildFetcher(SpawnProcess):
- __slots__ = ("pkg",)
+ __slots__ = ("fetchonly", "pkg",)
def start(self):
fetch_env = settings.environ()
fetch_env["PORTAGE_NICENESS"] = "0"
- fetch_env["PORTAGE_PARALLEL_FETCHONLY"] = "1"
+ if self.fetchonly:
+ fetch_env["PORTAGE_PARALLEL_FETCHONLY"] = "1"
ebuild_binary = os.path.join(
settings["EBUILD_BIN_PATH"], "ebuild")
__slots__ = ("args_set", "background", "find_blockers",
"ldpath_mtimes", "logger", "opts", "pkg", "pkg_count",
- "settings", "world_atom") + \
+ "prefetcher", "settings", "world_atom") + \
("_build_dir", "_buildpkg", "_ebuild_path", "_tree")
def start(self):
- args_set = self.args_set
- find_blockers = self.find_blockers
- ldpath_mtimes = self.ldpath_mtimes
logger = self.logger
opts = self.opts
pkg = self.pkg
- pkg_count = self.pkg_count
- scheduler = self.scheduler
settings = self.settings
world_atom = self.world_atom
root_config = pkg.root_config
- root = root_config.root
- system_set = root_config.sets["system"]
- world_set = root_config.sets["world"]
- vartree = root_config.trees["vartree"]
tree = "porttree"
self._tree = tree
portdb = root_config.trees[tree].dbapi
- debug = settings.get("PORTAGE_DEBUG") == "1"
- features = self.settings.features
settings["EMERGE_FROM"] = pkg.type_name
settings.backup_changes("EMERGE_FROM")
settings.reset()
ebuild_path = portdb.findname(self.pkg.cpv)
self._ebuild_path = ebuild_path
- #buildsyspkg: Check if we need to _force_ binary package creation
- issyspkg = "buildsyspkg" in features and \
- system_set.findAtomForPackage(pkg) and \
- not opts.buildpkg
+ prefetcher = self.prefetcher
+ if prefetcher is None:
+ pass
+ elif not prefetcher.isAlive():
+ prefetcher.cancel()
+ elif prefetcher.poll() is None:
- if opts.fetchonly:
- if opts.pretend:
+ waiting_msg = "Fetching files " + \
+ "in the background. " + \
+ "To view fetch progress, run `tail -f " + \
+ "/var/log/emerge-fetch.log` in another " + \
+ "terminal."
+ msg_prefix = colorize("GOOD", " * ")
+ from textwrap import wrap
+ waiting_msg = "".join("%s%s\n" % (msg_prefix, line) \
+ for line in wrap(waiting_msg, 65))
+ if not self.background:
+ writemsg(waiting_msg, noiselevel=-1)
+
+ self._current_task = prefetcher
+ prefetcher.addExitListener(self._prefetch_exit)
+ return
+
+ self._prefetch_exit(prefetcher)
+
+ def _prefetch_exit(self, prefetcher):
+
+ opts = self.opts
+ pkg = self.pkg
+ settings = self.settings
+
+ if opts.fetchonly and opts.pretend:
fetcher = EbuildFetchPretend(
fetch_all=opts.fetch_all_uri,
pkg=pkg, settings=settings)
retval = fetcher.execute()
self.returncode = retval
self.wait()
+ return
- else:
- fetcher = EbuildFetcher(pkg=pkg, scheduler=scheduler)
- self._start_task(fetcher, self._fetchonly_exit)
+ fetch_log = None
+ if self.background:
+ fetch_log = self.scheduler.fetch.log_file
+
+ fetcher = EbuildFetcher(fetchonly=opts.fetchonly,
+ background=self.background, logfile=fetch_log,
+ pkg=pkg, scheduler=self.scheduler)
+ if self.background:
+ fetcher.addExitListener(self._fetch_exit)
+ self._current_task = fetcher
+ self.scheduler.fetch.schedule(fetcher)
+ else:
+ self._start_task(fetcher, self._fetch_exit)
+
+ def _fetch_exit(self, fetcher):
+
+ opts = self.opts
+ pkg = self.pkg
+
+ if opts.fetchonly:
+ if self._final_exit(fetcher) != os.EX_OK:
+ eerror("!!! Fetch for %s failed, continuing..." % pkg.cpv,
+ phase="unpack", key=pkg.cpv)
+ self.wait()
+ return
+
+ if self._default_exit(fetcher) != os.EX_OK:
+ self.wait()
return
+ logger = self.logger
+ opts = self.opts
+ pkg_count = self.pkg_count
+ scheduler = self.scheduler
+ settings = self.settings
+ features = settings.features
+ ebuild_path = self._ebuild_path
+ system_set = pkg.root_config.sets["system"]
+
self._build_dir = EbuildBuildDir(pkg=pkg, settings=settings)
self._build_dir.lock()
(pkg_count.curval, pkg_count.maxval, pkg.cpv)
logger.log(msg, short_msg=short_msg)
+ #buildsyspkg: Check if we need to _force_ binary package creation
+ issyspkg = "buildsyspkg" in features and \
+ system_set.findAtomForPackage(pkg) and \
+ not opts.buildpkg
+
if opts.buildpkg or issyspkg:
self._buildpkg = True
scheduler=scheduler, settings=settings)
self._start_task(build, self._build_exit)
- def _fetchonly_exit(self, fetcher):
- if self._final_exit(fetcher) != os.EX_OK:
- pkg = self.pkg
- eerror("!!! Fetch for %s failed, continuing..." % pkg.cpv,
- phase="unpack", key=pkg.cpv)
- self.wait()
-
def _unlock_builddir(self):
portage.elog.elog_process(self.pkg.cpv, self.settings)
self._build_dir.unlock()
self._start_task(ebuild_phases, self._default_final_exit)
+class EbuildMetadataPhase(SubProcess):
+
+ """
+ Asynchronous interface for the ebuild "depend" phase which is
+ used to extract metadata from the ebuild.
+ """
+
+ __slots__ = ("cpv", "ebuild_path", "fd_pipes", "metadata_callback",
+ "ebuild_mtime", "portdb", "repo_path", "settings") + \
+ ("files", "_raw_metadata")
+
+ _file_names = ("ebuild",)
+ _files_dict = slot_dict_class(_file_names, prefix="")
+ _bufsize = SpawnProcess._bufsize
+ _metadata_fd = 9
+
+ def start(self):
+ settings = self.settings
+ settings.reset()
+ ebuild_path = self.ebuild_path
+ debug = settings.get("PORTAGE_DEBUG") == "1"
+ master_fd = None
+ slave_fd = None
+ fd_pipes = None
+ if self.fd_pipes is not None:
+ fd_pipes = self.fd_pipes.copy()
+ else:
+ fd_pipes = {}
+
+ fd_pipes.setdefault(0, sys.stdin.fileno())
+ fd_pipes.setdefault(1, sys.stdout.fileno())
+ fd_pipes.setdefault(2, sys.stderr.fileno())
+
+ # flush any pending output
+ for fd in fd_pipes.itervalues():
+ if fd == sys.stdout.fileno():
+ sys.stdout.flush()
+ if fd == sys.stderr.fileno():
+ sys.stderr.flush()
+
+ fd_pipes_orig = fd_pipes.copy()
+ self.files = self._files_dict()
+ files = self.files
+
+ master_fd, slave_fd = os.pipe()
+ fcntl.fcntl(master_fd, fcntl.F_SETFL,
+ fcntl.fcntl(master_fd, fcntl.F_GETFL) | os.O_NONBLOCK)
+
+ fd_pipes[self._metadata_fd] = slave_fd
+
+ retval = portage.doebuild(ebuild_path, "depend",
+ settings["ROOT"], settings, debug,
+ mydbapi=self.portdb, tree="porttree",
+ fd_pipes=fd_pipes, returnpid=True)
+
+ os.close(slave_fd)
+
+ if isinstance(retval, int):
+ # doebuild failed before spawning
+ os.close(master_fd)
+ self.returncode = retval
+ self.wait()
+ return
+
+ self.pid = retval[0]
+ portage.process.spawned_pids.remove(self.pid)
+
+ self._raw_metadata = []
+ files.ebuild = os.fdopen(master_fd, 'r')
+ self._reg_id = self.scheduler.register(files.ebuild.fileno(),
+ PollConstants.POLLIN, self._output_handler)
+ self.registered = True
+
+ def _output_handler(self, fd, event):
+ files = self.files
+ self._raw_metadata.append(files.ebuild.read())
+ if not self._raw_metadata[-1]:
+ for f in files.values():
+ f.close()
+ self.registered = False
+ self.scheduler.unregister(self._reg_id)
+ self.wait()
+
+ if self.returncode == os.EX_OK:
+ metadata = izip(portage.auxdbkeys,
+ "".join(self._raw_metadata).splitlines())
+ self.metadata_callback(self.cpv, self.ebuild_path,
+ self.repo_path, metadata, self.ebuild_mtime)
+
+ return self.registered
+
class EbuildPhase(SubProcess):
__slots__ = ("fd_pipes", "phase", "pkg",
mydbapi=mydbapi, tree=tree,
fd_pipes=fd_pipes, returnpid=True)
+ os.close(slave_fd)
+
+ if isinstance(retval, int):
+ # doebuild failed before spawning
+ os.close(master_fd)
+ self.returncode = retval
+ self.wait()
+ return
+
self.pid = retval[0]
+ portage.process.spawned_pids.remove(self.pid)
if logfile:
files.log = open(logfile, 'a')
else:
output_handler = self._dummy_handler
- os.close(slave_fd)
files.ebuild = os.fdopen(master_fd, 'r')
self._reg_id = self.scheduler.register(files.ebuild.fileno(),
PollConstants.POLLIN, output_handler)
for f in files.values():
f.close()
self.registered = False
- self._wait()
+ self.scheduler.unregister(self._reg_id)
+ self.wait()
return self.registered
def _dummy_handler(self, fd, event):
for f in files.values():
f.close()
self.registered = False
- self._wait()
+ self.scheduler.unregister(self._reg_id)
+ self.wait()
return self.registered
def _set_returncode(self, wait_retval):
from textwrap import wrap
waiting_msg = "".join("%s%s\n" % (msg_prefix, line) \
for line in wrap(waiting_msg, 65))
- writemsg(waiting_msg, noiselevel=-1)
+ if not self.background:
+ writemsg(waiting_msg, noiselevel=-1)
self._current_task = prefetcher
prefetcher.addExitListener(self._prefetch_exit)
__slots__ = ("args_set",
"binpkg_opts", "build_opts", "emerge_opts",
"failed_fetches", "find_blockers", "logger", "mtimedb", "pkg",
- "pkg_count", "prefetcher", "settings", "world_atom") + \
+ "pkg_count", "pkg_to_replace", "prefetcher",
+ "settings", "world_atom") + \
("_install_task",)
def start(self):
find_blockers=find_blockers,
ldpath_mtimes=ldpath_mtimes, logger=logger,
opts=build_opts, pkg=pkg, pkg_count=pkg_count,
- settings=settings, scheduler=scheduler,
- world_atom=world_atom)
+ prefetcher=self.prefetcher, scheduler=scheduler,
+ settings=settings, world_atom=world_atom)
self._install_task = build
self._start_task(build, self._ebuild_exit)
class SequentialTaskQueue(SlotObject):
- __slots__ = ("max_jobs", "running_tasks", "_task_queue")
+ __slots__ = ("max_jobs", "running_tasks", "_task_queue", "_scheduling")
def __init__(self, **kwargs):
SlotObject.__init__(self, **kwargs)
def add(self, task):
self._task_queue.append(task)
+ self.schedule()
def addFront(self, task):
self._task_queue.appendleft(task)
+ self.schedule()
def schedule(self):
if not self:
return False
+ if self._scheduling:
+ # Ignore any recursive schedule() calls triggered via
+ # self._task_exit().
+ return False
+
+ self._scheduling = True
+
task_queue = self._task_queue
running_tasks = self.running_tasks
max_jobs = self.max_jobs
task = task_queue.popleft()
cancelled = getattr(task, "cancelled", None)
if not cancelled:
+ task.addExitListener(self._task_exit)
task.start()
running_tasks.add(task)
state_changed = True
+ self._scheduling = False
+
return state_changed
+ def _task_exit(self, task):
+ self.schedule()
+
def clear(self):
self._task_queue.clear()
running_tasks = self.running_tasks
def __len__(self):
return len(self._task_queue) + len(self.running_tasks)
-class Scheduler(object):
+class PollLoop(object):
+
+ def __init__(self):
+ self._max_jobs = 1
+ self._max_load = None
+ self._jobs = 0
+ self._poll_event_handlers = {}
+ self._poll_event_handler_ids = {}
+ # Increment id for each new handler.
+ self._event_handler_id = 0
+ try:
+ self._poll = select.poll()
+ except AttributeError:
+ self._poll = PollSelectAdapter()
+
+ def _can_add_job(self):
+ jobs = self._jobs
+ max_jobs = self._max_jobs
+ max_load = self._max_load
+
+ if self._jobs >= self._max_jobs:
+ return False
+
+ if max_load is not None and max_jobs > 1 and self._jobs > 1:
+ try:
+ avg1, avg5, avg15 = os.getloadavg()
+ except OSError, e:
+ writemsg("!!! getloadavg() failed: %s\n" % (e,),
+ noiselevel=-1)
+ del e
+ return False
+
+ if avg1 >= max_load:
+ return False
+
+ return True
+
+ def _poll_loop(self):
+
+ event_handlers = self._poll_event_handlers
+ poll = self._poll.poll
+ state_change = 0
+
+ while event_handlers:
+ for f, event in poll():
+ handler, reg_id = event_handlers[f]
+ if not handler(f, event):
+ state_change += 1
+
+ if not state_change:
+ raise AssertionError("tight loop")
+
+ def _register(self, f, eventmask, handler):
+ """
+ @rtype: Integer
+ @return: A unique registration id, for use in schedule() or
+ unregister() calls.
+ """
+ if f in self._poll_event_handlers:
+ raise AssertionError("fd %d is already registered" % f)
+ self._event_handler_id += 1
+ reg_id = self._event_handler_id
+ self._poll_event_handler_ids[reg_id] = f
+ self._poll_event_handlers[f] = (handler, reg_id)
+ self._poll.register(f, eventmask)
+ return reg_id
+
+ def _unregister(self, reg_id):
+ f = self._poll_event_handler_ids[reg_id]
+ self._poll.unregister(f)
+ del self._poll_event_handlers[f]
+ del self._poll_event_handler_ids[reg_id]
+
+ def _schedule(self, wait_id):
+ """
+ Schedule until wait_id is not longer registered
+ for poll() events.
+ @type wait_id: int
+ @param wait_id: a task id to wait for
+ """
+ event_handlers = self._poll_event_handlers
+ handler_ids = self._poll_event_handler_ids
+ poll = self._poll.poll
+
+ while wait_id in handler_ids:
+ for f, event in poll():
+ handler, reg_id = event_handlers[f]
+ handler(f, event)
+
+class Scheduler(PollLoop):
_opts_ignore_blockers = \
frozenset(["--buildpkgonly",
_fetch_log = EPREFIX + "/var/log/emerge-fetch.log"
class _iface_class(SlotObject):
- __slots__ = ("fetch", "register", "schedule")
+ __slots__ = ("fetch", "register", "schedule", "unregister")
class _fetch_iface_class(SlotObject):
__slots__ = ("log_file", "schedule")
def __init__(self, settings, trees, mtimedb, myopts,
spinner, mergelist, favorites, digraph):
+ PollLoop.__init__(self)
self.settings = settings
self.target_root = settings["ROOT"]
self.trees = trees
schedule=self._schedule_fetch)
self._sched_iface = self._iface_class(
fetch=fetch_iface, register=self._register,
- schedule=self._schedule)
- self._poll_event_handlers = {}
- self._poll_event_handler_ids = {}
- # Increment id for each new handler.
- self._event_handler_id = 0
-
- try:
- self._poll = select.poll()
- except AttributeError:
- self._poll = PollSelectAdapter()
+ schedule=self._schedule, unregister=self._unregister)
self._task_queues = self._task_queues_class()
for k in self._task_queues.allowed_keys:
self._max_load = myopts.get("--load-average")
self._set_digraph(digraph)
- self._jobs = 0
features = self.settings.features
if "parallel-fetch" in features and \
except EnvironmentError:
pass
+ self._running_portage = None
+ portage_match = self._running_root.trees["vartree"].dbapi.match(
+ portage.const.PORTAGE_PACKAGE_ATOM)
+ if portage_match:
+ cpv = portage_match.pop()
+ self._running_portage = self._pkg(cpv, "installed",
+ self._running_root, installed=True)
+
def _set_max_jobs(self, max_jobs):
self._max_jobs = max_jobs
self._task_queues.jobs.max_jobs = max_jobs
elif pkg.type_name == "ebuild":
- prefetcher = EbuildFetcher(logfile=self._fetch_log, pkg=pkg,
+ prefetcher = EbuildFetcher(fetchonly=1,
+ logfile=self._fetch_log, pkg=pkg,
scheduler=self._sched_iface)
elif pkg.type_name == "binary" and \
EPREFIX == BPREFIX and \
portage.match_from_list(
portage.const.PORTAGE_PACKAGE_ATOM, [pkg]):
+ if self._running_portage:
+ return cmp(pkg, self._running_portage) != 0
return True
return False
return
self._completed_tasks.add(pkg)
+ pkg_to_replace = merge.merge.pkg_to_replace
+ if pkg_to_replace is not None:
+ # When a package is replaced, mark it's uninstall
+ # task complete (if any).
+ uninst_hash_key = \
+ ("installed", pkg.root, pkg_to_replace.cpv, "uninstall")
+ self._completed_tasks.add(uninst_hash_key)
if pkg.installed:
return
if self._is_restart_scheduled():
self._set_max_jobs(1)
- pkg_queue = self._pkg_queue
- failed_pkgs = self._failed_pkgs
- task_queues = self._task_queues
- max_jobs = self._max_jobs
- max_load = self._max_load
- background = max_jobs > 1
+ while not self._failed_pkgs and \
+ self._schedule_tasks():
+ self._poll_loop()
- while pkg_queue and not failed_pkgs:
+ while self._jobs:
+ self._poll_loop()
- if self._jobs >= max_jobs:
- self._schedule_main()
- continue
+ def _schedule_tasks(self):
+ """
+ @rtype: bool
+ @returns: True if tasks remain to schedule, False otherwise.
+ """
- if max_load is not None and max_jobs > 1 and self._jobs > 1:
- try:
- avg1, avg5, avg15 = os.getloadavg()
- except OSError, e:
- writemsg("!!! getloadavg() failed: %s\n" % (e,),
- noiselevel=-1)
- del e
- self._schedule_main()
- continue
+ task_queues = self._task_queues
+ background = self._max_jobs > 1
- if avg1 >= max_load:
- self._schedule_main()
- continue
+ while self._can_add_job():
- pkg = self._choose_pkg()
+ if not self._pkg_queue:
+ return False
+ pkg = self._choose_pkg()
if pkg is None:
- self._schedule_main()
- continue
+ return True
if not pkg.installed:
self._pkg_count.curval += 1
else:
task.addExitListener(self._build_exit)
task_queues.jobs.add(task)
-
- while self._jobs:
- self._schedule_main(wait=True)
-
- def _schedule_main(self, wait=False):
-
- event_handlers = self._poll_event_handlers
- poll = self._poll.poll
- max_jobs = self._max_jobs
-
- state_change = 0
-
- if self._schedule_tasks():
- state_change += 1
-
- while event_handlers:
- jobs = self._jobs
-
- for f, event in poll():
- handler, reg_id = event_handlers[f]
- if not handler(f, event):
- state_change += 1
- self._unregister(reg_id)
-
- if jobs == self._jobs:
- continue
-
- if self._schedule_tasks():
- state_change += 1
-
- if not wait and self._jobs < max_jobs:
- break
-
- if not state_change:
- raise AssertionError("tight loop")
-
- def _schedule_tasks(self):
- state_change = 0
- for x in self._task_queues.values():
- if x.schedule():
- state_change += 1
- return bool(state_change)
+ return True
def _task(self, pkg, background):
+ pkg_to_replace = None
+ if pkg.operation != "uninstall":
+ vardb = pkg.root_config.trees["vartree"].dbapi
+ previous_cpv = vardb.match(pkg.slot_atom)
+ if previous_cpv:
+ previous_cpv = previous_cpv.pop()
+ pkg_to_replace = self._pkg(previous_cpv,
+ "installed", pkg.root_config, installed=True)
+
task = MergeListItem(args_set=self._args_set,
background=background, binpkg_opts=self._binpkg_opts,
build_opts=self._build_opts,
failed_fetches=self._failed_fetches,
find_blockers=self._find_blockers(pkg), logger=self._logger,
mtimedb=self._mtimedb, pkg=pkg, pkg_count=self._pkg_count.copy(),
+ pkg_to_replace=pkg_to_replace,
prefetcher=self._prefetchers.get(pkg),
scheduler=self._sched_iface,
settings=self._allocate_config(pkg.root),
return True
return False
- def _register(self, f, eventmask, handler):
- """
- @rtype: Integer
- @return: A unique registration id, for use in schedule() or
- unregister() calls.
- """
- self._event_handler_id += 1
- reg_id = self._event_handler_id
- self._poll_event_handler_ids[reg_id] = f
- self._poll_event_handlers[f] = (handler, reg_id)
- self._poll.register(f, eventmask)
- return reg_id
-
- def _unregister(self, reg_id):
- f = self._poll_event_handler_ids[reg_id]
- self._poll.unregister(f)
- del self._poll_event_handlers[f]
- del self._poll_event_handler_ids[reg_id]
- self._schedule_tasks()
-
- def _schedule(self, wait_id):
- """
- Schedule until wait_id is not longer registered
- for poll() events.
- @type wait_id: int
- @param wait_id: a task id to wait for
- """
- event_handlers = self._poll_event_handlers
- handler_ids = self._poll_event_handler_ids
- poll = self._poll.poll
-
- self._schedule_tasks()
-
- while wait_id in handler_ids:
- for f, event in poll():
- handler, reg_id = event_handlers[f]
- if not handler(f, event):
- self._unregister(reg_id)
-
def _world_atom(self, pkg):
"""
Add the package to the world file, but only if
finally:
world_set.unlock()
+ def _pkg(self, cpv, type_name, root_config, installed=False):
+ """
+ Get a package instance from the cache, or create a new
+ one if necessary. Raises KeyError from aux_get if it
+ failures for some reason (package does not exist or is
+ corrupt).
+ """
+ operation = "merge"
+ if installed:
+ operation = "nomerge"
+
+ db = root_config.trees[
+ depgraph.pkg_tree_map[type_name]].dbapi
+ metadata = izip(Package.metadata_keys,
+ db.aux_get(cpv, Package.metadata_keys))
+ pkg = Package(cpv=cpv, metadata=metadata,
+ root_config=root_config, installed=installed)
+ if type_name == "ebuild":
+ settings = self.pkgsettings[root_config.root]
+ settings.setcpv(pkg)
+ pkg.metadata["USE"] = settings["PORTAGE_USE"]
+
+ if self._digraph and \
+ self._digraph.contains(pkg):
+ for existing_instance in self._digraph.order:
+ if existing_instance == pkg:
+ pkg = existing_instance
+ break
+
+ return pkg
+
+class MetadataRegen(PollLoop):
+
+ class _sched_iface_class(SlotObject):
+ __slots__ = ("register", "schedule", "unregister")
+
+ def __init__(self, portdb, max_jobs=None, max_load=None):
+ PollLoop.__init__(self)
+ self._portdb = portdb
+
+ if max_jobs is None:
+ max_jobs = 1
+
+ self._max_jobs = max_jobs
+ self._max_load = max_load
+ self._sched_iface = self._sched_iface_class(
+ register=self._register,
+ schedule=self._schedule,
+ unregister=self._unregister)
+
+ self._valid_pkgs = set()
+ self._process_iter = self._iter_metadata_processes()
+
+ def _iter_metadata_processes(self):
+ portdb = self._portdb
+ valid_pkgs = self._valid_pkgs
+ every_cp = portdb.cp_all()
+ every_cp.sort(reverse=True)
+
+ while every_cp:
+ cp = every_cp.pop()
+ portage.writemsg_stdout("Processing %s\n" % cp)
+ cpv_list = portdb.cp_list(cp)
+ for cpv in cpv_list:
+ valid_pkgs.add(cpv)
+ ebuild_path, repo_path = portdb.findname2(cpv)
+ metadata_process = portdb._metadata_process(
+ cpv, ebuild_path, repo_path)
+ if metadata_process is None:
+ continue
+ yield metadata_process
+
+ def run(self):
+
+ portdb = self._portdb
+ from portage.cache.cache_errors import CacheError
+ dead_nodes = {}
+
+ for mytree in portdb.porttrees:
+ try:
+ dead_nodes[mytree] = set(portdb.auxdb[mytree].iterkeys())
+ except CacheError, e:
+ portage.writemsg("Error listing cache entries for " + \
+ "'%s': %s, continuing...\n" % (mytree, e), noiselevel=-1)
+ del e
+ dead_nodes = None
+ break
+
+ while self._schedule_tasks():
+ self._poll_loop()
+
+ while self._jobs:
+ self._poll_loop()
+
+ if dead_nodes:
+ for y in self._valid_pkgs:
+ for mytree in portdb.porttrees:
+ if portdb.findname2(y, mytree=mytree)[0]:
+ dead_nodes[mytree].discard(y)
+
+ for mytree, nodes in dead_nodes.iteritems():
+ auxdb = portdb.auxdb[mytree]
+ for y in nodes:
+ try:
+ del auxdb[y]
+ except (KeyError, CacheError):
+ pass
+
+ def _schedule_tasks(self):
+ """
+ @rtype: bool
+ @returns: True if there may be remaining tasks to schedule,
+ False otherwise.
+ """
+ while self._can_add_job():
+ try:
+ metadata_process = self._process_iter.next()
+ except StopIteration:
+ return False
+
+ self._jobs += 1
+ metadata_process.scheduler = self._sched_iface
+ metadata_process.addExitListener(self._metadata_exit)
+ metadata_process.start()
+ return True
+
+ def _metadata_exit(self, metadata_process):
+ self._jobs -= 1
+ if metadata_process.returncode != os.EX_OK:
+ self._valid_pkgs.discard(metadata_process.cpv)
+ portage.writemsg("Error processing %s, continuing...\n" % \
+ (metadata_process.cpv,))
+ self._schedule_tasks()
+
class UninstallFailure(portage.exception.PortageException):
"""
An instance of this class is raised by unmerge() when
sys.stdout.flush()
os.umask(old_umask)
-def action_regen(settings, portdb):
+def action_regen(settings, portdb, max_jobs, max_load):
xterm_titles = "notitles" not in settings.features
emergelog(xterm_titles, " === regen")
#regenerate cache entries
except:
pass
sys.stdout.flush()
- mynodes = portdb.cp_all()
- from portage.cache.cache_errors import CacheError
- dead_nodes = {}
- for mytree in portdb.porttrees:
- try:
- dead_nodes[mytree] = set(portdb.auxdb[mytree].iterkeys())
- except CacheError, e:
- portage.writemsg("Error listing cache entries for " + \
- "'%s': %s, continuing...\n" % (mytree, e), noiselevel=-1)
- del e
- dead_nodes = None
- break
- for x in mynodes:
- mymatches = portdb.cp_list(x)
- portage.writemsg_stdout("Processing %s\n" % x)
- for y in mymatches:
- try:
- foo = portdb.aux_get(y,["DEPEND"])
- except (KeyError, portage.exception.PortageException), e:
- portage.writemsg(
- "Error processing %(cpv)s, continuing... (%(e)s)\n" % \
- {"cpv":y,"e":str(e)}, noiselevel=-1)
- if dead_nodes:
- for mytree in portdb.porttrees:
- if portdb.findname2(y, mytree=mytree)[0]:
- dead_nodes[mytree].discard(y)
- if dead_nodes:
- for mytree, nodes in dead_nodes.iteritems():
- auxdb = portdb.auxdb[mytree]
- for y in nodes:
- try:
- del auxdb[y]
- except (KeyError, CacheError):
- pass
+
+ regen = MetadataRegen(portdb, max_jobs=max_jobs, max_load=max_load)
+ regen.run()
+
portage.writemsg_stdout("done!\n")
def action_config(settings, trees, myopts, myfiles):
def emerge_main():
global portage # NFC why this is necessary now - genone
+ portage._disable_legacy_globals()
# Disable color until we're sure that it should be enabled (after
# EMERGE_DEFAULT_OPTS has been parsed).
portage.output.havecolor = 0
action_metadata(settings, portdb, myopts)
elif myaction=="regen":
validate_ebuild_environment(trees)
- action_regen(settings, portdb)
+ action_regen(settings, portdb, myopts.get("--jobs"),
+ myopts.get("--load-average"))
# HELP action
elif "config"==myaction:
validate_ebuild_environment(trees)