From 4d4b89ed3852c88a8642c6452a6a8e1addd49d95 Mon Sep 17 00:00:00 2001 From: Fabian Groffen Date: Sun, 13 Jul 2008 11:28:59 +0000 Subject: [PATCH] Merged from trunk 11017:11035 | 11018 | Fix a few typos found with pyflakes. Thanks to Cardoe for | | zmedico | reporting. | | 11019 | Fix SpawnProcess.start() to return early and notify exit | | zmedico | listeners if the spawn call fails. | | 11020 | For subclasses of Task, explicitly assign __hash__ = | | zmedico | Task.__hash__ for compatibility with python-2.6_beta. Thanks | | | to Arfrever for reporting. | | 11021 | * Add a digraph.get() method, analogous to dict.get(), used | | zmedico | to simplify Scheduler._pkg(). * Define digraph.__contains__ | | | for containment tests via the "in" operator. | | 11022 | Create a test case for the poll loop which uses the loop to | | zmedico | read data from a pipe and assert that the data written to | | | the pipe is identical to the data read from the pipe. In | | | order to implement this test, several useful classes have | | | been added: * PipeReader Reads output from one or more files | | | and saves it in memory, for retrieval via the getvalue() | | | method. This is driven by the scheduler's poll() loop, so it | | | runs entirely within the current process. * QueueScheduler | | | Add instances of SequentialTaskQueue and then call run(). | | | The run() method returns when no tasks remain. * | | | TaskScheduler A simple way to handle scheduling of | | | AsynchrousTask instances. Simply add tasks and call run(). | | | The run() method returns when no tasks remain. | | 11023 | Derive a PtyReaderTestCase from PipeReaderTestCase. This | | zmedico | test is expected to fail on some operating systems such as | | | Darwin that do not support poll() on pty devices. | | 11024 | Split out a portage._create_pty_or_pipe() function to avoid | | zmedico | duplicate code. | | 11025 | Implement a runtime test that checks whether it's possible | | zmedico | to use poll() to read data from a pty device. The test only | | | runs the first time that the can_poll_pty() function is | | | called (currently happens when the first PollScheduler | | | instance is created). The test is known to fail on Darwin | | | and we're not sure how many other kernels might have a | | | similar problem, this runtime test allows us to avoid having | | | to maintain special cases for various kernels. | | 11026 | Set self._mergelist inside Scheduler._calc_resume_list() | | zmedico | since it needs to be set when calling _set_digraph(). | | 11027 | Add missing self.wait() calls when tasks return | | zmedico | unsuccessfully. | | 11028 | Remember to close file descriptors before returning early | | zmedico | from can_poll_pty(). | | 11029 | Rename PollScheduler._schedule() to _schedule_wait(). | | zmedico | | | 11030 | * Wrap calls to PollScheduler._schedule_tasks() with a | | zmedico | _schedule() method that automatically returns early from any | | | recursive calls that the _schedule_tasks() call might | | | trigger. This makes _schedule() safe to call from inside | | | exit listeners. * Call _schedule() inside | | | Scheduler._job_exit() to trigger a new job when possible. * | | | Fix SequentialTaskQueue._task_exit() to avoid a potential | | | KeyError if the task had already been removed from | | | self.running_tasks for some reason. | | 11031 | Decrement the job count prior to merging a given package | | zmedico | (after it's been built or extracted). This will allow | | | additional jobs to spawn sooner in some cases when there may | | | be idle resources. | | 11032 | Fix Scheduler._schedule_tasks() so that it never schedules a | | zmedico | task when a package has already failed. | | 11033 | In order to avoid having BlockerDB instantiate a new | | zmedico | FakeVartree for each merge, add a FakeVartree.sync() method | | | and so that the same FakeVartree can be reused. Package | | | counters and timestamps are used to validate package | | | instances during each sync() call. | | 11034 | Fix broken binpkg path comparison in in binarytree.inject() | | zmedico | so that duplicate entries are correctly from the Packages | | | index. | | 11035 | Changes to pym/_emerge/__init__.py from the previous commit | | zmedico | weren't meant to be committed yet. | svn path=/main/branches/prefix/; revision=11036 --- pym/_emerge/__init__.py | 585 ++++++++++++++++++++----- pym/portage/__init__.py | 90 ++-- pym/portage/dbapi/bintree.py | 2 +- pym/portage/tests/process/__init__.py | 3 + pym/portage/tests/process/__test__ | 0 pym/portage/tests/process/test_poll.py | 69 +++ pym/repoman/utilities.py | 8 +- 7 files changed, 617 insertions(+), 140 deletions(-) create mode 100644 pym/portage/tests/process/__init__.py create mode 100644 pym/portage/tests/process/__test__ create mode 100644 pym/portage/tests/process/test_poll.py diff --git a/pym/_emerge/__init__.py b/pym/_emerge/__init__.py index 9d8adfbd9..8cc3cf817 100644 --- a/pym/_emerge/__init__.py +++ b/pym/_emerge/__init__.py @@ -1029,6 +1029,7 @@ class FakeVartree(portage.vartree): global updates are necessary (updates are performed when necessary if there is not a matching ebuild in the tree).""" def __init__(self, root_config, pkg_cache=None, acquire_lock=1): + self._root_config = root_config if pkg_cache is None: pkg_cache = {} real_vartree = root_config.trees["vartree"] @@ -1131,6 +1132,80 @@ class FakeVartree(portage.vartree): pkg, self.dbapi, self._global_updates) return self._aux_get(pkg, wants) + def sync(self, acquire_lock=1): + """ + Call this method to synchronize state with the real vardb + after one or more packages may have been installed or + uninstalled. + """ + vdb_path = os.path.join(self.root, portage.VDB_PATH) + try: + # At least the parent needs to exist for the lock file. + portage.util.ensure_dirs(vdb_path) + except portage.exception.PortageException: + pass + vdb_lock = None + try: + if acquire_lock and os.access(vdb_path, os.W_OK): + vdb_lock = portage.locks.lockdir(vdb_path) + self._sync() + finally: + if vdb_lock: + portage.locks.unlockdir(vdb_lock) + + def _sync(self): + + real_vardb = self._root_config.trees["vartree"].dbapi + current_cpv_set = frozenset(real_vardb.cpv_all()) + pkg_vardb = self.dbapi + aux_get_history = self._aux_get_history + + # Remove any packages that have been uninstalled. + for pkg in list(pkg_vardb): + if pkg.cpv not in current_cpv_set: + pkg_vardb.cpv_remove(pkg) + aux_get_history.discard(pkg.cpv) + + # Validate counters and timestamps. + slot_counters = {} + root = self.root + validation_keys = ["COUNTER", "_mtime_"] + for cpv in current_cpv_set: + + pkg_hash_key = ("installed", root, cpv, "nomerge") + pkg = pkg_vardb.get(pkg_hash_key) + if pkg is not None: + counter, mtime = real_vardb.aux_get(cpv, validation_keys) + + if counter != pkg.metadata["COUNTER"] or \ + mtime != pkg.mtime: + pkg_vardb.cpv_remove(pkg) + aux_get_history.discard(pkg.cpv) + pkg = None + + if pkg is None: + pkg = self._pkg(cpv) + + other_counter = slot_counters.get(pkg.slot_atom) + if other_counter is not None: + if other_counter > pkg.counter: + continue + + slot_counters[pkg.slot_atom] = pkg.counter + pkg_vardb.cpv_inject(pkg) + + real_vardb.flush_cache() + + def _pkg(self, cpv): + root_config = self._root_config + real_vardb = root_config.trees["vartree"].dbapi + pkg = Package(cpv=cpv, installed=True, + metadata=izip(Package.metadata_keys, + real_vardb.aux_get(cpv, Package.metadata_keys)), + root_config=root_config, + type_name="installed") + return pkg + def grab_global_updates(portdir): from portage.update import grab_updates, parse_updates updpath = os.path.join(portdir, "profiles", "updates") @@ -1315,6 +1390,8 @@ class Task(SlotObject): return str(self._get_hash_key()) class Blocker(Task): + + __hash__ = Task.__hash__ __slots__ = ("root", "atom", "cp", "satisfied") def __init__(self, **kwargs): @@ -1329,6 +1406,8 @@ class Blocker(Task): return self._hash_key class Package(Task): + + __hash__ = Task.__hash__ __slots__ = ("built", "cpv", "depth", "installed", "metadata", "onlydeps", "operation", "root_config", "type_name", @@ -1521,13 +1600,17 @@ class AsynchronousTask(SlotObject): """ __slots__ = ("background", "cancelled", "returncode") + \ - ("_exit_listeners",) + ("_exit_listeners", "_start_listeners") def start(self): """ Start an asynchronous task and then return as soon as possible. """ - pass + self._start() + self._start_hook() + + def _start(self): + raise NotImplementedError(self) def isAlive(self): return self.returncode is None @@ -1552,6 +1635,25 @@ class AsynchronousTask(SlotObject): self.cancelled = True self.wait() + def addStartListener(self, f): + """ + The function will be called with one argument, a reference to self. + """ + if self._start_listeners is None: + self._start_listeners = [] + self._start_listeners.append(f) + + def removeStartListener(self, f): + self._start_listeners.remove(f) + + def _start_hook(self): + if self._start_listeners is not None: + start_listeners = self._start_listeners + self._start_listeners = None + + for f in start_listeners: + f(self) + def addExitListener(self, f): """ The function will be called with one argument, a reference to self. @@ -1581,7 +1683,73 @@ class AsynchronousTask(SlotObject): for f in exit_listeners: f(self) - self._exit_listeners = None + +class PipeReader(AsynchronousTask): + + """ + Reads output from one or more files and saves it in memory, + for retrieval via the getvalue() method. This is driven by + the scheduler's poll() loop, so it runs entirely within the + current process. + """ + + __slots__ = ("input_files", "scheduler",) + \ + ("pid", "registered", "_reg_ids", "_read_data") + + _bufsize = 4096 + + def _start(self): + self._reg_ids = set() + self._read_data = [] + for k, f in self.input_files.iteritems(): + fcntl.fcntl(f.fileno(), fcntl.F_SETFL, + fcntl.fcntl(f.fileno(), fcntl.F_GETFL) | os.O_NONBLOCK) + self._reg_ids.add(self.scheduler.register(f.fileno(), + PollConstants.POLLIN, self._output_handler)) + self.registered = True + + def isAlive(self): + return self.registered + + def _wait(self): + if self.returncode is not None: + return self.returncode + if self.registered: + self.scheduler.schedule(self._reg_ids) + self.returncode = os.EX_OK + return self.returncode + + def getvalue(self): + """Retrieve the entire contents""" + return "".join(self._read_data) + + def close(self): + """Free the memory buffer.""" + self._read_data = None + + def _output_handler(self, fd, event): + files = self.input_files + for f in files.itervalues(): + if fd == f.fileno(): + break + + buf = array.array('B') + try: + buf.fromfile(f, self._bufsize) + except EOFError: + pass + + if buf: + self._read_data.append(buf.tostring()) + else: + for f in files.values(): + f.close() + self.registered = False + for reg_id in self._reg_ids: + self.scheduler.unregister(reg_id) + self.wait() + + return self.registered class CompositeTask(AsynchronousTask): @@ -1712,7 +1880,7 @@ class TaskSequence(CompositeTask): def add(self, task): self._task_queue.append(task) - def start(self): + def _start(self): self._start_next_task() def cancel(self): @@ -1725,7 +1893,7 @@ class TaskSequence(CompositeTask): def _task_exit_handler(self, task): if self._default_exit(task) != os.EX_OK: - pass + self.wait() elif self._task_queue: self._start_next_task() else: @@ -1825,7 +1993,7 @@ class SpawnProcess(SubProcess): _files_dict = slot_dict_class(_file_names, prefix="") _bufsize = 4096 - def start(self): + def _start(self): if self.cancelled: return @@ -1888,10 +2056,19 @@ class SpawnProcess(SubProcess): retval = portage.process.spawn(self.args, **kwargs) + os.close(slave_fd) + + if isinstance(retval, int): + # spawn failed + os.close(master_fd) + self.returncode = retval + self.wait() + return + self.pid = retval[0] portage.process.spawned_pids.remove(self.pid) - os.close(slave_fd) + files.process = os.fdopen(master_fd, 'r') self._reg_id = self.scheduler.register(files.process.fileno(), PollConstants.POLLIN, output_handler) @@ -1944,7 +2121,7 @@ class EbuildFetcher(SpawnProcess): __slots__ = ("fetchonly", "pkg",) - def start(self): + def _start(self): root_config = self.pkg.root_config portdb = root_config.trees["porttree"].dbapi @@ -1966,7 +2143,7 @@ class EbuildFetcher(SpawnProcess): self.args = fetch_args self.env = fetch_env - SpawnProcess.start(self) + SpawnProcess._start(self) class EbuildBuildDir(SlotObject): @@ -2050,7 +2227,7 @@ class EbuildBuild(CompositeTask): "prefetcher", "settings", "world_atom") + \ ("_build_dir", "_buildpkg", "_ebuild_path", "_tree") - def start(self): + def _start(self): logger = self.logger opts = self.opts @@ -2194,6 +2371,7 @@ class EbuildBuild(CompositeTask): def _build_exit(self, build): if self._default_exit(build) != os.EX_OK: self._unlock_builddir() + self.wait() return opts = self.opts @@ -2280,7 +2458,7 @@ class EbuildExecuter(CompositeTask): _phases = ("setup", "unpack", "compile", "test", "install") - def start(self): + def _start(self): pkg = self.pkg scheduler = self.scheduler tree = "porttree" @@ -2294,6 +2472,7 @@ class EbuildExecuter(CompositeTask): def _clean_phase_exit(self, clean_phase): if self._default_exit(clean_phase) != os.EX_OK: + self.wait() return pkg = self.pkg @@ -2337,7 +2516,7 @@ class EbuildMetadataPhase(SubProcess): _bufsize = SpawnProcess._bufsize _metadata_fd = 9 - def start(self): + def _start(self): settings = self.settings settings.reset() ebuild_path = self.ebuild_path @@ -2422,7 +2601,7 @@ class EbuildPhase(SubProcess): _files_dict = slot_dict_class(_file_names, prefix="") _bufsize = 4096 - def start(self): + def _start(self): root_config = self.pkg.root_config tree = self.tree mydbapi = root_config.trees[tree].dbapi @@ -2458,35 +2637,12 @@ class EbuildPhase(SubProcess): settings.get("EBUILD_EXIT_STATUS_FILE")) if logfile: - if portage._disable_openpty: - master_fd, slave_fd = os.pipe() - else: - from pty import openpty - try: - master_fd, slave_fd = openpty() - got_pty = True - except EnvironmentError, e: - portage._disable_openpty = True - portage.writemsg("openpty failed: '%s'\n" % str(e), - noiselevel=-1) - del e - master_fd, slave_fd = os.pipe() - - if got_pty: - # Disable post-processing of output since otherwise weird - # things like \n -> \r\n transformations may occur. - import termios - mode = termios.tcgetattr(slave_fd) - mode[1] &= ~termios.OPOST - termios.tcsetattr(slave_fd, termios.TCSANOW, mode) + got_pty, master_fd, slave_fd = \ + portage._create_pty_or_pipe(copy_term_size=fd_pipes_orig[1]) fcntl.fcntl(master_fd, fcntl.F_SETFL, fcntl.fcntl(master_fd, fcntl.F_GETFL) | os.O_NONBLOCK) - if got_pty and os.isatty(fd_pipes_orig[1]): - from portage.output import get_term_size, set_term_size - rows, columns = get_term_size() - set_term_size(rows, columns, slave_fd) fd_pipes[0] = fd_pipes_orig[0] fd_pipes[1] = slave_fd fd_pipes[2] = slave_fd @@ -2602,7 +2758,7 @@ class EbuildBinpkg(EbuildPhase): """ __slots__ = ("_binpkg_tmpfile",) - def start(self): + def _start(self): self.phase = "package" self.tree = "porttree" pkg = self.pkg @@ -2621,7 +2777,7 @@ class EbuildBinpkg(EbuildPhase): settings.backup_changes("PORTAGE_BINPKG_TMPFILE") try: - EbuildPhase.start(self) + EbuildPhase._start(self) finally: settings.pop("PORTAGE_BINPKG_TMPFILE", None) @@ -2676,6 +2832,7 @@ class EbuildMerge(SlotObject): class PackageUninstall(Task): + __hash__ = Task.__hash__ __slots__ = ("ldpath_mtimes", "opts", "pkg", "settings") def _get_hash_key(self): @@ -2701,7 +2858,7 @@ class Binpkg(CompositeTask): ("_bintree", "_build_dir", "_ebuild_path", "_fetched_pkg", "_image_dir", "_infloc", "_pkg_path", "_tree", "_verify") - def start(self): + def _start(self): pkg = self.pkg settings = self.settings @@ -2792,6 +2949,7 @@ class Binpkg(CompositeTask): self.wait() return elif self._default_exit(fetcher) != os.EX_OK: + self.wait() return verifier = None @@ -2805,7 +2963,8 @@ class Binpkg(CompositeTask): def _verifier_exit(self, verifier): if verifier is not None and \ self._default_exit(verifier) != os.EX_OK: - return + self.wait() + return logger = self.logger pkg = self.pkg @@ -2836,6 +2995,7 @@ class Binpkg(CompositeTask): def _clean_exit(self, clean_phase): if self._default_exit(clean_phase) != os.EX_OK: self._unlock_builddir() + self.wait() return dir_path = self._build_dir.dir_path @@ -2908,6 +3068,7 @@ class Binpkg(CompositeTask): def _setup_exit(self, setup_phase): if self._default_exit(setup_phase) != os.EX_OK: self._unlock_builddir() + self.wait() return extractor = BinpkgExtractorAsync(background=self.background, @@ -2958,7 +3119,7 @@ class BinpkgFetcher(SpawnProcess): pkg = self.pkg self.pkg_path = pkg.root_config.trees["bintree"].getname(pkg.cpv) - def start(self): + def _start(self): if self.cancelled: return @@ -3017,7 +3178,7 @@ class BinpkgFetcher(SpawnProcess): self.args = fetch_args self.env = fetch_env - SpawnProcess.start(self) + SpawnProcess._start(self) def _set_returncode(self, wait_retval): SpawnProcess._set_returncode(self, wait_retval) @@ -3051,7 +3212,7 @@ class BinpkgFetcher(SpawnProcess): class BinpkgVerifier(AsynchronousTask): __slots__ = ("pkg",) - def start(self): + def _start(self): """ Note: Unlike a normal AsynchronousTask.start() method, this one does all work is synchronously. The returncode @@ -3090,14 +3251,14 @@ class BinpkgExtractorAsync(SpawnProcess): _shell_binary = portage.const.BASH_BINARY - def start(self): + def _start(self): self.args = [self._shell_binary, "-c", "bzip2 -dqc -- %s | tar -xp -C %s -f -" % \ (portage._shell_quote(self.pkg_path), portage._shell_quote(self.image_dir))] self.env = self.pkg.root_config.settings.environ() - SpawnProcess.start(self) + SpawnProcess._start(self) class MergeListItem(CompositeTask): @@ -3113,7 +3274,7 @@ class MergeListItem(CompositeTask): "settings", "world_atom") + \ ("_install_task",) - def start(self): + def _start(self): pkg = self.pkg build_opts = self.build_opts @@ -3226,7 +3387,7 @@ class PackageMerge(AsynchronousTask): __slots__ = ("merge",) - def start(self): + def _start(self): self.returncode = self.merge.merge() self.wait() @@ -3435,20 +3596,31 @@ class BlockerDB(object): self._root_config = root_config self._vartree = root_config.trees["vartree"] self._portdb = root_config.trees["porttree"].dbapi - - self._dep_check_trees = { self._vartree.root : { - "porttree" : self._vartree, - "vartree" : self._vartree, - }} + + self._dep_check_trees = None + self._fake_vartree = None + + def _get_fake_vartree(self, acquire_lock=0): + fake_vartree = self._fake_vartree + if fake_vartree is None: + fake_vartree = FakeVartree(self._root_config, + acquire_lock=acquire_lock) + self._fake_vartree = fake_vartree + self._dep_check_trees = { self._vartree.root : { + "porttree" : fake_vartree, + "vartree" : fake_vartree, + }} + else: + fake_vartree.sync(acquire_lock=acquire_lock) + return fake_vartree def findInstalledBlockers(self, new_pkg, acquire_lock=0): blocker_cache = BlockerCache(self._vartree.root, self._vartree.dbapi) dep_keys = ["DEPEND", "RDEPEND", "PDEPEND"] - dep_check_trees = self._dep_check_trees settings = self._vartree.settings stale_cache = set(blocker_cache) - fake_vartree = \ - FakeVartree(self._root_config, acquire_lock=acquire_lock) + fake_vartree = self._get_fake_vartree(acquire_lock=acquire_lock) + dep_check_trees = self._dep_check_trees vardb = fake_vartree.dbapi installed_pkgs = list(vardb) @@ -3606,6 +3778,19 @@ class PackageVirtualDbapi(portage.dbapi): return True return False + def get(self, item, default=None): + cpv = getattr(item, "cpv", None) + if cpv is None: + if len(item) != 4: + return default + type_name, root, cpv, operation = item + + existing = self._cpv_map.get(cpv) + if existing is not None and \ + existing == item: + return existing + return default + def match_pkgs(self, atom): return [self._cpv_map[cpv] for cpv in self.match(atom)] @@ -7756,7 +7941,8 @@ class PollSelectAdapter(PollConstants): class SequentialTaskQueue(SlotObject): - __slots__ = ("max_jobs", "running_tasks", "_task_queue", "_scheduling") + __slots__ = ("auto_schedule", "max_jobs", "running_tasks") + \ + ("_task_queue", "_scheduling") def __init__(self, **kwargs): SlotObject.__init__(self, **kwargs) @@ -7767,11 +7953,13 @@ class SequentialTaskQueue(SlotObject): def add(self, task): self._task_queue.append(task) - self.schedule() + if self.auto_schedule: + self.schedule() def addFront(self, task): self._task_queue.appendleft(task) - self.schedule() + if self.auto_schedule: + self.schedule() def schedule(self): @@ -7794,16 +7982,15 @@ class SequentialTaskQueue(SlotObject): if hasattr(task, "registered") and task.registered: continue 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: + running_tasks.add(task) task.addExitListener(self._task_exit) task.start() - running_tasks.add(task) state_changed = True self._scheduling = False @@ -7811,7 +7998,9 @@ class SequentialTaskQueue(SlotObject): return state_changed def _task_exit(self, task): - self.schedule() + self.running_tasks.discard(task) + if self.auto_schedule: + self.schedule() def clear(self): self._task_queue.clear() @@ -7826,9 +8015,74 @@ class SequentialTaskQueue(SlotObject): def __len__(self): return len(self._task_queue) + len(self.running_tasks) -class PollLoop(object): +_can_poll_pty = None - def __init__(self): +def can_poll_pty(): + """ + Test if it's possible to use poll() on a pty device. This + is known to fail on Darwin. + @rtype: bool + @returns: True if poll() on a pty device succeeds, False otherwise. + """ + + global _can_poll_pty + if _can_poll_pty is not None: + return _can_poll_pty + + if not hasattr(select, "poll"): + _can_poll_pty = False + return _can_poll_pty + + got_pty, master_fd, slave_fd = \ + portage._create_pty_or_pipe(copy_term_size=sys.stdout.fileno()) + if not got_pty: + os.close(master_fd) + os.close(slave_fd) + _can_poll_pty = False + return _can_poll_pty + + test_string = 2 * "blah blah blah\n" + + master_file = os.fdopen(master_fd, 'r') + + task_scheduler = TaskScheduler(max_jobs=2, poll=select.poll()) + scheduler = task_scheduler.sched_iface + + producer = SpawnProcess( + args=["bash", "-c", "echo -n '%s'" % test_string], + fd_pipes={1:slave_fd}, scheduler=scheduler) + + consumer = PipeReader( + input_files={"producer" : master_file}, + scheduler=scheduler) + + task_scheduler.add(producer) + task_scheduler.add(consumer) + + def producer_start_cb(task): + os.close(slave_fd) + + producer.addStartListener(producer_start_cb) + task_scheduler.run() + _can_poll_pty = test_string == consumer.getvalue() + return _can_poll_pty + +def create_poll_instance(): + """ + Create an instance of select.poll, or an instance of + PollSelectAdapter there is no poll() implementation or + it is broken somehow. + """ + if can_poll_pty(): + return select.poll() + return PollSelectAdapter() + +class PollScheduler(object): + + class _sched_iface_class(SlotObject): + __slots__ = ("register", "schedule", "unregister") + + def __init__(self, poll=None): self._max_jobs = 1 self._max_load = None self._jobs = 0 @@ -7836,20 +8090,38 @@ class PollLoop(object): self._poll_event_handler_ids = {} # Increment id for each new handler. self._event_handler_id = 0 + if poll is None: + poll = create_poll_instance() + self._poll = poll + self._scheduling = False + + def _schedule(self): + """ + Calls _schedule_tasks() and automatically returns early from + any recursive calls to this method that the _schedule_tasks() + call might trigger. This makes _schedule() safe to call from + inside exit listeners. + """ + if self._scheduling: + return False + self._scheduling = True try: - self._poll = select.poll() - except AttributeError: - self._poll = PollSelectAdapter() + return self._schedule_tasks() + finally: + self._scheduling = False + + def _running_job_count(self): + return self._jobs def _can_add_job(self): - jobs = self._jobs max_jobs = self._max_jobs max_load = self._max_load - if self._jobs >= self._max_jobs: + if self._running_job_count() >= self._max_jobs: return False - if max_load is not None and max_jobs > 1 and self._jobs > 1: + if max_load is not None and max_jobs > 1 and \ + self._running_job_count() > 1: try: avg1, avg5, avg15 = os.getloadavg() except OSError, e: @@ -7899,7 +8171,7 @@ class PollLoop(object): del self._poll_event_handlers[f] del self._poll_event_handler_ids[reg_id] - def _schedule(self, wait_id): + def _schedule_wait(self, wait_ids): """ Schedule until wait_id is not longer registered for poll() events. @@ -7910,12 +8182,111 @@ class PollLoop(object): handler_ids = self._poll_event_handler_ids poll = self._poll.poll - while wait_id in handler_ids: + if isinstance(wait_ids, int): + wait_ids = frozenset([wait_ids]) + + while wait_ids.intersection(handler_ids): for f, event in poll(): handler, reg_id = event_handlers[f] handler(f, event) -class Scheduler(PollLoop): +class QueueScheduler(PollScheduler): + + """ + Add instances of SequentialTaskQueue and then call run(). The + run() method returns when no tasks remain. + """ + + def __init__(self, max_jobs=None, max_load=None, poll=None): + PollScheduler.__init__(self, poll=poll) + + 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_wait, + unregister=self._unregister) + + self._queues = [] + self._schedule_listeners = [] + + def add(self, q): + self._queues.append(q) + + def remove(self, q): + self._queues.remove(q) + + def run(self): + + while self._schedule(): + self._poll_loop() + + while self._running_job_count(): + self._poll_loop() + + def _schedule_tasks(self): + """ + @rtype: bool + @returns: True if there may be remaining tasks to schedule, + False otherwise. + """ + while self._can_add_job(): + n = self._max_jobs - self._running_job_count() + if n < 1: + break + + if not self._start_next_job(n): + return False + + for q in self._queues: + if q: + return True + return False + + def _running_job_count(self): + job_count = 0 + for q in self._queues: + job_count += len(q.running_tasks) + self._jobs = job_count + return job_count + + def _start_next_job(self, n=1): + started_count = 0 + for q in self._queues: + initial_job_count = len(q.running_tasks) + q.schedule() + final_job_count = len(q.running_tasks) + if final_job_count > initial_job_count: + started_count += (final_job_count - initial_job_count) + if started_count >= n: + break + return started_count + +class TaskScheduler(object): + + """ + A simple way to handle scheduling of AsynchrousTask instances. Simply + add tasks and call run(). The run() method returns when no tasks remain. + """ + + def __init__(self, max_jobs=None, max_load=None, poll=None): + self._queue = SequentialTaskQueue(max_jobs=max_jobs) + self._scheduler = QueueScheduler( + max_jobs=max_jobs, max_load=max_load, poll=poll) + self.sched_iface = self._scheduler.sched_iface + self.run = self._scheduler.run + self._scheduler.add(self._queue) + + def add(self, task): + self._queue.add(task) + + def run(self): + self._scheduler.schedule() + +class Scheduler(PollScheduler): _opts_ignore_blockers = \ frozenset(["--buildpkgonly", @@ -7957,7 +8328,7 @@ class Scheduler(PollLoop): def __init__(self, settings, trees, mtimedb, myopts, spinner, mergelist, favorites, digraph): - PollLoop.__init__(self) + PollScheduler.__init__(self) self.settings = settings self.target_root = settings["ROOT"] self.trees = trees @@ -7982,10 +8353,12 @@ class Scheduler(PollLoop): self.edebug = 1 self.pkgsettings = {} self._config_pool = {} + self._blocker_db = {} for root in trees: self.pkgsettings[root] = portage.config( clone=trees[root]["vartree"].settings) self._config_pool[root] = [] + self._blocker_db[root] = BlockerDB(trees[root]["root_config"]) self.curval = 0 self._logger = self._emerge_log_class( xterm_titles=("notitles" not in settings.features)) @@ -7993,11 +8366,12 @@ class Scheduler(PollLoop): schedule=self._schedule_fetch) self._sched_iface = self._iface_class( fetch=fetch_iface, register=self._register, - schedule=self._schedule, unregister=self._unregister) + schedule=self._schedule_wait, unregister=self._unregister) self._task_queues = self._task_queues_class() for k in self._task_queues.allowed_keys: - setattr(self._task_queues, k, SequentialTaskQueue()) + setattr(self._task_queues, k, + SequentialTaskQueue(auto_schedule=True)) self._prefetchers = weakref.WeakValueDictionary() self._pkg_queue = [] @@ -8146,7 +8520,7 @@ class Scheduler(PollLoop): import gc gc.collect() - blocker_db = BlockerDB(self.trees[new_pkg.root]["root_config"]) + blocker_db = self._blocker_db[new_pkg.root] blocker_dblinks = [] for blocking_pkg in blocker_db.findInstalledBlockers( @@ -8384,9 +8758,9 @@ class Scheduler(PollLoop): if not mergelist: break - mylist, dropped_tasks = self._calc_resume_list() + dropped_tasks = self._calc_resume_list() clear_caches(self.trees) - if not mylist: + if not self._mergelist: break if dropped_tasks: @@ -8404,10 +8778,9 @@ class Scheduler(PollLoop): _eerror(msg) del _eerror, msg del dropped_tasks - self._mergelist = mylist self._save_resume_list() self._pkg_count.curval = 0 - self._pkg_count.maxval = len([x for x in mylist \ + self._pkg_count.maxval = len([x for x in self._mergelist \ if isinstance(x, Package) and x.operation == "merge"]) self._logger.log(" *** Finished. Cleaning up...") @@ -8423,7 +8796,11 @@ class Scheduler(PollLoop): pass def _merge_exit(self, merge): - self._job_exit(merge.merge) + self._do_merge_exit(merge) + self._deallocate_config(merge.merge.settings) + self._schedule() + + def _do_merge_exit(self, merge): pkg = merge.merge.pkg if merge.returncode != os.EX_OK: self._failed_pkgs.append((pkg, merge.returncode)) @@ -8458,18 +8835,15 @@ class Scheduler(PollLoop): merge = PackageMerge(merge=build) merge.addExitListener(self._merge_exit) self._task_queues.merge.add(merge) - self._task_queues.merge.schedule() else: self._failed_pkgs.append((build.pkg, build.returncode)) - self._job_exit(build) + self._deallocate_config(build.settings) + self._jobs -= 1 + self._schedule() def _extract_exit(self, build): self._build_exit(build) - def _job_exit(self, job): - self._jobs -= 1 - self._deallocate_config(job.settings) - def _merge(self): self._add_prefetchers() @@ -8564,7 +8938,7 @@ class Scheduler(PollLoop): self._set_max_jobs(1) while not self._failed_pkgs and \ - self._schedule_tasks(): + self._schedule(): self._poll_loop() while self._jobs: @@ -8581,7 +8955,7 @@ class Scheduler(PollLoop): while self._can_add_job(): - if not self._pkg_queue: + if not self._pkg_queue or self._failed_pkgs: return False pkg = self._choose_pkg() @@ -8593,15 +8967,16 @@ class Scheduler(PollLoop): task = self._task(pkg, background) - self._jobs += 1 if pkg.installed: merge = PackageMerge(merge=task) merge.addExitListener(self._merge_exit) task_queues.merge.add(merge) elif pkg.built: + self._jobs += 1 task.addExitListener(self._extract_exit) task_queues.jobs.add(task) else: + self._jobs += 1 task.addExitListener(self._build_exit) task_queues.jobs.add(task) return True @@ -8691,8 +9066,10 @@ class Scheduler(PollLoop): mydepgraph.break_refs(mylist) mydepgraph.break_refs(dropped_tasks) mydepgraph.break_refs(mydepgraph.digraph.order) + + self._mergelist = mylist self._set_digraph(mydepgraph.digraph) - return (mylist, dropped_tasks) + return dropped_tasks def _show_list(self): myopts = self.myopts @@ -8760,22 +9137,16 @@ class Scheduler(PollLoop): 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 + if self._digraph is not None: + # Reuse existing instance when available. + pkg = self._digraph.get(pkg, pkg) return pkg -class MetadataRegen(PollLoop): - - class _sched_iface_class(SlotObject): - __slots__ = ("register", "schedule", "unregister") +class MetadataRegen(PollScheduler): def __init__(self, portdb, max_jobs=None, max_load=None): - PollLoop.__init__(self) + PollScheduler.__init__(self) self._portdb = portdb if max_jobs is None: @@ -8785,7 +9156,7 @@ class MetadataRegen(PollLoop): self._max_load = max_load self._sched_iface = self._sched_iface_class( register=self._register, - schedule=self._schedule, + schedule=self._schedule_wait, unregister=self._unregister) self._valid_pkgs = set() @@ -8826,7 +9197,7 @@ class MetadataRegen(PollLoop): dead_nodes = None break - while self._schedule_tasks(): + while self._schedule(): self._poll_loop() while self._jobs: @@ -8870,7 +9241,7 @@ class MetadataRegen(PollLoop): self._valid_pkgs.discard(metadata_process.cpv) portage.writemsg("Error processing %s, continuing...\n" % \ (metadata_process.cpv,)) - self._schedule_tasks() + self._schedule() class UninstallFailure(portage.exception.PortageException): """ diff --git a/pym/portage/__init__.py b/pym/portage/__init__.py index cc2d25de6..18af5af4c 100644 --- a/pym/portage/__init__.py +++ b/pym/portage/__init__.py @@ -438,6 +438,9 @@ class digraph(object): """Checks if the digraph contains mynode""" return node in self.nodes + def get(self, key, default=None): + return self.nodes.get(key, default) + def all_nodes(self): """Return a list of all nodes in the graph""" return self.order[:] @@ -507,6 +510,7 @@ class digraph(object): allnodes = all_nodes allzeros = leaf_nodes hasnode = contains + __contains__ = contains empty = is_empty copy = clone @@ -2918,6 +2922,54 @@ def _shell_quote(s): # stop trying to use it after the first failure. _disable_openpty = False +def _create_pty_or_pipe(copy_term_size=None): + """ + Try to create a pty and if then fails then create a normal + pipe instead. + + @param copy_term_size: If a tty file descriptor is given + then the term size will be copied to the pty. + @type copy_term_size: int + @rtype: tuple + @returns: A tuple of (is_pty, master_fd, slave_fd) where + is_pty is True if a pty was successfully allocated, and + False if a normal pipe was allocated. + """ + + got_pty = False + + global _disable_openpty + if _disable_openpty: + master_fd, slave_fd = os.pipe() + else: + from pty import openpty + try: + master_fd, slave_fd = openpty() + got_pty = True + except EnvironmentError, e: + _disable_openpty = True + writemsg("openpty failed: '%s'\n" % str(e), + noiselevel=-1) + del e + master_fd, slave_fd = os.pipe() + + if got_pty: + # Disable post-processing of output since otherwise weird + # things like \n -> \r\n transformations may occur. + import termios + mode = termios.tcgetattr(slave_fd) + mode[1] &= ~termios.OPOST + termios.tcsetattr(slave_fd, termios.TCSANOW, mode) + + if got_pty and \ + copy_term_size is not None and \ + os.isatty(copy_term_size): + from portage.output import get_term_size, set_term_size + rows, columns = get_term_size() + set_term_size(rows, columns, slave_fd) + + return (got_pty, master_fd, slave_fd) + # XXX This would be to replace getstatusoutput completely. # XXX Issue: cannot block execution. Deadlock condition. def spawn(mystring, mysettings, debug=0, free=0, droppriv=0, sesandbox=0, fakeroot=0, **keywords): @@ -3000,26 +3052,12 @@ def spawn(mystring, mysettings, debug=0, free=0, droppriv=0, sesandbox=0, fakero del keywords["logfile"] if 1 not in fd_pipes or 2 not in fd_pipes: raise ValueError(fd_pipes) - global _disable_openpty - if _disable_openpty: - master_fd, slave_fd = os.pipe() - else: - from pty import openpty - try: - master_fd, slave_fd = openpty() - got_pty = True - except EnvironmentError, e: - _disable_openpty = True - writemsg("openpty failed: '%s'\n" % str(e), noiselevel=1) - del e - master_fd, slave_fd = os.pipe() - if got_pty: - # Disable post-processing of output since otherwise weird - # things like \n -> \r\n transformations may occur. - import termios - mode = termios.tcgetattr(slave_fd) - mode[1] &= ~termios.OPOST - termios.tcsetattr(slave_fd, termios.TCSANOW, mode) + + fd_pipes.setdefault(0, sys.stdin.fileno()) + fd_pipes_orig = fd_pipes.copy() + + got_pty, master_fd, slave_fd = \ + _create_pty_or_pipe(copy_term_size=fd_pipes_orig[1]) # We must set non-blocking mode before we close the slave_fd # since otherwise the fcntl call can fail on FreeBSD (the child @@ -3030,12 +3068,6 @@ def spawn(mystring, mysettings, debug=0, free=0, droppriv=0, sesandbox=0, fakero 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() - if got_pty and os.isatty(fd_pipes_orig[1]): - from portage.output import get_term_size, set_term_size - rows, columns = get_term_size() - set_term_size(rows, columns, slave_fd) fd_pipes[0] = fd_pipes_orig[0] fd_pipes[1] = slave_fd fd_pipes[2] = slave_fd @@ -4800,7 +4832,8 @@ def prepare_build_dirs(myroot, mysettings, cleanup): writemsg("!!! Disabling logging.\n", noiselevel=-1) while "PORT_LOGDIR" in mysettings: del mysettings["PORT_LOGDIR"] - if "PORT_LOGDIR" in mysettings: + if "PORT_LOGDIR" in mysettings and \ + os.access(mysettings["PORT_LOGDIR"], os.W_OK): logid_path = os.path.join(mysettings["PORTAGE_BUILDDIR"], ".logid") if not os.path.exists(logid_path): f = open(logid_path, "w") @@ -5200,8 +5233,7 @@ def doebuild(myebuild, mydo, myroot, mysettings, debug=0, listonly=0, have_build_dirs = True # PORTAGE_LOG_FILE is set above by the prepare_build_dirs() call. logfile = mysettings.get("PORTAGE_LOG_FILE") - if logfile and not os.access(os.path.dirname(logfile), os.W_OK): - logfile = None + if have_build_dirs: env_file = os.path.join(mysettings["T"], "environment") env_stat = None diff --git a/pym/portage/dbapi/bintree.py b/pym/portage/dbapi/bintree.py index 6b3c392bb..bab956b57 100644 --- a/pym/portage/dbapi/bintree.py +++ b/pym/portage/dbapi/bintree.py @@ -832,7 +832,7 @@ class binarytree(object): d2 = pkgindex.packages[i] if d2["CPV"] != cpv: continue - if d2.get("PATH") == d.get("PATH"): + if d2.get("PATH", "") == d.get("PATH", ""): del pkgindex.packages[i] pkgindex.packages.append(d) self._update_pkgindex_header(pkgindex.header) diff --git a/pym/portage/tests/process/__init__.py b/pym/portage/tests/process/__init__.py new file mode 100644 index 000000000..a4a87a461 --- /dev/null +++ b/pym/portage/tests/process/__init__.py @@ -0,0 +1,3 @@ +# Copyright 1998-2008 Gentoo Foundation +# Distributed under the terms of the GNU General Public License v2 +# $Id: __init__.py 6870 2007-06-19 07:22:18Z zmedico $ diff --git a/pym/portage/tests/process/__test__ b/pym/portage/tests/process/__test__ new file mode 100644 index 000000000..e69de29bb diff --git a/pym/portage/tests/process/test_poll.py b/pym/portage/tests/process/test_poll.py new file mode 100644 index 000000000..8a52f4577 --- /dev/null +++ b/pym/portage/tests/process/test_poll.py @@ -0,0 +1,69 @@ +# Copyright 1998-2008 Gentoo Foundation +# Distributed under the terms of the GNU General Public License v2 +# $Id: test_spawn.py 8474 2007-11-09 03:35:38Z zmedico $ + +import errno, os, sys +import fcntl +import termios +import portage +from portage.output import get_term_size, set_term_size +from portage.tests import TestCase +from _emerge import PipeReader, SpawnProcess, TaskScheduler + +class PipeReaderTestCase(TestCase): + + def _create_pipe(self): + return os.pipe() + + def _assertEqual(self, test_string, consumer_value): + self.assertEqual(test_string, consumer_value) + + def testPipeReader(self): + """ + Use a poll loop to read data from a pipe and assert that + the data written to the pipe is identical to the data + read from the pipe. + """ + + test_string = 2 * "blah blah blah\n" + + master_fd, slave_fd = self._create_pipe() + master_file = os.fdopen(master_fd, 'r') + + task_scheduler = TaskScheduler(max_jobs=2) + scheduler = task_scheduler.sched_iface + + producer = SpawnProcess( + args=["bash", "-c", "echo -n '%s'" % test_string], + fd_pipes={1:slave_fd}, scheduler=scheduler) + + consumer = PipeReader( + input_files={"producer" : master_file}, + scheduler=scheduler) + + task_scheduler.add(producer) + task_scheduler.add(consumer) + + def producer_start_cb(task): + os.close(slave_fd) + + producer.addStartListener(producer_start_cb) + task_scheduler.run() + + self._assertEqual(test_string, consumer.getvalue()) + +class PtyReaderTestCase(PipeReaderTestCase): + + def _assertEqual(self, test_string, consumer_value): + if test_string != consumer_value: + # This test is expected to fail on some operating systems + # such as Darwin that do not support poll() on pty devices. + self.todo = True + self.assertEqual(test_string, consumer_value) + + def _create_pipe(self): + + got_pty, master_fd, slave_fd = \ + portage._create_pty_or_pipe(copy_term_size=sys.stdout.fileno()) + + return (master_fd, slave_fd) diff --git a/pym/repoman/utilities.py b/pym/repoman/utilities.py index 28df593cf..f6c50326b 100644 --- a/pym/repoman/utilities.py +++ b/pym/repoman/utilities.py @@ -7,6 +7,7 @@ scan""" import commands +import errno import itertools import logging import os @@ -14,6 +15,7 @@ import sys from portage import output from portage.output import red, green +from portage.process import find_binary from portage import exception from portage import util normalize_path = util.normalize_path @@ -74,7 +76,7 @@ def detect_vcs_conflicts(options, vcs): else: retval = os.system(vcs + " update " + " ".join(myupdates)) if retval != os.EX_OK: - logging.fatal("!!! " + cvs + " exited with an error. Terminating.") + logging.fatal("!!! " + vcs + " exited with an error. Terminating.") sys.exit(retval) @@ -103,8 +105,8 @@ def parse_use_local_desc(mylines, usedict=None): try: pkg, flag = pkg_flag[0].split(":") except ValueError: - raise exception,ParseError("line %d: Malformed input: '%s'" % \ - (linenum + 1, l.rstrip("\n"))) + raise exception.ParseError("line %d: Malformed input: '%s'" % \ + (line_num + 1, l.rstrip("\n"))) usedict.setdefault(pkg, set()) usedict[pkg].add(flag) return usedict -- 2.26.2